From 20911dfd7104ccd10dd0d40fa855d7d529aafb08 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Sun, 4 Oct 2026 17:33:39 +0000 Subject: [PATCH 1/2] feat(executor): execute the exact SQL entropy fallback; Example 2 status Port of the parked #565: SQL LN and a complete, unordered `SUM(column) OVER ()` window (SQLWindowSum) run natively, so the original Q2 SQL of #509 Example 2 executes and matches its recognized entropy form. Partitioned, ordered and finite frames stay refused. planner_layering_example2 records Example 2's status: the three queries lower to Cardinality / FrequencyEntropy / FrequencyL2 (the design's integer Q3 does not), the time filter is a scan predicate, plan_stages plans the workload with UnivMon offered to each statistic, and the exact candidates execute (without the time filter: the runtime has no now() yet). Co-Authored-By: Claude Opus 5.5 --- crates/executor/src/expressions/planner.rs | 25 +- .../executor/src/operators/aggregate/mod.rs | 35 +++ crates/executor/src/operators/mod.rs | 7 +- crates/executor/src/operators/unchecked.rs | 9 + crates/executor/src/physical_planner/mod.rs | 33 +++ crates/executor/tests/blocking_resources.rs | 16 ++ crates/executor/tests/physical_semantics.rs | 56 +++++ .../tests/planner_layering_example2.rs | 217 ++++++++++++++++++ .../tests/sql_frequency_entropy.rs | 76 +++++- 9 files changed, 463 insertions(+), 11 deletions(-) create mode 100644 crates/integration-tests/tests/planner_layering_example2.rs diff --git a/crates/executor/src/expressions/planner.rs b/crates/executor/src/expressions/planner.rs index e4d73e05a..492c82c21 100644 --- a/crates/executor/src/expressions/planner.rs +++ b/crates/executor/src/expressions/planner.rs @@ -127,15 +127,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 { @@ -656,7 +663,7 @@ fn validate(expr: &ScalarExpr, schema: &planner_types::ir::schema::Schema) -> Re 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/executor/src/operators/aggregate/mod.rs b/crates/executor/src/operators/aggregate/mod.rs index 69e6c10d0..e2b67442d 100644 --- a/crates/executor/src/operators/aggregate/mod.rs +++ b/crates/executor/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/executor/src/operators/mod.rs b/crates/executor/src/operators/mod.rs index 94d44dae0..5f521ac3d 100644 --- a/crates/executor/src/operators/mod.rs +++ b/crates/executor/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/executor/src/operators/unchecked.rs b/crates/executor/src/operators/unchecked.rs index 1b5dd7c4f..46e72eac3 100644 --- a/crates/executor/src/operators/unchecked.rs +++ b/crates/executor/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/executor/src/physical_planner/mod.rs b/crates/executor/src/physical_planner/mod.rs index e74a944e0..9f9d689ee 100644 --- a/crates/executor/src/physical_planner/mod.rs +++ b/crates/executor/src/physical_planner/mod.rs @@ -820,6 +820,39 @@ fn bind_operation(node: &PhysicalASAPDAGNode, inputs: &[SchemaRef]) -> Result { + use planner_types::ir::{ + operator::{WindowFrameBound, WindowFrameOffset}, + scalar::ScalarValue, + }; + 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/executor/tests/blocking_resources.rs b/crates/executor/tests/blocking_resources.rs index 23848db2b..41ddb8297 100644 --- a/crates/executor/tests/blocking_resources.rs +++ b/crates/executor/tests/blocking_resources.rs @@ -291,3 +291,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/executor/tests/physical_semantics.rs b/crates/executor/tests/physical_semantics.rs index 970798d67..3f8538de4 100644 --- a/crates/executor/tests/physical_semantics.rs +++ b/crates/executor/tests/physical_semantics.rs @@ -945,3 +945,59 @@ fn sql_sqrt_executes_numeric_and_null_arguments() { matches!(compiled.evaluate(&[Value::Float64(-1.0)]).unwrap(), Value::Float64(v) if v.is_nan()) ); } + +// 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 = ScalarExpr::FunctionCall { + name: "ln".into(), + args: vec![ScalarExpr::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/planner_layering_example2.rs b/crates/integration-tests/tests/planner_layering_example2.rs new file mode 100644 index 000000000..057cff8df --- /dev/null +++ b/crates/integration-tests/tests/planner_layering_example2.rs @@ -0,0 +1,217 @@ +//! #509 Example 2 status: one UnivMon over `src_ip` for distinct, entropy and +//! L2. Records what lowers, plans through `plan_stages` and executes exactly. +mod physical_common; +use asap_executor::values::Value; +use asap_frontend_sql::{lower_sql, SqlCatalog}; +use asap_logical_optimizer::pass1::replacement::Realization; +use asap_plan_selection::{plan_stages, PlanningModels}; +use asap_types::ir::operator::AggIntent; +use asap_types::ir::schema::{DataType, Field, Schema, SketchAlgorithm}; +use asap_types::ir::{NonASAPOp, OperatorNode, QueryRoot, ScalarExpr}; +use asap_types::types::AccuracyTarget; +use asap_types::workload::{DataArrival, DataWorkload, Evidence, EvidenceSource, Rate}; +use std::rc::Rc; + +const WINDOW: &str = "ts >= now() - INTERVAL '1 minute'"; +const Q1: &str = "SELECT COUNT(DISTINCT src_ip) FROM flows WHERE ts >= now() - INTERVAL '1 minute'"; +const Q2: &str = "SELECT -SUM(p * LN(p)) FROM (SELECT COUNT(*) * 1.0 / SUM(COUNT(*)) OVER () AS p FROM flows WHERE ts >= now() - INTERVAL '1 minute' GROUP BY src_ip)"; +/// The design's Q3: its Int64 `c * c` can overflow, so it is not recognized. +const Q3: &str = "SELECT SQRT(SUM(c * c)) FROM (SELECT src_ip, COUNT(*) AS c FROM flows WHERE ts >= now() - INTERVAL '1 minute' GROUP BY src_ip)"; +/// Q3 with a floating product, which the L2 rule accepts. +const Q3_FLOAT: &str = "SELECT SQRT(SUM(CAST(c AS DOUBLE) * CAST(c AS DOUBLE))) FROM (SELECT src_ip, COUNT(*) AS c FROM flows WHERE ts >= now() - INTERVAL '1 minute' GROUP BY src_ip)"; + +fn catalog() -> SqlCatalog { + SqlCatalog::new().with_table( + "flows", + Schema::new(vec![ + Field::plain("ts", DataType::Timestamp, false), + Field::plain("src_ip", DataType::Utf8, false), + ]), + ) +} +fn target(epsilon: f64) -> AccuracyTarget { + AccuracyTarget::EpsilonDelta { + epsilon, + delta: 0.01, + } +} +fn intent(node: &OperatorNode, matches: fn(&AggIntent) -> bool) -> bool { + OperatorNode::reachable(&Rc::new(node.clone())).iter().any(|n| { + matches!(n.non_asap(), Some(NonASAPOp::Aggregate { measures, .. }) if measures.iter().any(matches)) + }) +} +async fn workload() -> Vec<(Rc, AccuracyTarget)> { + let mut roots = vec![]; + for (sql, epsilon) in [(Q1, 0.02), (Q2, 0.05), (Q3_FLOAT, 0.01)] { + let root = lower_sql(sql, &catalog(), target(epsilon)).await.unwrap(); + roots.push((root, target(epsilon))); + } + roots +} + +// The SQL frontend names all three Example 2 computations; the design's +// integer Q3 stays relational because its overflow is observable. +#[tokio::test] +async fn example2_queries_lower_to_frequency_intents() { + let roots = workload().await; + assert!(intent(&roots[0].0, |m| matches!( + m, + AggIntent::Cardinality { .. } + ))); + assert!(intent(&roots[1].0, |m| matches!( + m, + AggIntent::FrequencyEntropy { .. } + ))); + assert!(intent(&roots[2].0, |m| matches!( + m, + AggIntent::FrequencyL2 { .. } + ))); + let design_q3 = lower_sql(Q3, &catalog(), target(0.01)).await.unwrap(); + assert!(!intent(&design_q3, |m| matches!( + m, + AggIntent::FrequencyL2 { .. } + ))); +} + +// `WHERE ts >= now() - INTERVAL '1 minute'` becomes a predicate on the scan, +// not a per-measure FILTER or a TimeRange; Pass 1 accepts the aggregates. +#[tokio::test] +async fn example2_time_filter_is_a_scan_predicate() { + for (root, _) in workload().await { + let nodes = OperatorNode::reachable(&root); + let scans: Vec<_> = nodes + .iter() + .filter_map(|n| match n.non_asap() { + Some(NonASAPOp::Scan { predicates, .. }) => Some(predicates.clone()), + _ => None, + }) + .collect(); + assert!(scans.iter().all(|p| p.len() == 1), "{WINDOW} on every scan"); + assert!(scans + .iter() + .all(|p| matches!(&p[0].0, ScalarExpr::Compare { .. }))); + assert!(nodes.iter().all(|n| !matches!( + n.non_asap(), + Some(NonASAPOp::TimeRange { .. } | NonASAPOp::Filter { .. }) + ))); + assert!(nodes.iter().all(|n| !matches!( + n.non_asap(), + Some(NonASAPOp::Aggregate { filters, .. }) if filters.iter().any(Option::is_some) + ))); + } +} + +// `plan_stages` plans the workload, and Pass 1 offers UnivMon to each of the +// three statistics. Sharing one UnivMon needs the summary-capability rule. +#[tokio::test] +async fn example2_plans_through_the_stage_pipeline() { + let roots = workload().await; + let targets: Vec<_> = roots.iter().map(|(_, t)| Some(t.clone())).collect(); + let data = DataWorkload { + arrival: DataArrival::ContinuouslyIngesting, + data_ingestion_interval: Evidence::default(), + ingestion_volume: Evidence::default(), + ingestion_rate: declared(Rate(100_000.0)), + input_cardinality: declared(10_000_000), + distribution: Evidence::default(), + }; + let run = plan_stages( + roots + .into_iter() + .enumerate() + .map(|(i, (root, _))| (i, QueryRoot::Operator(root))) + .collect(), + &targets, + &data, + PlanningModels::builtin(), + 0, + ) + .expect("Example 2 plans"); + let inventory = &run.stage1[0].inventory; + for statistic in [ + (|m: &AggIntent| matches!(m, AggIntent::Cardinality { .. })) as fn(&AggIntent) -> bool, + |m| matches!(m, AggIntent::FrequencyEntropy { .. }), + |m| matches!(m, AggIntent::FrequencyL2 { .. }), + ] { + let target = inventory + .targets + .iter() + .find(|t| matches!(t.target.non_asap(), Some(NonASAPOp::Aggregate { measures, .. }) if measures.iter().any(statistic))) + .expect("statistic target"); + assert!(target + .alternatives + .iter() + .any(|a| matches!(a, Realization::PassThrough))); + assert!(target.alternatives.iter().any( + |a| matches!(a, Realization::Sketch(kind) if kind.algorithm() == &SketchAlgorithm::UnivMon) + )); + } +} + +fn declared(value: T) -> Evidence { + Evidence { + value: Some(value), + source: EvidenceSource::Declared, + ..Default::default() + } +} + +// The exact candidate of each query executes natively without the time +// filter. With it, binding fails: the runtime has no `now()` yet. +#[tokio::test] +async fn example2_exact_candidates_execute() { + let rows: Vec<_> = ["a", "a", "b"] + .into_iter() + .map(|ip| vec![Value::Timestamp(0), Value::Utf8(ip.into())]) + .collect(); + let mut results = vec![]; + for (sql, epsilon) in [(Q1, 0.02), (Q2, 0.05), (Q3_FLOAT, 0.01)] { + let windowed = lower_sql(sql, &catalog(), target(epsilon)).await.unwrap(); + let message = bind_error(&windowed); + assert!(message.contains("CurrentTimestamp"), "{message}"); + let unwindowed = sql.replace(&format!(" WHERE {WINDOW}"), ""); + let root = lower_sql(&unwindowed, &catalog(), target(epsilon)) + .await + .unwrap(); + results.push(physical_common::execute_raw_rows(&root, rows.clone())); + } + assert!(matches!(results[0][..], [ref r] if matches!(r[..], [Value::Int64(2)]))); + let entropy = -(2.0_f64 / 3.0) * (2.0_f64 / 3.0).ln() - (1.0_f64 / 3.0) * (1.0_f64 / 3.0).ln(); + assert!(matches!(results[1][0][0], Value::Float64(v) if (v - entropy).abs() < 1e-12)); + assert!(matches!(results[2][0][0], Value::Float64(v) if (v - 5.0_f64.sqrt()).abs() < 1e-12)); +} + +/// The error binding `root` against an empty `flows` connector reports. +fn bind_error(root: &Rc) -> String { + use asap_executor::physical_planner::bind_with_data_sources; + use asap_executor::sources::{DataSources, MemorySource}; + use asap_types::ir::export::{NonASAPOpKind, PhysicalASAPOperatorPayload}; + use std::{collections::BTreeMap, sync::Arc}; + let wire = physical_common::compile_physical_asap_dag(root).unwrap(); + let (source, schema) = wire + .nodes + .iter() + .find_map(|node| match &node.payload { + PhysicalASAPOperatorPayload::Relational { + operator: NonASAPOpKind::Scan { source, .. }, + } => Some((source.clone(), node.output_schema.clone())), + _ => None, + }) + .unwrap(); + let mut sources = DataSources::default(); + sources + .register( + source, + Arc::new(MemorySource::new(Arc::new(schema), vec![]).unwrap()), + ) + .unwrap(); + bind_with_data_sources( + &wire, + BTreeMap::new(), + &[u64::from(wire.roots[0].0)], + &sources, + ) + .err() + .expect("binding fails") + .to_string() +} diff --git a/crates/integration-tests/tests/sql_frequency_entropy.rs b/crates/integration-tests/tests/sql_frequency_entropy.rs index a2b61c740..d10afe00b 100644 --- a/crates/integration-tests/tests/sql_frequency_entropy.rs +++ b/crates/integration-tests/tests/sql_frequency_entropy.rs @@ -1,7 +1,8 @@ //! The proposal's entropy idiom executes in native exact frequency operators, with SQL units. mod physical_common; use asap_executor::values::Value; -use asap_frontend_sql::{lower_sql, SqlCatalog}; +use asap_frontend_common::resolve_root; +use asap_frontend_sql::{lower_sql, SqlCatalog, SqlLowerer}; use asap_types::{ ir::schema::{DataType, Field, Schema}, types::AccuracyTarget, @@ -21,6 +22,13 @@ async fn entropy_rewrite_executes_nats_and_empty_population_guard() { let rewritten = lower_sql(sql, &catalog, AccuracyTarget::Exact) .await .unwrap(); + let root = resolve_root( + &SqlLowerer::new(&catalog) + .lower(sql, &AccuracyTarget::Exact) + .await + .unwrap(), + ) + .unwrap(); for (keys, expected) in [ (vec![], None), (vec!["a", "a"], Some(-0.0)), @@ -35,7 +43,13 @@ 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)) => { @@ -48,3 +62,63 @@ 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_executor::physical_planner::bind_with_data_sources; + use asap_executor::sources::{DataSources, MemorySource}; + use asap_types::ir::export::{NonASAPOpKind, PhysicalASAPOperatorPayload}; + use asap_types::ir::operator::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_physical_asap_dag(&root).unwrap(); + let scan_schema = Arc::new( + wire.nodes + .iter() + .find(|node| { + matches!( + node.payload, + PhysicalASAPOperatorPayload::Relational { + operator: 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.roots[0].0)], + &sources, + ) + .err() + .expect("unsupported window rejected"); + assert!( + error + .to_string() + .contains("native SQL window SUM requires the complete unordered relation"), + "{error}" + ); + } +} From 0b2194bbf15df6dbd75f7b1738d0a4fc69f27797 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Sun, 4 Oct 2026 18:24:54 +0000 Subject: [PATCH 2/2] fix: integrate with #593 and #594 plan_stages takes per-root demand since #594, and DataWorkload declares metric types since #593; Example 2's pipeline test runs each query once. Co-Authored-By: Claude Opus 5.5 --- .../tests/planner_layering_example2.rs | 20 ++++++++++++++++--- 1 file changed, 17 insertions(+), 3 deletions(-) diff --git a/crates/integration-tests/tests/planner_layering_example2.rs b/crates/integration-tests/tests/planner_layering_example2.rs index 057cff8df..1f1dc0de6 100644 --- a/crates/integration-tests/tests/planner_layering_example2.rs +++ b/crates/integration-tests/tests/planner_layering_example2.rs @@ -9,7 +9,10 @@ use asap_types::ir::operator::AggIntent; use asap_types::ir::schema::{DataType, Field, Schema, SketchAlgorithm}; use asap_types::ir::{NonASAPOp, OperatorNode, QueryRoot, ScalarExpr}; use asap_types::types::AccuracyTarget; -use asap_types::workload::{DataArrival, DataWorkload, Evidence, EvidenceSource, Rate}; +use asap_types::workload::{ + DataArrival, DataWorkload, Evidence, EvidenceSource, Predictability, QueryRecurrence, Rate, + RootDemand, +}; use std::rc::Rc; const WINDOW: &str = "ts >= now() - INTERVAL '1 minute'"; @@ -106,7 +109,17 @@ async fn example2_time_filter_is_a_scan_predicate() { #[tokio::test] async fn example2_plans_through_the_stage_pipeline() { let roots = workload().await; - let targets: Vec<_> = roots.iter().map(|(_, t)| Some(t.clone())).collect(); + let demand: Vec<_> = roots + .iter() + .map(|(_, t)| RootDemand { + accuracy: Some(t.clone()), + recurrence: QueryRecurrence::OneTime { + invocations: 1, + execute_at: None, + }, + predictability: Predictability::default(), + }) + .collect(); let data = DataWorkload { arrival: DataArrival::ContinuouslyIngesting, data_ingestion_interval: Evidence::default(), @@ -114,6 +127,7 @@ async fn example2_plans_through_the_stage_pipeline() { ingestion_rate: declared(Rate(100_000.0)), input_cardinality: declared(10_000_000), distribution: Evidence::default(), + metric_types: Default::default(), }; let run = plan_stages( roots @@ -121,7 +135,7 @@ async fn example2_plans_through_the_stage_pipeline() { .enumerate() .map(|(i, (root, _))| (i, QueryRoot::Operator(root))) .collect(), - &targets, + &demand, &data, PlanningModels::builtin(), 0,