Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
25 changes: 16 additions & 9 deletions crates/asap-physical-operators/src/expressions/planner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -126,15 +126,22 @@ pub(super) fn evaluate(
.collect::<Result<Vec<_>, _>>()?;
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 {
Expand Down Expand Up @@ -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]
Expand Down
35 changes: 35 additions & 0 deletions crates/asap-physical-operators/src/operators/aggregate/mod.rs
Original file line number Diff line number Diff line change
@@ -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<Self, Error> {
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<usize>,
Expand Down Expand Up @@ -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,
Expand Down
7 changes: 6 additions & 1 deletion crates/asap-physical-operators/src/operators/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -120,6 +120,9 @@ enum Kind {
groups: Vec<usize>,
window: Option<(i64, i64)>,
},
SQLWindowSum {
column: usize,
},
Aggregate {
groups: Vec<usize>,
measures: Vec<Reduction>,
Expand Down Expand Up @@ -281,6 +284,7 @@ impl PhysicalOperator<Batch, SchemaRef> for Operator {
| Kind::SeriesBinary { .. }
| Kind::SeriesHistogramQuantile { .. }
| Kind::SeriesRelabel { .. }
| Kind::SQLWindowSum { .. }
| Kind::Aggregate { .. }
| Kind::Window { .. }
| Kind::Join { .. }
Expand Down Expand Up @@ -333,6 +337,7 @@ impl PhysicalOperator<Batch, SchemaRef> for Operator {
Kind::Filter(_) => "Filter",
Kind::Limit { .. } => "Limit",
Kind::Sort { .. } => "Sort",
Kind::SQLWindowSum { .. } => "SQLWindowSum",
Kind::Aggregate { .. } => "Aggregate",
Kind::Window { .. } => "WindowAggregate",
Kind::SemiJoin { .. } => "SemiJoin",
Expand Down Expand Up @@ -384,7 +389,7 @@ impl PhysicalOperator<Batch, SchemaRef> 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),
Expand Down
9 changes: 9 additions & 0 deletions crates/asap-physical-operators/src/operators/unchecked.rs
Original file line number Diff line number Diff line change
Expand Up @@ -116,6 +116,15 @@ impl TryFrom<UncheckedOperator> 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"));
Expand Down
30 changes: 30 additions & 0 deletions crates/asap-physical-operators/src/physical_planner/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -798,6 +798,36 @@ fn bind_operation(node: &PostAsapDAGNode, inputs: &[SchemaRef]) -> Result<Operat
*offset as u64,
groups(input, partition_by)?,
),
NonASAPOpKind::SQLWindowFunc {
func: planner_types::pre_asap::WindowFuncKind::Sum,
args,
partition_by,
order_by,
frame: Some(frame),
output_name,
} => {
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,
Expand Down
16 changes: 16 additions & 0 deletions crates/asap-physical-operators/tests/blocking_resources.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
56 changes: 56 additions & 0 deletions crates/asap-physical-operators/tests/physical_semantics.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
));
}
}
62 changes: 62 additions & 0 deletions crates/integration-tests/tests/sql_frequency_entropy.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)) => {
Expand All @@ -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}"
);
}
}
1 change: 1 addition & 0 deletions docs/develop_docs/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)
20 changes: 18 additions & 2 deletions docs/develop_docs/planner-layering-status.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Loading