diff --git a/crates/asap-aware-mapping/src/frequency_rewrite.rs b/crates/asap-aware-mapping/src/frequency_rewrite.rs index 1ab027ea4..a8ac1ea53 100644 --- a/crates/asap-aware-mapping/src/frequency_rewrite.rs +++ b/crates/asap-aware-mapping/src/frequency_rewrite.rs @@ -22,6 +22,11 @@ fn expand( fn substitute(expr: &ScalarExpr, cols: &[ProjectItem]) -> Option { Some(match expr { ScalarExpr::Column(id) => cols.get(*id)?.expr.clone(), + ScalarExpr::Literal(_) => expr.clone(), + ScalarExpr::Negative { expr, semantics } => ScalarExpr::Negative { + expr: Box::new(substitute(expr, cols)?), + semantics: *semantics, + }, ScalarExpr::Cast { expr, to, @@ -64,12 +69,7 @@ fn uncast(expr: &ScalarExpr) -> &ScalarExpr { } pub(super) fn frequency_l2_rewrite(root: &Rc) -> Option> { - let NonASAPOp::Project { - cols, - child, - qualifier, - } = root.non_asap()? - else { + let NonASAPOp::Project { cols, child, .. } = root.non_asap()? else { return None; }; let [item] = cols.as_slice() else { @@ -119,6 +119,29 @@ pub(super) fn frequency_l2_rewrite(root: &Rc) -> Option, +) -> Option<(usize, asap_types::types::AccuracyTarget, Rc)> { let NonASAPOp::Aggregate { reduction: Reduction::Reduce(keys), measures, @@ -126,28 +149,20 @@ pub(super) fn frequency_l2_rewrite(root: &Rc) -> Option) -> Option Option<&ScalarExpr> { + let ScalarExpr::Arithmetic { + op: ArithmeticOpKind::Mul, + left, + right, + semantics: ExprSemantics::Sql, + } = product + else { + return None; + }; + for (probability, logarithm) in [(left, right), (right, left)] { + let ScalarExpr::FunctionCall { name, args } = logarithm.as_ref() else { + continue; + }; + if name.eq_ignore_ascii_case("ln") && args.as_slice() == [probability.as_ref().clone()] { + return Some(probability); + } + } + None +} + +fn unit_count_term(expr: &ScalarExpr) -> bool { + match uncast(expr) { + ScalarExpr::Column(1) => true, + ScalarExpr::Arithmetic { op: ArithmeticOpKind::Mul, left, right, semantics: ExprSemantics::Sql } => { + [(left, right), (right, left)].into_iter().any(|(count, scale)| matches!(uncast(count), ScalarExpr::Column(1)) && matches!(scale.as_ref(), ScalarExpr::Literal(ScalarValue::Float64(value)) if *value == 1.0)) + }, + _ => false, + } +} + +pub(super) fn frequency_entropy_rewrite(root: &Rc) -> Option> { + use asap_types::pre_asap::{WindowFrameBound, WindowFrameOffset, WindowFuncKind}; + let NonASAPOp::Project { cols, child, .. } = root.non_asap()? else { + return None; + }; + let [item] = cols.as_slice() else { + return None; + }; + let (expr, outer) = expand(item.expr.clone(), Rc::clone(child))?; + let ScalarExpr::Negative { + expr, + semantics: ExprSemantics::Sql, + } = expr + else { + return None; + }; + if !matches!(uncast(&expr), ScalarExpr::Column(0)) { + return None; + } + let NonASAPOp::Aggregate { + reduction: Reduction::Reduce(keys), + measures, + filters, having: None, - child: Rc::clone(child), + child, + .. + } = outer.non_asap()? + else { + return None; + }; + if keys.is_without() || !keys.keys().is_empty() || any_measure_filtered(filters) { + return None; + } + let [AggIntent::Sum { col: Some(col) }] = measures.as_slice() else { + return None; + }; + let (product, window) = expand(ScalarExpr::Column(*col), Rc::clone(child))?; + let probability = probability_term(&product)?; + let ScalarExpr::Arithmetic { + op: ArithmeticOpKind::Div, + left, + right, + semantics: ExprSemantics::Sql, + } = probability + else { + return None; + }; + if !unit_count_term(left) + || !matches!(uncast(right), ScalarExpr::Column(2)) + || probability.scalar_type(&window.schema).ok()?.0 != DataType::Float64 + { + return None; + } + let NonASAPOp::SQLWindowFunc { + func: WindowFuncKind::Sum, + args, + partition_by, + order_by, + frame: Some(frame), + child, + .. + } = window.non_asap()? + else { + return None; + }; + if partition_by.is_without() + || !partition_by.keys().is_empty() + || !order_by.is_empty() + || args.as_slice() != [ScalarExpr::Column(1)] + { + return None; + } + if !matches!( + frame.start_bound, + WindowFrameBound::Preceding(WindowFrameOffset::Scalar(ScalarValue::Null)) + ) || !matches!( + frame.end_bound, + WindowFrameBound::Following(WindowFrameOffset::Scalar(ScalarValue::Null)) + ) { + return None; + } + let (key, accuracy, input) = grouped_unit_count(child)?; + let nats = ScalarExpr::Arithmetic { + op: ArithmeticOpKind::Mul, + left: Box::new(ScalarExpr::Column(0)), + right: Box::new(ScalarExpr::Literal(ScalarValue::Float64( + std::f64::consts::LN_2, + ))), + semantics: ExprSemantics::Sql, + }; + // -SUM(p*LN(p)) is negative zero for a single-identity population. + let nats = ScalarExpr::Negative { + expr: Box::new(ScalarExpr::Arithmetic { + op: ArithmeticOpKind::Sub, + left: Box::new(ScalarExpr::Literal(ScalarValue::Float64(0.0))), + right: Box::new(nats), + semantics: ExprSemantics::Sql, + }), + semantics: ExprSemantics::Sql, + }; + sql_frequency_result( + root, + input, + AggIntent::FrequencyEntropy { + col: Some(key), + accuracy, + }, + "frequency_entropy", + nats, + ) +} + +// Both rules use an exact population guard: an approximate statistic may be +// zero even for nonempty input, and must not control SQL's NULL result. +fn sql_frequency_result( + root: &Rc, + input: Rc, + measure: AggIntent, + name: &str, + value: ScalarExpr, +) -> Option> { + use asap_types::{ir::Predicate, pre_asap::JoinKind, types::AccuracyTarget}; + let NonASAPOp::Project { qualifier, .. } = root.non_asap()? else { + return None; + }; + let aggregate = |measure, name: &str| { + OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Aggregate { + reduction: Reduction::Reduce(GroupKeys::none()), + measures: vec![measure], + output_names: vec![name.into()], + filters: vec![], + having: None, + child: Rc::clone(&input), + })) + .ok() + }; + let statistic = aggregate(measure, name)?; + let count = aggregate( + AggIntent::Count { + accuracy: AccuracyTarget::Exact, + }, + "population_count", + )?; + let child = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Join { + kind: JoinKind::Cross, + pred: Predicate(ScalarExpr::Literal(ScalarValue::Boolean(true))), + left: statistic, + right: count, })) .ok()?; - // L2 is positive for any nonempty unit-update population. Restore SQL SUM's - // NULL on an empty relation without introducing another count computation. let rewritten = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Project { cols: vec![ProjectItem { alias: Some(root.schema.fields.first()?.name.clone()), @@ -176,9 +361,9 @@ pub(super) fn frequency_l2_rewrite(root: &Rc) -> Option) -> Option) -> Vec { + if let Some(rewritten) = crate::frequency_rewrite::frequency_entropy_rewrite(target.root) { + return vec![ReplacementSubDAG { + strategy: "SemanticEquivalentRewriteStrategy", + replacement: Replacement::SubDAG(rewritten), + provenance: crate::replacement::ReplacementProvenance::LogicalRewrite, + rationale: "recognize SQL natural-log entropy with explicit bits-to-nats conversion and exact empty-population guard".into(), + }]; + } if let Some(rewritten) = crate::frequency_rewrite::frequency_l2_rewrite(target.root) { return vec![ReplacementSubDAG { strategy: "SemanticEquivalentRewriteStrategy", diff --git a/crates/frontend-sql/tests/frequency_entropy.rs b/crates/frontend-sql/tests/frequency_entropy.rs new file mode 100644 index 000000000..9d029c07c --- /dev/null +++ b/crates/frontend-sql/tests/frequency_entropy.rs @@ -0,0 +1,116 @@ +//! SQL entropy recognition keeps natural-log units and the original relational path. +use asap_aware_mapping::{ + replacement::{Replacement, ReplacementStrategy, TargetSubDAG}, + SemanticEquivalentRewriteStrategy, +}; +use asap_frontend_sql::{lower_sql, SqlCatalog}; +use asap_types::{ + ir::{NonASAPOp, OperatorNode}, + pre_asap::{AggIntent, DataType, Field, Schema}, + types::AccuracyTarget, +}; +fn catalog(nullable: bool) -> SqlCatalog { + SqlCatalog::new().with_table( + "flows", + Schema::new(vec![Field::plain("src_ip", DataType::Utf8, nullable)]), + ) +} +fn has_entropy(node: &OperatorNode) -> bool { + if let Some(NonASAPOp::Aggregate { measures, .. }) = node.non_asap() { + if measures + .iter() + .any(|m| matches!(m, AggIntent::FrequencyEntropy { .. })) + { + return true; + } + } + node.children().iter().any(|child| has_entropy(child)) +} +// The proposal's natural-log idiom adds an entropy intent and preserves output schema. +#[tokio::test] +async fn recognizes_sql_entropy_in_nats() { + let root = lower_sql("SELECT -SUM(p * LN(p)) AS entropy FROM (SELECT COUNT(*) * 1.0 / SUM(COUNT(*)) OVER () AS p FROM flows GROUP BY src_ip) f", &catalog(false), AccuracyTarget::Exact).await.unwrap(); + let replacements = SemanticEquivalentRewriteStrategy.replacements(&TargetSubDAG::new(&root)); + let rewritten = replacements + .iter() + .find_map(|r| match &r.replacement { + Replacement::SubDAG(node) if has_entropy(node) => Some(node), + _ => None, + }) + .expect("entropy candidate"); + assert_eq!(root.schema, rewritten.schema); + assert!(!has_entropy(&root)); +} + +// Changes to units, normalization, window coverage or the counted population are not entropy rewrites. +#[tokio::test] +async fn declines_non_equivalent_entropy_shapes() { + for (sql, nullable) in [ + ("SELECT -SUM(p*LN(p)) FROM (SELECT COUNT(*)*2.0/SUM(COUNT(*)) OVER () AS p FROM flows GROUP BY src_ip) f", false), + ("SELECT -SUM(p*LOG2(p)) FROM (SELECT COUNT(*)*1.0/SUM(COUNT(*)) OVER () AS p FROM flows GROUP BY src_ip) f", false), + ("SELECT -SUM(p*LN(p)) FROM (SELECT COUNT(*)*1.0/SUM(COUNT(*)) OVER (PARTITION BY src_ip) AS p FROM flows GROUP BY src_ip) f", false), + ("SELECT -SUM(p*LN(p)) FROM (SELECT COUNT(*)*1.0/SUM(COUNT(*)) OVER (ORDER BY src_ip) AS p FROM flows GROUP BY src_ip) f", false), + ("SELECT -SUM(p*LN(p)) FROM (SELECT COUNT(*)*1.0/SUM(COUNT(*)) OVER () AS p FROM flows GROUP BY src_ip HAVING COUNT(*) > 1) f", false), + ("SELECT -SUM(p*LN(p)) FROM (SELECT COUNT(*)*1.0/SUM(COUNT(*)) OVER () AS p FROM flows GROUP BY src_ip) f", true), + ] { + let root = lower_sql(sql, &catalog(nullable), AccuracyTarget::Exact).await.unwrap(); + let candidates = SemanticEquivalentRewriteStrategy.replacements(&TargetSubDAG::new(&root)); + assert!(!candidates.iter().any(|r| matches!(&r.replacement, Replacement::SubDAG(n) if has_entropy(n)))); + } +} + +// Default search retains both the original exact SQL graph and the recognized entropy graph. +#[tokio::test] +async fn entropy_search_preserves_relational_alternative() { + use asap_aware_mapping::replacement::{default_strategies, search_workload_with_targets}; + let root = lower_sql("SELECT -SUM(p*LN(p)) FROM (SELECT COUNT(*)*1.0/SUM(COUNT(*)) OVER () AS p FROM flows GROUP BY src_ip) f", &catalog(false), AccuracyTarget::Exact).await.unwrap(); + let space = search_workload_with_targets( + vec![(0, root, Some(AccuracyTarget::Exact))], + &default_strategies(), + &asap_aware_mapping::accuracy::DefaultAccuracyModel, + ); + let inventory = space.enumerate_candidate_dags(1000).unwrap(); + assert!(inventory + .candidates + .iter() + .any(|candidate| has_entropy(&candidate[0].1))); + assert!(inventory + .candidates + .iter() + .any(|candidate| !has_entropy(&candidate[0].1))); +} + +// Entropy receives the source query budget while its empty-population guard stays exact. +#[tokio::test] +async fn entropy_accuracy_and_population_guard_are_separate() { + let target = AccuracyTarget::EpsilonDelta { + epsilon: 0.05, + delta: 0.01, + }; + let root = lower_sql("SELECT -SUM(p*LN(p)) FROM (SELECT COUNT(*)*1.0/SUM(COUNT(*)) OVER () AS p FROM flows GROUP BY src_ip) f", &catalog(false), target.clone()).await.unwrap(); + let candidates = SemanticEquivalentRewriteStrategy.replacements(&TargetSubDAG::new(&root)); + let Replacement::SubDAG(node) = &candidates[0].replacement else { + panic!("rewrite"); + }; + let NonASAPOp::Project { child, .. } = node.expect_non_asap() else { + panic!("project"); + }; + let NonASAPOp::Join { left, right, .. } = child.expect_non_asap() else { + panic!("guarded statistic"); + }; + let NonASAPOp::Aggregate { measures, .. } = left.expect_non_asap() else { + panic!("entropy"); + }; + assert!( + matches!(&measures[0], AggIntent::FrequencyEntropy { accuracy, .. } if *accuracy == target) + ); + let NonASAPOp::Aggregate { measures, .. } = right.expect_non_asap() else { + panic!("count"); + }; + assert!(matches!( + measures.as_slice(), + [AggIntent::Count { + accuracy: AccuracyTarget::Exact + }] + )); +} diff --git a/crates/frontend-sql/tests/frequency_l2.rs b/crates/frontend-sql/tests/frequency_l2.rs index 8c32b6340..cde8a9c4c 100644 --- a/crates/frontend-sql/tests/frequency_l2.rs +++ b/crates/frontend-sql/tests/frequency_l2.rs @@ -74,7 +74,10 @@ async fn frequency_l2_preserves_accuracy_target() { let NonASAPOp::Project { child, .. } = node.expect_non_asap() else { panic!("project"); }; - let NonASAPOp::Aggregate { measures, .. } = child.expect_non_asap() else { + let NonASAPOp::Join { left, .. } = child.expect_non_asap() else { + panic!("guarded statistic"); + }; + let NonASAPOp::Aggregate { measures, .. } = left.expect_non_asap() else { panic!("aggregate"); }; assert!(matches!(&measures[0], AggIntent::FrequencyL2 { accuracy, .. } if *accuracy == target)); @@ -100,3 +103,32 @@ async fn default_search_keeps_exact_sql_and_frequency_alternatives() { .iter() .any(|candidate| !has_l2(&candidate[0].1))); } + +// An approximate L2 estimate must never decide whether SQL returns NULL. +#[tokio::test] +async fn frequency_empty_input_guard_uses_an_exact_count() { + let target = AccuracyTarget::EpsilonDelta { + epsilon: 0.01, + delta: 0.01, + }; + let root = lower_sql("SELECT SQRT(SUM(c*c)) FROM (SELECT src_ip, CAST(COUNT(*) AS DOUBLE) AS c FROM flows GROUP BY src_ip) f", &catalog(false), target).await.unwrap(); + let candidates = SemanticEquivalentRewriteStrategy.replacements(&TargetSubDAG::new(&root)); + let Replacement::SubDAG(node) = &candidates[0].replacement else { + panic!("rewrite"); + }; + let NonASAPOp::Project { child, .. } = node.expect_non_asap() else { + panic!("project"); + }; + let NonASAPOp::Join { right, .. } = child.expect_non_asap() else { + panic!("exact population guard"); + }; + let NonASAPOp::Aggregate { measures, .. } = right.expect_non_asap() else { + panic!("count"); + }; + assert!(matches!( + measures.as_slice(), + [AggIntent::Count { + accuracy: AccuracyTarget::Exact + }] + )); +} diff --git a/crates/integration-tests/tests/sql_frequency_entropy.rs b/crates/integration-tests/tests/sql_frequency_entropy.rs new file mode 100644 index 000000000..3fcd628d3 --- /dev/null +++ b/crates/integration-tests/tests/sql_frequency_entropy.rs @@ -0,0 +1,55 @@ +//! The proposal's entropy idiom executes in native exact frequency operators, with SQL units. +mod physical_common; +use asap_aware_mapping::{ + replacement::{Replacement, ReplacementStrategy, TargetSubDAG}, + SemanticEquivalentRewriteStrategy, +}; +use asap_frontend_sql::{lower_sql, SqlCatalog}; +use asap_physical_operators::values::Value; +use asap_types::{ + pre_asap::{DataType, Field, Schema}, + types::AccuracyTarget, +}; + +// Native wire execution produces nats, NULL for no population and SQL's negative zero for one identity. +#[tokio::test] +async fn entropy_rewrite_executes_nats_and_empty_population_guard() { + let catalog = SqlCatalog::new().with_table( + "flows", + Schema::new(vec![ + Field::plain("src_ip", DataType::Utf8, false), + Field::plain("keep", DataType::Bool, false), + ]), + ); + let root = lower_sql("SELECT -SUM(p*LN(p)) AS entropy FROM (SELECT COUNT(*)*1.0/SUM(COUNT(*)) OVER () AS p FROM flows WHERE keep GROUP BY src_ip) f", &catalog, AccuracyTarget::Exact).await.unwrap(); + let candidates = SemanticEquivalentRewriteStrategy.replacements(&TargetSubDAG::new(&root)); + let Replacement::SubDAG(rewritten) = &candidates[0].replacement else { + panic!("entropy rewrite"); + }; + for (keys, expected) in [ + (vec![], None), + (vec!["a", "a"], Some(-0.0)), + (vec!["a", "a", "b", "b"], Some(std::f64::consts::LN_2)), + ( + vec!["a", "a", "a", "b"], + Some(-0.75_f64 * 0.75_f64.ln() - 0.25_f64 * 0.25_f64.ln()), + ), + ] { + let mut rows: Vec<_> = keys + .into_iter() + .map(|key| vec![Value::Utf8(key.into()), Value::Bool(true)]) + .collect(); + rows.push(vec![Value::Utf8("discard".into()), Value::Bool(false)]); + let actual = physical_common::execute_raw_rows(rewritten, rows); + match (&actual[0][0], expected) { + (Value::Null, None) => {} + (Value::Float64(value), Some(expected)) => { + assert!((value - expected).abs() < 1e-12); + if expected == 0.0 { + assert_eq!(value.to_bits(), (-0.0_f64).to_bits()); + } + } + other => panic!("wrong SQL entropy: {other:?}"), + } + } +} diff --git a/docs/develop_docs/planner-layering-status.md b/docs/develop_docs/planner-layering-status.md index b0b6c7715..a24c98d89 100644 --- a/docs/develop_docs/planner-layering-status.md +++ b/docs/develop_docs/planner-layering-status.md @@ -6,7 +6,7 @@ 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. | Floating-point SQL frequency L2 products now have a conservative logical rewrite. Integer products and the entropy idiom remain unrecognized. Preserve alias lineage, filters, NULL groups, empty inputs, count overflow and entropy units when adding recognition. | +| Language frontends and common logical IR | SQL/PromQL/MetricsQL lower to unified operators and scalars. | Floating-point SQL frequency L2 products and the normalized natural-log entropy idiom now have conservative logical rewrites. Integer L2 products remain unrecognized. 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. | @@ -74,3 +74,21 @@ intent column list retains the existing sample-value convention. `integration-tests/tests/sql_cardinality.rs` lowers `COUNT(DISTINCT src_ip)` and executes its exact path through raw scan predicates and native wire binding. Filtered aggregate measures remain outside native binding's existing scope. + +## SQL entropy follow-up acceptance + +The proposal's `-SUM(p * LN(p))` form now exposes `FrequencyEntropy` when `p` +is a grouped unit count divided by the full, unpartitioned count window. +Recognition requires one nonnullable Boolean/Int64/Utf8 identity, no +measure filters/HAVING, no ordering, and an unbounded window in both directions. +The rewrite converts the core intent's bits to nats using `ln(2)`, retains +SQL's negative zero for a single identity and uses an exact population count +to restore empty-input NULL. L2 now uses the same exact population guard, so a +zero approximate estimate cannot decide whether SQL returns NULL. + +Frontend tests cover recognition, non-equivalent probability/window/unit +shapes, candidate retention and accuracy propagation. Native wire execution +checks filtering, nats, empty population, negative zero and unequal frequencies. +The original entropy SQL graph remains an alternative, but native SQL window +binding for that original graph is still a runtime gap at this step. Neither +recognition nor the exact path proves UnivMon's probabilistic accuracy bound.