1pub mod channel_derivation;
8pub mod lookup_join;
10pub mod predicate_split;
12pub mod streaming_optimizer;
14
15#[allow(clippy::disallowed_types)] use std::collections::HashMap;
17use std::sync::Arc;
18
19use arrow::datatypes::{Field, Schema, SchemaRef};
20use datafusion::logical_expr::LogicalPlan;
21use datafusion::prelude::SessionContext;
22use sqlparser::ast::{ObjectName, Select, SetExpr, Statement, TableFactor};
23
24use crate::parser::aggregation_parser::analyze_aggregates;
25use crate::parser::analytic_parser::{
26 analyze_analytic_functions, analyze_window_frames, FrameBound,
27};
28use crate::parser::join_parser::{analyze_joins, JoinAnalysis, JoinType, MultiJoinAnalysis};
29use crate::parser::lookup_table::{validate_properties, LookupTableProperties};
30use crate::parser::order_analyzer::analyze_order_by;
31use crate::parser::{
32 CreateLookupTableStatement, CreateSinkStatement, CreateSourceStatement, EmitClause, SinkFrom,
33 StreamingStatement, WindowFunction, WindowRewriter,
34};
35use crate::translator::{
36 AnalyticWindowConfig, HavingFilterConfig, JoinOperatorConfig, OrderOperatorConfig,
37 WindowFrameConfig, WindowOperatorConfig,
38};
39
40#[derive(Debug, Clone)]
42pub struct LookupTableInfo {
43 pub name: String,
45 pub columns: Vec<(String, String)>,
47 pub primary_key: Vec<String>,
49 pub properties: LookupTableProperties,
51 pub arrow_schema: SchemaRef,
53 pub raw_options: HashMap<String, String>,
55}
56
57pub struct StreamingPlanner {
59 sources: HashMap<String, SourceInfo>,
61 sinks: HashMap<String, SinkInfo>,
63 lookup_tables: HashMap<String, LookupTableInfo>,
65 windowed_views: std::collections::HashSet<String>,
67}
68
69#[derive(Debug, Clone)]
71pub struct SourceInfo {
72 pub name: String,
74 pub watermark_column: Option<String>,
76 pub options: HashMap<String, String>,
78}
79
80#[derive(Debug, Clone)]
82pub struct SinkInfo {
83 pub name: String,
85 pub from: String,
87 pub options: HashMap<String, String>,
89}
90
91fn is_inline_unnest(factor: &TableFactor) -> bool {
92 match factor {
93 TableFactor::UNNEST { .. } => true,
94 TableFactor::Table {
95 name,
96 args: Some(_),
97 ..
98 }
99 | TableFactor::Function { name, .. } => {
100 name.0.len() == 1 && name.to_string().eq_ignore_ascii_case("unnest")
101 }
102 _ => false,
103 }
104}
105
106fn has_implicit_multi_source(select: &Select) -> bool {
107 select
108 .from
109 .iter()
110 .filter(|from| !is_inline_unnest(&from.relation))
111 .count()
112 > 1
113}
114
115#[derive(Debug)]
117#[allow(clippy::large_enum_variant)]
118pub enum StreamingPlan {
119 RegisterSource(SourceInfo),
121
122 RegisterSink(SinkInfo),
124
125 Query(QueryPlan),
127
128 Standard(Box<Statement>),
130
131 RegisterLookupTable(LookupTableInfo),
133
134 DropLookupTable {
136 name: String,
138 },
139}
140
141#[derive(Debug)]
143pub struct QueryPlan {
144 pub name: Option<String>,
146 pub window_config: Option<WindowOperatorConfig>,
148 pub join_config: Option<Vec<JoinOperatorConfig>>,
150 pub order_config: Option<OrderOperatorConfig>,
152 pub analytic_config: Option<AnalyticWindowConfig>,
154 pub having_config: Option<HavingFilterConfig>,
156 pub frame_config: Option<WindowFrameConfig>,
158 pub emit_clause: Option<EmitClause>,
160 pub statement: Box<Statement>,
162}
163
164#[derive(Debug, Clone, PartialEq, Eq)]
169pub struct StateBackedJoinAdmission {
170 left_table: String,
171 right_table: String,
172 left_keys: Vec<String>,
173 right_keys: Vec<String>,
174 join_type: JoinType,
175}
176
177impl StateBackedJoinAdmission {
178 pub fn try_new(
183 left_table: impl Into<String>,
184 right_table: impl Into<String>,
185 left_keys: Vec<String>,
186 right_keys: Vec<String>,
187 left_outer: bool,
188 ) -> Result<Self, String> {
189 let left_table = left_table.into();
190 let right_table = right_table.into();
191 if left_table.is_empty() || right_table.is_empty() {
192 return Err("state-backed join relations cannot be empty".into());
193 }
194 if left_keys.is_empty()
195 || left_keys.len() != right_keys.len()
196 || left_keys.iter().chain(&right_keys).any(String::is_empty)
197 {
198 return Err("state-backed join keys must be non-empty with matching arity".into());
199 }
200 Ok(Self {
201 left_table,
202 right_table,
203 left_keys,
204 right_keys,
205 join_type: if left_outer {
206 JoinType::Left
207 } else {
208 JoinType::Inner
209 },
210 })
211 }
212
213 fn matches(&self, step: &JoinAnalysis) -> bool {
214 let mut left_keys = Vec::with_capacity(1 + step.additional_key_columns.len());
215 let mut right_keys = Vec::with_capacity(1 + step.additional_key_columns.len());
216 left_keys.push(step.left_key_column.as_str());
217 right_keys.push(step.right_key_column.as_str());
218 for (left, right) in &step.additional_key_columns {
219 left_keys.push(left);
220 right_keys.push(right);
221 }
222 self.left_table == step.left_table
223 && self.right_table == step.right_table
224 && self.join_type == step.join_type
225 && self.left_keys.iter().map(String::as_str).eq(left_keys)
226 && self.right_keys.iter().map(String::as_str).eq(right_keys)
227 && !step.is_asof_join
228 && !step.is_temporal_join
229 && step.time_bound.is_none()
230 }
231}
232
233impl StreamingPlanner {
234 #[must_use]
236 pub fn new() -> Self {
237 Self {
238 sources: HashMap::new(),
239 sinks: HashMap::new(),
240 lookup_tables: HashMap::new(),
241 windowed_views: std::collections::HashSet::new(),
242 }
243 }
244
245 pub fn plan(&mut self, statement: &StreamingStatement) -> Result<StreamingPlan, PlanningError> {
251 self.plan_internal(statement, None)
252 }
253
254 pub fn plan_state_backed_join(
263 &mut self,
264 statement: &StreamingStatement,
265 admission: &StateBackedJoinAdmission,
266 ) -> Result<StreamingPlan, PlanningError> {
267 if !matches!(
268 statement,
269 StreamingStatement::CreateContinuousQuery { .. }
270 | StreamingStatement::CreateStream { .. }
271 ) {
272 return Err(PlanningError::InvalidQuery(
273 "state-backed join admission is valid only for a named streaming query".into(),
274 ));
275 }
276 self.plan_internal(statement, Some(admission))
277 }
278
279 fn plan_internal(
280 &mut self,
281 statement: &StreamingStatement,
282 state_backed_join: Option<&StateBackedJoinAdmission>,
283 ) -> Result<StreamingPlan, PlanningError> {
284 match statement {
285 StreamingStatement::CreateSource(source) => self.plan_create_source(source),
286 StreamingStatement::CreateSink(sink) => self.plan_create_sink(sink),
287 StreamingStatement::CreateContinuousQuery {
288 name,
289 query,
290 emit_clause,
291 ..
292 }
293 | StreamingStatement::CreateStream {
294 name,
295 query,
296 emit_clause,
297 ..
298 } => self.plan_continuous_query(name, query, emit_clause.as_ref(), state_backed_join),
299 StreamingStatement::Standard(stmt) => self.plan_standard_statement(stmt),
300 StreamingStatement::CreateLookupTable(lt) => self.plan_create_lookup_table(lt),
301 StreamingStatement::DropLookupTable { name, if_exists } => {
302 self.plan_drop_lookup_table(name, *if_exists)
303 }
304 StreamingStatement::DropSource { .. }
305 | StreamingStatement::DropSink { .. }
306 | StreamingStatement::DropStream { .. }
307 | StreamingStatement::DropMaterializedView { .. }
308 | StreamingStatement::Show(_)
309 | StreamingStatement::Describe { .. }
310 | StreamingStatement::Explain { .. }
311 | StreamingStatement::CreateMaterializedView { .. }
312 | StreamingStatement::InsertInto { .. }
313 | StreamingStatement::AlterSource { .. }
314 | StreamingStatement::Checkpoint
315 | StreamingStatement::RestoreCheckpoint { .. }
316 | StreamingStatement::Subscribe(_)
317 | StreamingStatement::DeclareCursorForSubscribe { .. } => {
318 Err(PlanningError::UnsupportedSql(format!(
321 "Statement type {:?} is handled by the database layer, not the planner",
322 std::mem::discriminant(statement)
323 )))
324 }
325 }
326 }
327
328 pub fn unregister_query(&mut self, name: &str) {
331 self.windowed_views.remove(name);
332 }
333
334 #[must_use]
336 pub fn has_query(&self, name: &str) -> bool {
337 self.windowed_views.contains(name)
338 }
339
340 pub fn unregister_source(&mut self, name: &str) {
342 self.sources.remove(name);
343 }
344
345 pub fn unregister_sink(&mut self, name: &str) {
347 self.sinks.remove(name);
348 }
349
350 pub fn unregister_lookup_table(&mut self, name: &str) {
352 self.lookup_tables.remove(name);
353 }
354
355 fn plan_create_source(
357 &mut self,
358 source: &CreateSourceStatement,
359 ) -> Result<StreamingPlan, PlanningError> {
360 let name = object_name_to_string(&source.name);
361
362 if !source.or_replace && !source.if_not_exists && self.sources.contains_key(&name) {
364 return Err(PlanningError::InvalidQuery(format!(
365 "Source '{}' already exists",
366 name
367 )));
368 }
369
370 let watermark_column = source.watermark.as_ref().map(|w| w.column.value.clone());
372
373 let info = SourceInfo {
374 name: name.clone(),
375 watermark_column,
376 options: source.with_options.clone(),
377 };
378
379 self.sources.insert(name, info.clone());
381
382 Ok(StreamingPlan::RegisterSource(info))
383 }
384
385 fn plan_create_sink(
387 &mut self,
388 sink: &CreateSinkStatement,
389 ) -> Result<StreamingPlan, PlanningError> {
390 let name = object_name_to_string(&sink.name);
391
392 if !sink.or_replace && !sink.if_not_exists && self.sinks.contains_key(&name) {
394 return Err(PlanningError::InvalidQuery(format!(
395 "Sink '{}' already exists",
396 name
397 )));
398 }
399
400 let from = match &sink.from {
402 SinkFrom::Table(table) => object_name_to_string(table),
403 SinkFrom::Query(_) => format!("{}_query", name),
404 };
405
406 let info = SinkInfo {
407 name: name.clone(),
408 from,
409 options: sink.with_options.clone(),
410 };
411
412 self.sinks.insert(name, info.clone());
414
415 Ok(StreamingPlan::RegisterSink(info))
416 }
417
418 fn plan_continuous_query(
420 &mut self,
421 name: &ObjectName,
422 query: &StreamingStatement,
423 emit_clause: Option<&EmitClause>,
424 state_backed_join: Option<&StateBackedJoinAdmission>,
425 ) -> Result<StreamingPlan, PlanningError> {
426 let stmt = match query {
428 StreamingStatement::Standard(stmt) => stmt.as_ref().clone(),
429 _ => {
430 return Err(PlanningError::InvalidQuery(
431 "Continuous query must contain a SELECT statement".to_string(),
432 ))
433 }
434 };
435
436 let query_plan = self.analyze_query(&stmt, emit_clause, state_backed_join)?;
438
439 let view_name = object_name_to_string(name);
442 if query_plan.window_config.is_some() {
443 self.windowed_views.insert(view_name);
444 } else {
445 self.windowed_views.remove(&view_name);
446 }
447
448 Ok(StreamingPlan::Query(QueryPlan {
449 name: Some(object_name_to_string(name)),
450 window_config: query_plan.window_config,
451 join_config: query_plan.join_config,
452 order_config: query_plan.order_config,
453 analytic_config: query_plan.analytic_config,
454 having_config: query_plan.having_config,
455 frame_config: query_plan.frame_config,
456 emit_clause: emit_clause.cloned(),
457 statement: Box::new(stmt),
458 }))
459 }
460
461 #[allow(clippy::unused_self)] fn plan_standard_statement(&self, stmt: &Statement) -> Result<StreamingPlan, PlanningError> {
464 if let Statement::Query(query) = stmt {
466 if let SetExpr::Select(select) = query.body.as_ref() {
467 if has_implicit_multi_source(select) {
468 return Err(PlanningError::InvalidQuery(
469 "implicit multi-source joins are unsupported; use one bounded INNER JOIN"
470 .to_string(),
471 ));
472 }
473 let window_function = Self::extract_window_from_select(select);
475
476 let join_analysis = analyze_joins(select).map_err(|e| {
478 PlanningError::InvalidQuery(format!("Join analysis failed: {e}"))
479 })?;
480
481 if let Some(ref multi) = join_analysis {
482 validate_streaming_joins(multi, &self.lookup_tables, None)?;
483 }
484
485 let order_analysis = analyze_order_by(stmt);
487 let order_config = OrderOperatorConfig::from_analysis(&order_analysis)
488 .map_err(PlanningError::InvalidQuery)?;
489
490 let analytic_analysis = analyze_analytic_functions(stmt);
492 let analytic_config =
493 analytic_analysis.map(|a| AnalyticWindowConfig::from_analysis(&a));
494
495 let agg_analysis = analyze_aggregates(stmt);
497 let having_config = agg_analysis.having_expr.map(HavingFilterConfig::new);
498
499 let frame_analysis = analyze_window_frames(stmt);
501 let frame_config = frame_analysis
502 .as_ref()
503 .map(WindowFrameConfig::from_analysis);
504
505 if let Some(fa) = &frame_analysis {
507 for f in &fa.functions {
508 if matches!(f.end_bound, FrameBound::UnboundedFollowing) {
509 return Err(PlanningError::InvalidQuery(
510 "UNBOUNDED FOLLOWING is not supported in streaming window frames"
511 .to_string(),
512 ));
513 }
514 }
515 }
516
517 let has_streaming_features = window_function.is_some()
518 || join_analysis.is_some()
519 || order_config.is_some()
520 || analytic_config.is_some()
521 || having_config.is_some()
522 || frame_config.is_some();
523
524 if has_streaming_features {
525 let window_config = match window_function {
526 Some(w) => Some(
527 WindowOperatorConfig::from_window_function(&w)
528 .map_err(|e| PlanningError::InvalidQuery(e.to_string()))?,
529 ),
530 None => None,
531 };
532
533 let join_config =
534 join_analysis.map(|m| JoinOperatorConfig::from_multi_analysis(&m));
535
536 return Ok(StreamingPlan::Query(QueryPlan {
537 name: None,
538 window_config,
539 join_config,
540 order_config,
541 analytic_config,
542 having_config,
543 frame_config,
544 emit_clause: None,
545 statement: Box::new(stmt.clone()),
546 }));
547 }
548 }
549 }
550
551 Ok(StreamingPlan::Standard(Box::new(stmt.clone())))
553 }
554
555 fn analyze_query(
557 &self,
558 stmt: &Statement,
559 emit_clause: Option<&EmitClause>,
560 state_backed_join: Option<&StateBackedJoinAdmission>,
561 ) -> Result<QueryAnalysis, PlanningError> {
562 let mut analysis = QueryAnalysis::default();
563
564 if let Statement::Query(query) = stmt {
565 if let SetExpr::Select(select) = query.body.as_ref() {
566 if has_implicit_multi_source(select) {
567 return Err(PlanningError::InvalidQuery(
568 "implicit multi-source joins are unsupported; use one bounded INNER JOIN"
569 .to_string(),
570 ));
571 }
572 if let Some(window) = Self::extract_window_from_select(select) {
574 let mut config = WindowOperatorConfig::from_window_function(&window)
575 .map_err(|e| PlanningError::InvalidQuery(e.to_string()))?;
576
577 if let Some(emit) = emit_clause {
579 config = config
580 .with_emit_clause(emit)
581 .map_err(|e| PlanningError::InvalidQuery(e.to_string()))?;
582 }
583
584 analysis.window_config = Some(config);
585 }
586
587 if let Some(multi) = analyze_joins(select).map_err(|e| {
589 PlanningError::InvalidQuery(format!("Join analysis failed: {e}"))
590 })? {
591 validate_streaming_joins(&multi, &self.lookup_tables, state_backed_join)?;
592 analysis.join_config = Some(JoinOperatorConfig::from_multi_analysis(&multi));
593 }
594 }
595 }
596
597 let order_analysis = analyze_order_by(stmt);
599 analysis.order_config = OrderOperatorConfig::from_analysis(&order_analysis)
600 .map_err(PlanningError::InvalidQuery)?;
601
602 if let Some(analytic) = analyze_analytic_functions(stmt) {
604 analysis.analytic_config = Some(AnalyticWindowConfig::from_analysis(&analytic));
605 }
606
607 let agg_analysis = analyze_aggregates(stmt);
609 analysis.having_config = agg_analysis.having_expr.map(HavingFilterConfig::new);
610
611 if let Some(frame_analysis) = analyze_window_frames(stmt) {
613 for f in &frame_analysis.functions {
615 if matches!(f.end_bound, FrameBound::UnboundedFollowing) {
616 return Err(PlanningError::InvalidQuery(
617 "UNBOUNDED FOLLOWING is not supported in streaming window frames"
618 .to_string(),
619 ));
620 }
621 }
622 analysis.frame_config = Some(WindowFrameConfig::from_analysis(&frame_analysis));
623 }
624
625 Ok(analysis)
626 }
627
628 fn extract_window_from_select(select: &sqlparser::ast::Select) -> Option<WindowFunction> {
630 use sqlparser::ast::GroupByExpr;
632 match &select.group_by {
633 GroupByExpr::Expressions(exprs, _modifiers) => {
634 for group_by_expr in exprs {
635 if let Ok(Some(window)) = WindowRewriter::extract_window_function(group_by_expr)
636 {
637 return Some(window);
638 }
639 }
640 }
641 GroupByExpr::All(_) => {}
642 }
643 None
644 }
645
646 fn plan_create_lookup_table(
648 &mut self,
649 lt: &CreateLookupTableStatement,
650 ) -> Result<StreamingPlan, PlanningError> {
651 let name = object_name_to_string(<.name);
652
653 if !lt.or_replace && !lt.if_not_exists && self.lookup_tables.contains_key(&name) {
654 return Err(PlanningError::InvalidQuery(format!(
655 "Lookup table '{}' already exists",
656 name
657 )));
658 }
659
660 let columns: Vec<(String, String)> = lt
661 .columns
662 .iter()
663 .map(|c| (c.name.value.clone(), c.data_type.to_string()))
664 .collect();
665
666 let properties = validate_properties(<.with_options).map_err(|e| {
667 PlanningError::InvalidQuery(format!("Invalid lookup table properties: {e}"))
668 })?;
669
670 let arrow_fields: Vec<Field> = lt
672 .columns
673 .iter()
674 .map(|c| {
675 let dt = crate::translator::streaming_ddl::sql_type_to_arrow(&c.data_type)
676 .map_err(|e| PlanningError::InvalidQuery(e.to_string()))?;
677 let nullable = !c
678 .options
679 .iter()
680 .any(|opt| matches!(opt.option, sqlparser::ast::ColumnOption::NotNull));
681 Ok(Field::new(&c.name.value, dt, nullable))
682 })
683 .collect::<Result<_, PlanningError>>()?;
684 let arrow_schema = Arc::new(Schema::new(arrow_fields));
685
686 let info = LookupTableInfo {
687 name: name.clone(),
688 columns,
689 primary_key: lt.primary_key.clone(),
690 properties,
691 arrow_schema,
692 raw_options: lt.with_options.clone(),
693 };
694
695 self.lookup_tables.insert(name, info.clone());
696
697 Ok(StreamingPlan::RegisterLookupTable(info))
698 }
699
700 fn plan_drop_lookup_table(
702 &mut self,
703 name: &ObjectName,
704 if_exists: bool,
705 ) -> Result<StreamingPlan, PlanningError> {
706 let name_str = object_name_to_string(name);
707
708 if !if_exists && !self.lookup_tables.contains_key(&name_str) {
709 return Err(PlanningError::InvalidQuery(format!(
710 "Lookup table '{}' does not exist",
711 name_str
712 )));
713 }
714
715 self.lookup_tables.remove(&name_str);
716
717 Ok(StreamingPlan::DropLookupTable { name: name_str })
718 }
719
720 #[must_use]
722 pub fn get_source(&self, name: &str) -> Option<&SourceInfo> {
723 self.sources.get(name)
724 }
725
726 #[must_use]
728 pub fn get_sink(&self, name: &str) -> Option<&SinkInfo> {
729 self.sinks.get(name)
730 }
731
732 #[must_use]
734 pub fn list_sources(&self) -> Vec<&SourceInfo> {
735 self.sources.values().collect()
736 }
737
738 #[must_use]
740 pub fn list_sinks(&self) -> Vec<&SinkInfo> {
741 self.sinks.values().collect()
742 }
743
744 #[must_use]
746 pub fn get_lookup_table(&self, name: &str) -> Option<&LookupTableInfo> {
747 self.lookup_tables.get(name)
748 }
749
750 #[must_use]
752 pub fn list_lookup_tables(&self) -> Vec<&LookupTableInfo> {
753 self.lookup_tables.values().collect()
754 }
755
756 #[must_use]
758 pub fn lookup_tables_cloned(&self) -> HashMap<String, LookupTableInfo> {
759 self.lookup_tables.clone()
760 }
761
762 #[allow(clippy::unused_self)] pub async fn to_logical_plan(
773 &self,
774 plan: &QueryPlan,
775 ctx: &SessionContext,
776 ) -> Result<LogicalPlan, PlanningError> {
777 let sql = plan.statement.to_string();
781 ctx.state()
782 .create_logical_plan(&sql)
783 .await
784 .map_err(PlanningError::DataFusion)
785 }
786}
787
788impl Default for StreamingPlanner {
789 fn default() -> Self {
790 Self::new()
791 }
792}
793
794#[derive(Debug, Default)]
796#[allow(clippy::struct_field_names)]
797struct QueryAnalysis {
798 window_config: Option<WindowOperatorConfig>,
799 join_config: Option<Vec<JoinOperatorConfig>>,
800 order_config: Option<OrderOperatorConfig>,
801 analytic_config: Option<AnalyticWindowConfig>,
802 having_config: Option<HavingFilterConfig>,
803 frame_config: Option<WindowFrameConfig>,
804}
805
806fn object_name_to_string(name: &ObjectName) -> String {
808 match name.0.as_slice() {
809 [sqlparser::ast::ObjectNamePart::Identifier(ident)] => ident.value.clone(),
810 _ => name.to_string(),
811 }
812}
813
814fn validate_streaming_joins(
816 multi: &MultiJoinAnalysis,
817 lookup_tables: &HashMap<String, LookupTableInfo>,
818 state_backed_join: Option<&StateBackedJoinAdmission>,
819) -> Result<(), PlanningError> {
820 if multi.joins.len() != 1 {
821 return Err(PlanningError::InvalidQuery(
822 "multi-way streaming joins require explicitly named two-way stages".to_string(),
823 ));
824 }
825 for step in &multi.joins {
826 if step.time_bound.is_some_and(|bound| bound.is_zero()) {
827 return Err(PlanningError::InvalidQuery(
828 "streaming interval joins require a positive finite time bound".to_string(),
829 ));
830 }
831 if step.time_bound.is_some()
832 && !step.is_asof_join
833 && !step.is_temporal_join
834 && !matches!(step.join_type, JoinType::Inner)
835 {
836 return Err(PlanningError::InvalidQuery(format!(
837 "streaming interval joins support only INNER joins; {:?} requires durable per-row matched metadata",
838 step.join_type,
839 )));
840 }
841 if step.is_bounded() {
842 continue;
843 }
844 let left_lookup = lookup_tables.contains_key(&step.left_table);
845 let right_lookup = lookup_tables.contains_key(&step.right_table);
846 if !left_lookup && !right_lookup {
847 if state_backed_join.is_some_and(|admission| admission.matches(step)) {
848 continue;
849 }
850 return Err(PlanningError::InvalidQuery(format!(
851 "unbounded join between streaming sources '{}' and '{}'; \
852 add a temporal predicate or use a lookup table",
853 step.left_table, step.right_table,
854 )));
855 }
856 }
857 Ok(())
858}
859
860#[derive(Debug, thiserror::Error)]
862pub enum PlanningError {
863 UnsupportedSql(String),
865
866 InvalidQuery(String),
868
869 SourceNotFound(String),
871
872 SinkNotFound(String),
874
875 DataFusion(#[from] datafusion_common::DataFusionError),
877}
878
879impl std::fmt::Display for PlanningError {
880 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
881 match self {
882 Self::UnsupportedSql(msg) => write!(f, "Unsupported SQL: {msg}"),
883 Self::InvalidQuery(msg) => write!(f, "Invalid query: {msg}"),
884 Self::SourceNotFound(name) => write!(f, "Source not found: {name}"),
885 Self::SinkNotFound(name) => write!(f, "Sink not found: {name}"),
886 Self::DataFusion(e) => {
887 let translated = crate::error::translate_datafusion_error(&e.to_string());
888 write!(f, "{translated}")
889 }
890 }
891 }
892}
893
894#[cfg(test)]
895mod tests {
896 use super::*;
897 use crate::parser::StreamingParser;
898
899 #[test]
900 fn test_plan_create_source() {
901 let mut planner = StreamingPlanner::new();
902 let statements =
903 StreamingParser::parse_sql("CREATE SOURCE events (id INT, name VARCHAR)").unwrap();
904
905 let plan = planner.plan(&statements[0]).unwrap();
906 match plan {
907 StreamingPlan::RegisterSource(info) => {
908 assert_eq!(info.name, "events");
909 }
910 _ => panic!("Expected RegisterSource plan"),
911 }
912 }
913
914 #[test]
915 fn test_plan_create_sink() {
916 let mut planner = StreamingPlanner::new();
917 let statements = StreamingParser::parse_sql("CREATE SINK output FROM events").unwrap();
918
919 let plan = planner.plan(&statements[0]).unwrap();
920 match plan {
921 StreamingPlan::RegisterSink(info) => {
922 assert_eq!(info.name, "output");
923 assert_eq!(info.from, "events");
924 }
925 _ => panic!("Expected RegisterSink plan"),
926 }
927 }
928
929 #[test]
930 fn test_plan_duplicate_source() {
931 let mut planner = StreamingPlanner::new();
932
933 let statements =
935 StreamingParser::parse_sql("CREATE SOURCE events (id INT, name VARCHAR)").unwrap();
936 planner.plan(&statements[0]).unwrap();
937
938 let result = planner.plan(&statements[0]);
940 assert!(result.is_err());
941 }
942
943 #[test]
944 fn test_plan_source_if_not_exists() {
945 let mut planner = StreamingPlanner::new();
946
947 let statements =
949 StreamingParser::parse_sql("CREATE SOURCE events (id INT, name VARCHAR)").unwrap();
950 planner.plan(&statements[0]).unwrap();
951
952 let statements =
954 StreamingParser::parse_sql("CREATE SOURCE IF NOT EXISTS events (id INT, name VARCHAR)")
955 .unwrap();
956 let result = planner.plan(&statements[0]);
957 assert!(result.is_ok());
958 }
959
960 #[test]
961 fn test_plan_source_or_replace() {
962 let mut planner = StreamingPlanner::new();
963
964 let statements =
966 StreamingParser::parse_sql("CREATE SOURCE events (id INT, name VARCHAR)").unwrap();
967 planner.plan(&statements[0]).unwrap();
968
969 let statements =
971 StreamingParser::parse_sql("CREATE OR REPLACE SOURCE events (id INT, name VARCHAR)")
972 .unwrap();
973 let result = planner.plan(&statements[0]);
974 assert!(result.is_ok());
975 }
976
977 #[test]
978 fn test_plan_source_with_watermark() {
979 let mut planner = StreamingPlanner::new();
980 let statements = StreamingParser::parse_sql(
981 "CREATE SOURCE events (
982 id INT,
983 ts TIMESTAMP,
984 WATERMARK FOR ts AS ts - INTERVAL '5' SECOND
985 )",
986 )
987 .unwrap();
988
989 let plan = planner.plan(&statements[0]).unwrap();
990 match plan {
991 StreamingPlan::RegisterSource(info) => {
992 assert_eq!(info.name, "events");
993 assert_eq!(info.watermark_column, Some("ts".to_string()));
994 }
995 _ => panic!("Expected RegisterSource plan"),
996 }
997 }
998
999 #[test]
1000 fn test_plan_standard_select() {
1001 let mut planner = StreamingPlanner::new();
1002 let statements = StreamingParser::parse_sql("SELECT * FROM events").unwrap();
1003
1004 let plan = planner.plan(&statements[0]).unwrap();
1005 match plan {
1006 StreamingPlan::Standard(_) => {}
1007 _ => panic!("Expected Standard plan for simple SELECT"),
1008 }
1009 }
1010
1011 #[test]
1012 fn test_list_sources_and_sinks() {
1013 let mut planner = StreamingPlanner::new();
1014
1015 let s1 = StreamingParser::parse_sql("CREATE SOURCE src1 (id INT)").unwrap();
1017 let s2 = StreamingParser::parse_sql("CREATE SOURCE src2 (id INT)").unwrap();
1018 planner.plan(&s1[0]).unwrap();
1019 planner.plan(&s2[0]).unwrap();
1020
1021 let k1 = StreamingParser::parse_sql("CREATE SINK sink1 FROM src1").unwrap();
1023 planner.plan(&k1[0]).unwrap();
1024
1025 assert_eq!(planner.list_sources().len(), 2);
1026 assert_eq!(planner.list_sinks().len(), 1);
1027 assert!(planner.get_source("src1").is_some());
1028 assert!(planner.get_sink("sink1").is_some());
1029 }
1030
1031 #[test]
1032 fn test_plan_query_with_window() {
1033 let mut planner = StreamingPlanner::new();
1034 let statements = StreamingParser::parse_sql(
1035 "SELECT COUNT(*) FROM events GROUP BY TUMBLE(event_time, INTERVAL '5' MINUTE)",
1036 )
1037 .unwrap();
1038
1039 let plan = planner.plan(&statements[0]).unwrap();
1040 match plan {
1041 StreamingPlan::Query(query_plan) => {
1042 assert!(query_plan.window_config.is_some());
1043 let config = query_plan.window_config.unwrap();
1044 assert_eq!(config.time_column, "event_time");
1045 assert_eq!(config.size.as_secs(), 300);
1046 }
1047 _ => panic!("Expected Query plan"),
1048 }
1049 }
1050
1051 #[test]
1052 fn test_plan_query_with_join() {
1053 let mut planner = StreamingPlanner::new();
1054 let statements = StreamingParser::parse_sql(
1055 "SELECT * FROM orders o JOIN payments p ON o.order_id = p.order_id \
1056 AND p.ts BETWEEN o.ts AND o.ts + INTERVAL '1' HOUR",
1057 )
1058 .unwrap();
1059
1060 let plan = planner.plan(&statements[0]).unwrap();
1061 match plan {
1062 StreamingPlan::Query(query_plan) => {
1063 assert!(query_plan.join_config.is_some());
1064 let configs = query_plan.join_config.unwrap();
1065 assert_eq!(configs.len(), 1);
1066 assert_eq!(configs[0].left_key(), "order_id");
1067 assert_eq!(configs[0].right_key(), "order_id");
1068 }
1069 _ => panic!("Expected Query plan"),
1070 }
1071 }
1072
1073 #[test]
1074 fn test_plan_rejects_unbounded_streaming_join() {
1075 let mut planner = StreamingPlanner::new();
1076 let statements = StreamingParser::parse_sql(
1077 "CREATE STREAM joined AS SELECT * FROM orders o JOIN payments p \
1078 ON o.order_id = p.order_id",
1079 )
1080 .unwrap();
1081 let err = planner.plan(&statements[0]).unwrap_err();
1082 let msg = format!("{err}");
1083 assert!(msg.contains("unbounded join"), "got: {msg}");
1084 assert!(msg.contains("lookup table"), "got: {msg}");
1085 }
1086
1087 #[test]
1088 fn test_plan_rejects_unbounded_standard_select_join() {
1089 let statements = StreamingParser::parse_sql(
1090 "SELECT * FROM orders o JOIN payments p ON o.order_id = p.order_id",
1091 )
1092 .unwrap();
1093 assert!(matches!(
1094 statements.first(),
1095 Some(StreamingStatement::Standard(_))
1096 ));
1097
1098 let error = StreamingPlanner::new().plan(&statements[0]).unwrap_err();
1099 assert!(error.to_string().contains("unbounded join"));
1100 }
1101
1102 #[test]
1103 fn exact_state_backed_certificate_allows_only_its_join() {
1104 let statements = StreamingParser::parse_sql(
1105 "CREATE STREAM joined AS SELECT * FROM orders o JOIN payments p \
1106 ON o.order_id = p.order_id",
1107 )
1108 .unwrap();
1109 let admission = StateBackedJoinAdmission::try_new(
1110 "orders",
1111 "payments",
1112 vec!["order_id".into()],
1113 vec!["order_id".into()],
1114 false,
1115 )
1116 .unwrap();
1117 StreamingPlanner::new()
1118 .plan_state_backed_join(&statements[0], &admission)
1119 .unwrap();
1120
1121 let wrong_key = StateBackedJoinAdmission::try_new(
1122 "orders",
1123 "payments",
1124 vec!["tenant_id".into()],
1125 vec!["tenant_id".into()],
1126 false,
1127 )
1128 .unwrap();
1129 let error = StreamingPlanner::new()
1130 .plan_state_backed_join(&statements[0], &wrong_key)
1131 .unwrap_err();
1132 assert!(error.to_string().contains("unbounded join"));
1133 }
1134
1135 #[test]
1136 fn composite_state_backed_certificate_binds_ordered_key_mapping() {
1137 let statements = StreamingParser::parse_sql(
1138 "CREATE STREAM joined AS SELECT * FROM orders o JOIN payments p \
1139 ON o.tenant_id = p.account_id AND o.order_id = p.payment_order_id",
1140 )
1141 .unwrap();
1142 let exact = StateBackedJoinAdmission::try_new(
1143 "orders",
1144 "payments",
1145 vec!["tenant_id".into(), "order_id".into()],
1146 vec!["account_id".into(), "payment_order_id".into()],
1147 false,
1148 )
1149 .unwrap();
1150 StreamingPlanner::new()
1151 .plan_state_backed_join(&statements[0], &exact)
1152 .unwrap();
1153
1154 for (left_keys, right_keys) in [
1155 (
1156 vec!["order_id".into(), "tenant_id".into()],
1157 vec!["payment_order_id".into(), "account_id".into()],
1158 ),
1159 (
1160 vec!["tenant_id".into(), "order_id".into()],
1161 vec!["account_id".into(), "order_id".into()],
1162 ),
1163 ] {
1164 let mismatch = StateBackedJoinAdmission::try_new(
1165 "orders", "payments", left_keys, right_keys, false,
1166 )
1167 .unwrap();
1168 let error = StreamingPlanner::new()
1169 .plan_state_backed_join(&statements[0], &mismatch)
1170 .unwrap_err();
1171 assert!(error.to_string().contains("unbounded join"));
1172 }
1173 }
1174
1175 #[test]
1176 fn state_backed_certificate_binds_inner_vs_left() {
1177 let statements = StreamingParser::parse_sql(
1178 "CREATE STREAM joined AS SELECT * FROM orders o LEFT JOIN payments p \
1179 ON o.order_id = p.order_id",
1180 )
1181 .unwrap();
1182 let left = StateBackedJoinAdmission::try_new(
1183 "orders",
1184 "payments",
1185 vec!["order_id".into()],
1186 vec!["order_id".into()],
1187 true,
1188 )
1189 .unwrap();
1190 StreamingPlanner::new()
1191 .plan_state_backed_join(&statements[0], &left)
1192 .unwrap();
1193
1194 let inner = StateBackedJoinAdmission::try_new(
1195 "orders",
1196 "payments",
1197 vec!["order_id".into()],
1198 vec!["order_id".into()],
1199 false,
1200 )
1201 .unwrap();
1202 assert!(StreamingPlanner::new()
1203 .plan_state_backed_join(&statements[0], &inner)
1204 .is_err());
1205 }
1206
1207 #[test]
1208 fn state_backed_certificate_does_not_relax_right_or_multiway_joins() {
1209 let admission = StateBackedJoinAdmission::try_new(
1210 "orders",
1211 "payments",
1212 vec!["order_id".into()],
1213 vec!["order_id".into()],
1214 false,
1215 )
1216 .unwrap();
1217 for sql in [
1218 "CREATE STREAM joined AS SELECT * FROM orders o RIGHT JOIN payments p \
1219 ON o.order_id = p.order_id",
1220 "CREATE STREAM joined AS SELECT * FROM orders o JOIN payments p \
1221 ON o.order_id = p.order_id JOIN refunds r ON p.order_id = r.order_id",
1222 ] {
1223 let statements = StreamingParser::parse_sql(sql).unwrap();
1224 assert!(StreamingPlanner::new()
1225 .plan_state_backed_join(&statements[0], &admission)
1226 .is_err());
1227 }
1228 }
1229
1230 #[test]
1231 fn test_plan_rejects_implicit_cross_join() {
1232 let statements = StreamingParser::parse_sql("SELECT * FROM orders, payments").unwrap();
1233 let error = StreamingPlanner::new().plan(&statements[0]).unwrap_err();
1234 assert!(error.to_string().contains("implicit multi-source"));
1235 }
1236
1237 #[test]
1238 fn test_plan_allows_unnest_after_one_stream_source() {
1239 let statements = StreamingParser::parse_sql(
1240 "SELECT event_id, tag FROM events, UNNEST(make_array('a', 'b')) AS tags(tag)",
1241 )
1242 .unwrap();
1243 StreamingPlanner::new().plan(&statements[0]).unwrap();
1244 }
1245
1246 #[test]
1247 fn test_plan_rejects_unbounded_join_between_windowed_views() {
1248 let mut planner = StreamingPlanner::new();
1249 for sql in [
1251 "CREATE STREAM price_1m AS SELECT TUMBLE(ts, INTERVAL '1' MINUTE) AS bucket, \
1252 AVG(p) AS price FROM trades GROUP BY TUMBLE(ts, INTERVAL '1' MINUTE)",
1253 "CREATE STREAM sent_1m AS SELECT TUMBLE(ts, INTERVAL '1' MINUTE) AS bucket, \
1254 AVG(s) AS ms FROM posts GROUP BY TUMBLE(ts, INTERVAL '1' MINUTE)",
1255 ] {
1256 let st = StreamingParser::parse_sql(sql).unwrap();
1257 planner.plan(&st[0]).unwrap();
1258 }
1259 let st = StreamingParser::parse_sql(
1263 "CREATE STREAM joined AS SELECT a.bucket, a.price, b.ms \
1264 FROM price_1m a JOIN sent_1m b ON a.bucket = b.bucket",
1265 )
1266 .unwrap();
1267 let error = planner.plan(&st[0]).unwrap_err();
1268 assert!(error.to_string().contains("unbounded join"));
1269 }
1270
1271 #[test]
1272 fn test_plan_rejects_non_inner_interval_joins() {
1273 for join_type in [
1274 "LEFT",
1275 "RIGHT",
1276 "FULL",
1277 "LEFT SEMI",
1278 "LEFT ANTI",
1279 "RIGHT SEMI",
1280 "RIGHT ANTI",
1281 ] {
1282 let sql = format!(
1283 "SELECT * FROM orders o {join_type} JOIN payments p \
1284 ON o.order_id = p.order_id \
1285 AND p.ts BETWEEN o.ts AND o.ts + INTERVAL '1' HOUR"
1286 );
1287 let statements = StreamingParser::parse_sql(&sql).unwrap();
1288 let error = StreamingPlanner::new().plan(&statements[0]).unwrap_err();
1289 let message = error.to_string();
1290 assert!(
1291 message.contains("only INNER joins"),
1292 "{join_type} interval join was not rejected: {message}"
1293 );
1294 }
1295 }
1296
1297 #[test]
1298 fn test_plan_rejects_zero_interval_bound() {
1299 let statements = StreamingParser::parse_sql(
1300 "SELECT * FROM orders o JOIN payments p ON o.order_id = p.order_id \
1301 AND p.ts BETWEEN o.ts AND o.ts + INTERVAL '0' SECOND",
1302 )
1303 .unwrap();
1304 let error = StreamingPlanner::new().plan(&statements[0]).unwrap_err();
1305 assert!(error.to_string().contains("positive finite time bound"));
1306 }
1307
1308 #[test]
1309 fn test_plan_query_with_lag() {
1310 let mut planner = StreamingPlanner::new();
1311 let statements = StreamingParser::parse_sql(
1312 "SELECT price, LAG(price) OVER (PARTITION BY symbol ORDER BY ts) AS prev FROM trades",
1313 )
1314 .unwrap();
1315
1316 let plan = planner.plan(&statements[0]).unwrap();
1317 match plan {
1318 StreamingPlan::Query(query_plan) => {
1319 assert!(query_plan.analytic_config.is_some());
1320 let config = query_plan.analytic_config.unwrap();
1321 assert_eq!(config.functions.len(), 1);
1322 assert_eq!(config.partition_columns, vec!["symbol".to_string()]);
1323 }
1324 _ => panic!("Expected Query plan with analytic config"),
1325 }
1326 }
1327
1328 #[test]
1329 fn test_plan_query_with_having() {
1330 let mut planner = StreamingPlanner::new();
1331 let statements = StreamingParser::parse_sql(
1332 "SELECT symbol, COUNT(*) AS cnt FROM trades \
1333 GROUP BY symbol, TUMBLE(ts, INTERVAL '5' MINUTE) \
1334 HAVING COUNT(*) > 10",
1335 )
1336 .unwrap();
1337
1338 let plan = planner.plan(&statements[0]).unwrap();
1339 match plan {
1340 StreamingPlan::Query(query_plan) => {
1341 assert!(query_plan.window_config.is_some());
1342 assert!(query_plan.having_config.is_some());
1343 let config = query_plan.having_config.unwrap();
1344 assert!(
1345 config.predicate().contains("COUNT(*)"),
1346 "predicate was: {}",
1347 config.predicate()
1348 );
1349 }
1350 _ => panic!("Expected Query plan with having config"),
1351 }
1352 }
1353
1354 #[test]
1355 fn test_plan_query_without_having() {
1356 let mut planner = StreamingPlanner::new();
1357 let statements = StreamingParser::parse_sql(
1358 "SELECT COUNT(*) FROM events GROUP BY TUMBLE(event_time, INTERVAL '5' MINUTE)",
1359 )
1360 .unwrap();
1361
1362 let plan = planner.plan(&statements[0]).unwrap();
1363 match plan {
1364 StreamingPlan::Query(query_plan) => {
1365 assert!(query_plan.having_config.is_none());
1366 }
1367 _ => panic!("Expected Query plan"),
1368 }
1369 }
1370
1371 #[test]
1372 fn test_plan_having_only_produces_query_plan() {
1373 let mut planner = StreamingPlanner::new();
1375 let statements = StreamingParser::parse_sql(
1376 "SELECT category, SUM(amount) FROM orders GROUP BY category HAVING SUM(amount) > 1000",
1377 )
1378 .unwrap();
1379
1380 let plan = planner.plan(&statements[0]).unwrap();
1381 match plan {
1382 StreamingPlan::Query(query_plan) => {
1383 assert!(query_plan.having_config.is_some());
1384 assert!(query_plan.window_config.is_none());
1385 }
1386 _ => panic!("Expected Query plan for HAVING-only query"),
1387 }
1388 }
1389
1390 #[test]
1391 fn test_plan_having_compound_predicate() {
1392 let mut planner = StreamingPlanner::new();
1393 let statements = StreamingParser::parse_sql(
1394 "SELECT symbol, COUNT(*) AS cnt, SUM(vol) AS total \
1395 FROM trades GROUP BY symbol \
1396 HAVING COUNT(*) >= 5 AND SUM(vol) > 10000",
1397 )
1398 .unwrap();
1399
1400 let plan = planner.plan(&statements[0]).unwrap();
1401 match plan {
1402 StreamingPlan::Query(query_plan) => {
1403 let config = query_plan.having_config.unwrap();
1404 let pred = config.predicate();
1405 assert!(pred.contains("AND"), "predicate was: {pred}");
1406 }
1407 _ => panic!("Expected Query plan"),
1408 }
1409 }
1410
1411 #[test]
1412 fn test_plan_query_with_lead() {
1413 let mut planner = StreamingPlanner::new();
1414 let statements = StreamingParser::parse_sql(
1415 "SELECT LEAD(price, 2) OVER (ORDER BY ts) AS next2 FROM trades",
1416 )
1417 .unwrap();
1418
1419 let plan = planner.plan(&statements[0]).unwrap();
1420 match plan {
1421 StreamingPlan::Query(query_plan) => {
1422 assert!(query_plan.analytic_config.is_some());
1423 let config = query_plan.analytic_config.unwrap();
1424 assert!(config.has_lookahead());
1425 assert_eq!(config.functions[0].offset, 2);
1426 }
1427 _ => panic!("Expected Query plan with analytic config"),
1428 }
1429 }
1430
1431 #[test]
1434 fn test_plan_single_join_produces_vec_of_one() {
1435 let mut planner = StreamingPlanner::new();
1436 let statements = StreamingParser::parse_sql(
1437 "SELECT * FROM a JOIN b ON a.id = b.a_id \
1438 AND b.ts BETWEEN a.ts AND a.ts + INTERVAL '1' HOUR",
1439 )
1440 .unwrap();
1441
1442 let plan = planner.plan(&statements[0]).unwrap();
1443 match plan {
1444 StreamingPlan::Query(qp) => {
1445 let configs = qp.join_config.unwrap();
1446 assert_eq!(configs.len(), 1);
1447 }
1448 _ => panic!("Expected Query plan"),
1449 }
1450 }
1451
1452 #[test]
1453 fn test_plan_rejects_multi_way_interval_join() {
1454 let mut planner = StreamingPlanner::new();
1455 let statements = StreamingParser::parse_sql(
1456 "SELECT * FROM a JOIN b ON a.id = b.a_id \
1457 AND b.ts BETWEEN a.ts AND a.ts + INTERVAL '1' HOUR \
1458 JOIN c ON b.id = c.b_id \
1459 AND c.ts BETWEEN b.ts AND b.ts + INTERVAL '1' HOUR",
1460 )
1461 .unwrap();
1462
1463 let error = planner.plan(&statements[0]).unwrap_err();
1464 assert!(error
1465 .to_string()
1466 .contains("explicitly named two-way stages"));
1467 }
1468
1469 #[test]
1470 fn test_plan_rejects_mixed_multi_way_join() {
1471 let mut planner = StreamingPlanner::new();
1472 let _ = planner.plan(
1475 &StreamingParser::parse_sql(
1476 "CREATE LOOKUP TABLE customers (id BIGINT NOT NULL, name VARCHAR, \
1477 PRIMARY KEY (id)) WITH (connector = 'parquet', path = '/tmp/x.parquet')",
1478 )
1479 .unwrap()[0],
1480 );
1481 let statements = StreamingParser::parse_sql(
1482 "SELECT * FROM orders o \
1483 JOIN payments p ON o.id = p.order_id \
1484 AND p.ts BETWEEN o.ts AND o.ts + INTERVAL '1' HOUR \
1485 JOIN customers c ON p.cust_id = c.id",
1486 )
1487 .unwrap();
1488
1489 let error = planner.plan(&statements[0]).unwrap_err();
1490 assert!(error
1491 .to_string()
1492 .contains("explicitly named two-way stages"));
1493 }
1494
1495 #[test]
1496 fn test_plan_backward_compat_no_join() {
1497 let mut planner = StreamingPlanner::new();
1498 let statements = StreamingParser::parse_sql("SELECT * FROM orders").unwrap();
1499
1500 let plan = planner.plan(&statements[0]).unwrap();
1501 match plan {
1502 StreamingPlan::Standard(_) => {} _ => panic!("Expected Standard plan for simple SELECT"),
1504 }
1505 }
1506
1507 #[test]
1510 fn test_plan_query_with_rows_frame() {
1511 let mut planner = StreamingPlanner::new();
1512 let statements = StreamingParser::parse_sql(
1513 "SELECT AVG(price) OVER (ORDER BY ts \
1514 ROWS BETWEEN 9 PRECEDING AND CURRENT ROW) AS ma FROM trades",
1515 )
1516 .unwrap();
1517
1518 let plan = planner.plan(&statements[0]).unwrap();
1519 match plan {
1520 StreamingPlan::Query(qp) => {
1521 assert!(qp.frame_config.is_some());
1522 let fc = qp.frame_config.unwrap();
1523 assert_eq!(fc.functions.len(), 1);
1524 assert_eq!(fc.functions[0].source_column, "price");
1525 }
1526 _ => panic!("Expected Query plan with frame_config"),
1527 }
1528 }
1529
1530 #[test]
1531 fn test_plan_frame_with_partition() {
1532 let mut planner = StreamingPlanner::new();
1533 let statements = StreamingParser::parse_sql(
1534 "SELECT AVG(price) OVER (PARTITION BY symbol ORDER BY ts \
1535 ROWS BETWEEN 4 PRECEDING AND CURRENT ROW) AS ma FROM trades",
1536 )
1537 .unwrap();
1538
1539 let plan = planner.plan(&statements[0]).unwrap();
1540 match plan {
1541 StreamingPlan::Query(qp) => {
1542 let fc = qp.frame_config.unwrap();
1543 assert_eq!(fc.partition_columns, vec!["symbol".to_string()]);
1544 assert_eq!(fc.order_columns, vec!["ts".to_string()]);
1545 }
1546 _ => panic!("Expected Query plan with frame_config"),
1547 }
1548 }
1549
1550 #[test]
1551 fn test_plan_no_frame_is_standard() {
1552 let mut planner = StreamingPlanner::new();
1553 let statements = StreamingParser::parse_sql("SELECT * FROM trades").unwrap();
1554
1555 let plan = planner.plan(&statements[0]).unwrap();
1556 match plan {
1557 StreamingPlan::Standard(_) => {} _ => panic!("Expected Standard plan for simple SELECT"),
1559 }
1560 }
1561
1562 #[test]
1563 fn test_plan_unbounded_following_rejected() {
1564 let mut planner = StreamingPlanner::new();
1565 let statements = StreamingParser::parse_sql(
1566 "SELECT SUM(amount) OVER (ORDER BY id \
1567 ROWS BETWEEN CURRENT ROW AND UNBOUNDED FOLLOWING) AS rest \
1568 FROM orders",
1569 )
1570 .unwrap();
1571
1572 let result = planner.plan(&statements[0]);
1573 assert!(result.is_err());
1574 let err = result.unwrap_err().to_string();
1575 assert!(err.contains("UNBOUNDED FOLLOWING"), "error was: {err}");
1576 }
1577}