From 8a0a8d5978717c846a45d0f10522b3f1a8ea3f5f Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Sat, 3 Oct 2026 15:34:32 +0000 Subject: [PATCH] feat(runtime): execute exact SQL entropy fallback --- .../src/expressions/planner.rs | 25 +++++--- .../src/operators/aggregate/mod.rs | 35 +++++++++++ .../src/operators/mod.rs | 7 ++- .../src/operators/unchecked.rs | 9 +++ .../src/physical_planner/mod.rs | 30 +++++++++ .../tests/blocking_resources.rs | 16 +++++ .../tests/physical_semantics.rs | 56 +++++++++++++++++ .../tests/sql_frequency_entropy.rs | 62 +++++++++++++++++++ docs/develop_docs/README.md | 1 + docs/develop_docs/planner-layering-status.md | 20 +++++- 10 files changed, 249 insertions(+), 12 deletions(-) diff --git a/crates/asap-physical-operators/src/expressions/planner.rs b/crates/asap-physical-operators/src/expressions/planner.rs index 11bea51ae..c854666de 100644 --- a/crates/asap-physical-operators/src/expressions/planner.rs +++ b/crates/asap-physical-operators/src/expressions/planner.rs @@ -126,15 +126,22 @@ pub(super) fn evaluate( .collect::, _>>()?; return Ok(Value::Float64(promql_function(name, &values)?)); } - if name.eq_ignore_ascii_case("sqrt") { - return match evaluate(&args[0], row, schema)? { - Value::Null => Ok(Value::Null), - Value::Float64(value) => Ok(Value::Float64(value.sqrt())), - Value::Int64(value) => Ok(Value::Float64((value as f64).sqrt())), - _ => Err(Error::Invalid( - "SQL sqrt requires a numeric argument".into(), - )), + if name.eq_ignore_ascii_case("sqrt") || name.eq_ignore_ascii_case("ln") { + let value = match evaluate(&args[0], row, schema)? { + Value::Null => return Ok(Value::Null), + Value::Float64(value) => value, + Value::Int64(value) => value as f64, + _ => { + return Err(Error::Invalid( + "SQL math function requires a numeric argument".into(), + )) + } }; + return Ok(Value::Float64(if name.eq_ignore_ascii_case("sqrt") { + value.sqrt() + } else { + value.ln() + })); } if name == "promql_drop_metric_name" { let Value::Utf8(encoded) = evaluate(&args[0], row, schema)? else { @@ -655,7 +662,7 @@ fn validate(expr: &ScalarExpr, schema: &planner_types::pre_asap::Schema) -> Resu Ok(()) } ScalarExpr::FunctionCall { name, args } => { - if name.eq_ignore_ascii_case("sqrt") { + if name.eq_ignore_ascii_case("sqrt") || name.eq_ignore_ascii_case("ln") { if args.len() != 1 || !matches!( args[0] diff --git a/crates/asap-physical-operators/src/operators/aggregate/mod.rs b/crates/asap-physical-operators/src/operators/aggregate/mod.rs index 11b5d6049..b1ada24d8 100644 --- a/crates/asap-physical-operators/src/operators/aggregate/mod.rs +++ b/crates/asap-physical-operators/src/operators/aggregate/mod.rs @@ -1,5 +1,20 @@ use super::*; impl Operator { + /// A SQL SUM window over the complete, unordered input relation. + pub fn sql_window_sum(input: SchemaRef, column: usize, name: String) -> Result { + let (dtype, _) = plain(&input, column)?; + if !matches!(dtype, DataType::Int64 | DataType::Float64) { + return Err(invalid("SQL window SUM requires a numeric column")); + } + let mut output = (*input).clone(); + output.fields.push(result_field(&name, dtype.clone(), true)); + Ok(Self { + kind: Kind::SQLWindowSum { column }, + inputs: vec![input], + output: Arc::new(output), + }) + } + pub fn aggregate( input: SchemaRef, groups: Vec, @@ -171,6 +186,26 @@ pub(super) fn execute<'a>( Ok(futures::stream::once(async move { let (rows, _memory) = collect_rows(input, &context).await?; let result = match &operator.kind { + Kind::SQLWindowSum { column } => { + let mut work = Cooperative::new(&context); + let total = reduce_one( + &rows, + &Reduction::Sum(*column), + &operator.inputs[0], + &mut work, + &context, + ) + .await?; + let mut workspace = Workspace::new(&context)?; + let mut result = Vec::with_capacity(rows.len()); + for mut row in rows { + work.checkpoint().await?; + row.push(total.clone()); + workspace.grow(row_bytes(&row))?; + result.push(row); + } + result + } Kind::Window { intent, coordinate, diff --git a/crates/asap-physical-operators/src/operators/mod.rs b/crates/asap-physical-operators/src/operators/mod.rs index 517019118..be3155160 100644 --- a/crates/asap-physical-operators/src/operators/mod.rs +++ b/crates/asap-physical-operators/src/operators/mod.rs @@ -120,6 +120,9 @@ enum Kind { groups: Vec, window: Option<(i64, i64)>, }, + SQLWindowSum { + column: usize, + }, Aggregate { groups: Vec, measures: Vec, @@ -281,6 +284,7 @@ impl PhysicalOperator for Operator { | Kind::SeriesBinary { .. } | Kind::SeriesHistogramQuantile { .. } | Kind::SeriesRelabel { .. } + | Kind::SQLWindowSum { .. } | Kind::Aggregate { .. } | Kind::Window { .. } | Kind::Join { .. } @@ -333,6 +337,7 @@ impl PhysicalOperator for Operator { Kind::Filter(_) => "Filter", Kind::Limit { .. } => "Limit", Kind::Sort { .. } => "Sort", + Kind::SQLWindowSum { .. } => "SQLWindowSum", Kind::Aggregate { .. } => "Aggregate", Kind::Window { .. } => "WindowAggregate", Kind::SemiJoin { .. } => "SemiJoin", @@ -384,7 +389,7 @@ impl PhysicalOperator for Operator { Kind::Filter(_) => filter::execute(self, inputs, context), Kind::Limit { .. } => limit::execute(self, inputs, context), Kind::Sort { .. } => sort::execute(self, inputs, context), - Kind::Window { .. } | Kind::Aggregate { .. } => { + Kind::SQLWindowSum { .. } | Kind::Window { .. } | Kind::Aggregate { .. } => { aggregate::execute(self, inputs, context) } Kind::Join { .. } | Kind::SemiJoin { .. } => joins::execute(self, inputs, context), diff --git a/crates/asap-physical-operators/src/operators/unchecked.rs b/crates/asap-physical-operators/src/operators/unchecked.rs index 1b5dd7c4f..46e72eac3 100644 --- a/crates/asap-physical-operators/src/operators/unchecked.rs +++ b/crates/asap-physical-operators/src/operators/unchecked.rs @@ -116,6 +116,15 @@ impl TryFrom for Operator { groups, window, } => Operator::window(input(0)?, *intent, coordinate, value, groups, window)?, + Kind::SQLWindowSum { column } => { + let name = output + .fields + .last() + .ok_or_else(|| invalid("SQL window SUM output missing"))? + .name + .clone(); + Operator::sql_window_sum(input(0)?, column, name)? + } Kind::Aggregate { groups, measures } => { if groups.len() + measures.len() != output.fields.len() { return Err(invalid("aggregate width mismatch")); diff --git a/crates/asap-physical-operators/src/physical_planner/mod.rs b/crates/asap-physical-operators/src/physical_planner/mod.rs index 9372589ca..fb7af68fd 100644 --- a/crates/asap-physical-operators/src/physical_planner/mod.rs +++ b/crates/asap-physical-operators/src/physical_planner/mod.rs @@ -798,6 +798,36 @@ fn bind_operation(node: &PostAsapDAGNode, inputs: &[SchemaRef]) -> Result { + use planner_types::pre_asap::{ScalarValue, WindowFrameBound, WindowFrameOffset}; + let [WireScalarExpr::Column(column)] = args.as_slice() else { + return Err(invalid("SQL window SUM requires one column")); + }; + if partition_by.is_without() + || !partition_by.keys().is_empty() + || !order_by.is_empty() + || !matches!( + frame.start_bound, + WindowFrameBound::Preceding(WindowFrameOffset::Scalar(ScalarValue::Null)) + ) + || !matches!( + frame.end_bound, + WindowFrameBound::Following(WindowFrameOffset::Scalar(ScalarValue::Null)) + ) + { + return Err(invalid( + "native SQL window SUM requires the complete unordered relation", + )); + } + Operator::sql_window_sum(input.clone(), *column, output_name.clone()) + } NonASAPOpKind::Aggregate { reduction, measures, diff --git a/crates/asap-physical-operators/tests/blocking_resources.rs b/crates/asap-physical-operators/tests/blocking_resources.rs index 873eef8d5..94d3eabb0 100644 --- a/crates/asap-physical-operators/tests/blocking_resources.rs +++ b/crates/asap-physical-operators/tests/blocking_resources.rs @@ -292,3 +292,19 @@ fn frequency_dictionary_enforces_memory_budget() { assert_eq!(run.retained_bytes(), 0); } } + +// A complete SUM window accounts for its expanded output and frees memory on failure. +#[test] +fn complete_window_sum_enforces_workspace_budget() { + let sources = source(64); + let run = context(12_000); + let inputs = sources.execute(&[0], run.clone()).unwrap(); + let operator = Operator::sql_window_sum(schema(1), 0, "total".into()).unwrap(); + let mut output = operator.start(inputs, run.clone()).unwrap(); + assert!(matches!( + block_on(output.next()), + Some(Err(Error::MemoryLimit)) + )); + 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 2ce2bda66..5c7045af7 100644 --- a/crates/asap-physical-operators/tests/physical_semantics.rs +++ b/crates/asap-physical-operators/tests/physical_semantics.rs @@ -948,3 +948,59 @@ fn exact_cardinality_grouping_normalizes_float_identities() { [Value::Int64(2), Value::Int64(0)] )); } + +// A complete SQL SUM window keeps every row and appends one nullable total, including recovery. +#[test] +fn complete_sql_sum_window_preserves_rows_and_nulls() { + let input = schema(&[("v", DataType::Int64, true)]); + let operator = Operator::sql_window_sum(input.clone(), 0, "total".into()).unwrap(); + let operator: Operator = + serde_json::from_slice(&serde_json::to_vec(&operator).unwrap()).unwrap(); + for (rows, expected) in [ + (vec![], None), + (vec![vec![Value::Null]], None), + ( + vec![ + vec![Value::Int64(1)], + vec![Value::Null], + vec![Value::Int64(3)], + ], + Some(4), + ), + ] { + let original = rows.clone(); + let actual = unary(input.clone(), vec![rows], operator.clone()); + assert_eq!(actual.len(), original.len()); + for (row, original) in actual.iter().zip(original) { + assert_eq!(row[0].key().unwrap(), original[0].key().unwrap()); + match (&row[1], expected) { + (Value::Null, None) => {} + (Value::Int64(value), Some(expected)) => assert_eq!(*value, expected), + other => panic!("wrong complete-window sum: {other:?}"), + } + } + } +} + +// SQL LN preserves nullable numeric signatures and natural-log units. +#[test] +fn sql_ln_executes_numeric_and_null_arguments() { + for (dtype, value) in [ + (DataType::Int64, Value::Int64(2)), + (DataType::Float64, Value::Float64(2.0)), + ] { + let input = schema(&[("v", dtype, true)]); + let expression = QueryExpr::FunctionCall { + name: "ln".into(), + args: vec![QueryExpr::Column(0)], + }; + let compiled = CompiledExpression::compile(&expression, &input).unwrap(); + assert!( + matches!(compiled.evaluate(&[value]).unwrap(), Value::Float64(v) if v == std::f64::consts::LN_2) + ); + assert!(matches!( + compiled.evaluate(&[Value::Null]).unwrap(), + Value::Null + )); + } +} diff --git a/crates/integration-tests/tests/sql_frequency_entropy.rs b/crates/integration-tests/tests/sql_frequency_entropy.rs index 3fcd628d3..bf64a723c 100644 --- a/crates/integration-tests/tests/sql_frequency_entropy.rs +++ b/crates/integration-tests/tests/sql_frequency_entropy.rs @@ -40,7 +40,14 @@ async fn entropy_rewrite_executes_nats_and_empty_population_guard() { .map(|key| vec![Value::Utf8(key.into()), Value::Bool(true)]) .collect(); rows.push(vec![Value::Utf8("discard".into()), Value::Bool(false)]); + let original = physical_common::execute_raw_rows(&root, rows.clone()); let actual = physical_common::execute_raw_rows(rewritten, rows); + match (&original[0][0], &actual[0][0]) { + (Value::Null, Value::Null) => {} + (Value::Float64(a), Value::Float64(b)) => assert!((a - b).abs() < 1e-12), + other => panic!("original SQL differs from entropy rewrite: {other:?}"), + } + match (&actual[0][0], expected) { (Value::Null, None) => {} (Value::Float64(value), Some(expected)) => { @@ -53,3 +60,58 @@ async fn entropy_rewrite_executes_nats_and_empty_population_guard() { } } } + +// Native binding refuses partial, partitioned or ordered SUM windows instead of treating them as totals. +#[tokio::test] +async fn native_sql_sum_window_rejects_other_frames() { + use asap_physical_operators::physical_planner::bind_with_data_sources; + use asap_physical_operators::sources::{DataSources, MemorySource}; + use asap_types::pre_asap::Source; + use std::{collections::BTreeMap, sync::Arc}; + let schema = Schema::new(vec![Field::plain("src_ip", DataType::Int64, false)]); + let catalog = SqlCatalog::new().with_table("flows", schema.clone()); + for sql in [ + "SELECT SUM(src_ip) OVER (PARTITION BY src_ip) FROM flows", + "SELECT SUM(src_ip) OVER (ORDER BY src_ip) FROM flows", + "SELECT SUM(src_ip) OVER (ROWS BETWEEN 1 PRECEDING AND CURRENT ROW) FROM flows", + ] { + let root = lower_sql(sql, &catalog, AccuracyTarget::Exact) + .await + .unwrap(); + let wire = physical_common::compile_post_asap_dag(&root).unwrap(); + let scan_schema = Arc::new( + wire.nodes + .iter() + .find(|node| { + matches!( + node.payload, + asap_types::ir::export::PostAsapOperatorPayload::Relational { + operator: asap_types::ir::export::NonASAPOpKind::Scan { .. } + } + ) + }) + .unwrap() + .output_schema + .clone(), + ); + let mut sources = DataSources::default(); + sources + .register( + Source::Table { + table_ref: "flows".into(), + }, + Arc::new(MemorySource::new(scan_schema, vec![]).unwrap()), + ) + .unwrap(); + let error = + bind_with_data_sources(&wire, BTreeMap::new(), &[u64::from(wire.root.0)], &sources) + .err() + .expect("unsupported window rejected"); + assert!( + error + .to_string() + .contains("native SQL window SUM requires the complete unordered relation"), + "{error}" + ); + } +} diff --git a/docs/develop_docs/README.md b/docs/develop_docs/README.md index 817722adf..ed2585ef2 100644 --- a/docs/develop_docs/README.md +++ b/docs/develop_docs/README.md @@ -16,5 +16,6 @@ formats, evidence, and verification workflows. - [Physical handoff cost references](physical-handoff-costs.md), [storage operations](storage-operation-costs.md) - [Replacement explanations](replacement-explanations.md) - [Physical compile coverage for deployment computation](physical-compile-coverage.md) +- [Planner-layering implementation status and follow-up scopes](planner-layering-status.md) - [Planner vocabulary migration (#427)](planner-vocabulary-migration.md) diff --git a/docs/develop_docs/planner-layering-status.md b/docs/develop_docs/planner-layering-status.md index a24c98d89..eb6008d40 100644 --- a/docs/develop_docs/planner-layering-status.md +++ b/docs/develop_docs/planner-layering-status.md @@ -89,6 +89,22 @@ 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 +The original entropy SQL graph remains an alternative. Its native SQL window +binding is supplied by the next follow-up below. Neither recognition nor the exact path proves UnivMon's probabilistic accuracy bound. + +## Native entropy fallback follow-up acceptance + +Native binding now implements SQL `LN` and the exact `SUM(column) OVER ()` +window over a complete, unordered relation with unbounded start/end bounds. +The window appends its total to every original row, keeps input metadata, +propagates all-NULL totals, supports recovery and obeys workspace limits. +Partitioned, ordered and finite frames remain explicitly unsupported. This +operator executes one SQL relation; it is unrelated to the proposal's missing +streaming tumbling/sliding/EH summary-window planning. + +The entropy acceptance fixture now executes both the original SQL and its +frequency rewrite, comparing filtering, natural-log units, empty-input NULL +and unequal frequencies. Raw analytical cost lowering still excludes this +unordered SUM window; supplying native execution does not provide missing +cost evidence or extend the analytical adapter's existing ordered-window rule.