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::temporal::temporal_table_version_count;
36use crate::translator::{
37 AnalyticWindowConfig, JoinOperatorConfig, OrderOperatorConfig, WindowFrameConfig,
38 WindowOperatorConfig,
39};
40
41#[derive(Debug, Clone)]
43pub struct LookupTableInfo {
44 pub name: String,
46 pub columns: Vec<(String, String)>,
48 pub primary_key: Vec<String>,
50 pub properties: LookupTableProperties,
52 pub arrow_schema: SchemaRef,
54 pub raw_options: HashMap<String, String>,
56}
57
58pub struct StreamingPlanner {
60 sources: HashMap<String, SourceInfo>,
62 sinks: HashMap<String, SinkInfo>,
64 lookup_tables: HashMap<String, LookupTableInfo>,
66 windowed_views: std::collections::HashSet<String>,
68}
69
70#[derive(Debug, Clone)]
72pub struct SourceInfo {
73 pub name: String,
75 pub watermark_column: Option<String>,
77 pub primary_key: Vec<String>,
79}
80
81#[derive(Debug, Clone)]
83pub struct SinkInfo {
84 pub name: String,
86 pub from: String,
88}
89
90fn is_inline_unnest(factor: &TableFactor) -> bool {
91 match factor {
92 TableFactor::UNNEST { .. } => true,
93 TableFactor::Table {
94 name,
95 args: Some(_),
96 ..
97 }
98 | TableFactor::Function { name, .. } => {
99 name.0.len() == 1 && name.to_string().eq_ignore_ascii_case("unnest")
100 }
101 _ => false,
102 }
103}
104
105fn has_implicit_multi_source(select: &Select) -> bool {
106 select
107 .from
108 .iter()
109 .filter(|from| !is_inline_unnest(&from.relation))
110 .count()
111 > 1
112}
113
114#[derive(Debug)]
116#[allow(clippy::large_enum_variant)]
117pub enum StreamingPlan {
118 RegisterSource(SourceInfo),
120
121 RegisterSink(SinkInfo),
123
124 Query(QueryPlan),
126
127 Standard(Box<Statement>),
129
130 RegisterLookupTable(LookupTableInfo),
132
133 DropLookupTable {
135 name: String,
137 },
138}
139
140#[derive(Debug)]
142pub struct QueryPlan {
143 pub name: Option<String>,
145 pub window_config: Option<WindowOperatorConfig>,
147 pub join_config: Option<Vec<JoinOperatorConfig>>,
149 pub order_config: Option<OrderOperatorConfig>,
151 pub analytic_config: Option<AnalyticWindowConfig>,
153 pub frame_config: Option<WindowFrameConfig>,
155 pub emit_clause: Option<EmitClause>,
157 pub statement: Box<Statement>,
159}
160
161#[derive(Debug, Clone, PartialEq, Eq)]
166pub struct ChangelogEnrichAdmission {
167 left_table: String,
168 right_table: String,
169 left_keys: Vec<String>,
170 right_keys: Vec<String>,
171 join_type: JoinType,
172}
173
174impl ChangelogEnrichAdmission {
175 pub fn try_new(
180 left_table: impl Into<String>,
181 right_table: impl Into<String>,
182 left_keys: Vec<String>,
183 right_keys: Vec<String>,
184 left_outer: bool,
185 ) -> Result<Self, String> {
186 let left_table = left_table.into();
187 let right_table = right_table.into();
188 if left_table.is_empty() || right_table.is_empty() {
189 return Err("changelog enrichment relations cannot be empty".into());
190 }
191 if left_keys.is_empty()
192 || left_keys.len() != right_keys.len()
193 || left_keys.iter().chain(&right_keys).any(String::is_empty)
194 {
195 return Err("changelog enrichment keys must be non-empty with matching arity".into());
196 }
197 Ok(Self {
198 left_table,
199 right_table,
200 left_keys,
201 right_keys,
202 join_type: if left_outer {
203 JoinType::Left
204 } else {
205 JoinType::Inner
206 },
207 })
208 }
209
210 fn matches(&self, step: &JoinAnalysis) -> bool {
211 let mut left_keys = Vec::with_capacity(1 + step.additional_key_columns.len());
212 let mut right_keys = Vec::with_capacity(1 + step.additional_key_columns.len());
213 left_keys.push(step.left_key_column.as_str());
214 right_keys.push(step.right_key_column.as_str());
215 for (left, right) in &step.additional_key_columns {
216 left_keys.push(left);
217 right_keys.push(right);
218 }
219 self.left_table == step.left_table
220 && self.right_table == step.right_table
221 && self.join_type == step.join_type
222 && self.left_keys.iter().map(String::as_str).eq(left_keys)
223 && self.right_keys.iter().map(String::as_str).eq(right_keys)
224 && !step.is_temporal_join()
225 && step.time_bound.is_none()
226 }
227}
228
229impl StreamingPlanner {
230 #[must_use]
232 pub fn new() -> Self {
233 Self {
234 sources: HashMap::new(),
235 sinks: HashMap::new(),
236 lookup_tables: HashMap::new(),
237 windowed_views: std::collections::HashSet::new(),
238 }
239 }
240
241 pub fn plan(&mut self, statement: &StreamingStatement) -> Result<StreamingPlan, PlanningError> {
247 self.plan_internal(statement, None)
248 }
249
250 pub fn plan_changelog_enrich(
259 &mut self,
260 statement: &StreamingStatement,
261 admission: &ChangelogEnrichAdmission,
262 ) -> Result<StreamingPlan, PlanningError> {
263 if !matches!(
264 statement,
265 StreamingStatement::CreateContinuousQuery { .. }
266 | StreamingStatement::CreateStream { .. }
267 ) {
268 return Err(PlanningError::InvalidQuery(
269 "changelog enrichment admission is valid only for a named streaming query".into(),
270 ));
271 }
272 self.plan_internal(statement, Some(admission))
273 }
274
275 fn plan_internal(
276 &mut self,
277 statement: &StreamingStatement,
278 changelog_enrich: Option<&ChangelogEnrichAdmission>,
279 ) -> Result<StreamingPlan, PlanningError> {
280 match statement {
281 StreamingStatement::CreateSource(source) => self.plan_create_source(source),
282 StreamingStatement::CreateSink(sink) => self.plan_create_sink(sink),
283 StreamingStatement::CreateContinuousQuery {
284 name,
285 query,
286 emit_clause,
287 ..
288 }
289 | StreamingStatement::CreateStream {
290 name,
291 query,
292 emit_clause,
293 ..
294 } => self.plan_continuous_query(name, query, emit_clause.as_ref(), changelog_enrich),
295 StreamingStatement::Standard(stmt) => self.plan_standard_statement(stmt, None),
296 StreamingStatement::TemporalProbeQuery {
297 statement,
298 analysis,
299 } => self.plan_standard_statement(statement, Some(analysis)),
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 primary_key: source
377 .primary_key
378 .iter()
379 .map(|column| column.value.clone())
380 .collect(),
381 };
382
383 self.sources.insert(name, info.clone());
385
386 Ok(StreamingPlan::RegisterSource(info))
387 }
388
389 fn plan_create_sink(
391 &mut self,
392 sink: &CreateSinkStatement,
393 ) -> Result<StreamingPlan, PlanningError> {
394 let name = object_name_to_string(&sink.name);
395
396 if !sink.or_replace && !sink.if_not_exists && self.sinks.contains_key(&name) {
398 return Err(PlanningError::InvalidQuery(format!(
399 "Sink '{}' already exists",
400 name
401 )));
402 }
403
404 let from = match &sink.from {
406 SinkFrom::Table(table) => object_name_to_string(table),
407 SinkFrom::Query(_) => format!("{}_query", name),
408 };
409
410 let info = SinkInfo {
411 name: name.clone(),
412 from,
413 };
414
415 self.sinks.insert(name, info.clone());
417
418 Ok(StreamingPlan::RegisterSink(info))
419 }
420
421 fn plan_continuous_query(
423 &mut self,
424 name: &ObjectName,
425 query: &StreamingStatement,
426 emit_clause: Option<&EmitClause>,
427 changelog_enrich: Option<&ChangelogEnrichAdmission>,
428 ) -> Result<StreamingPlan, PlanningError> {
429 let (stmt, temporal_probe) = match query {
431 StreamingStatement::Standard(stmt) => (stmt.as_ref().clone(), None),
432 StreamingStatement::TemporalProbeQuery {
433 statement,
434 analysis,
435 } => (statement.as_ref().clone(), Some(analysis.as_ref())),
436 _ => {
437 return Err(PlanningError::InvalidQuery(
438 "Continuous query must contain a SELECT statement".to_string(),
439 ))
440 }
441 };
442
443 let query_plan =
445 self.analyze_query(&stmt, emit_clause, changelog_enrich, temporal_probe)?;
446
447 let view_name = object_name_to_string(name);
450 if query_plan.window_config.is_some() {
451 self.windowed_views.insert(view_name);
452 } else {
453 self.windowed_views.remove(&view_name);
454 }
455
456 Ok(StreamingPlan::Query(QueryPlan {
457 name: Some(object_name_to_string(name)),
458 window_config: query_plan.window_config,
459 join_config: query_plan.join_config,
460 order_config: query_plan.order_config,
461 analytic_config: query_plan.analytic_config,
462 frame_config: query_plan.frame_config,
463 emit_clause: emit_clause.cloned(),
464 statement: Box::new(stmt),
465 }))
466 }
467
468 #[allow(clippy::unused_self)] fn plan_standard_statement(
471 &self,
472 stmt: &Statement,
473 temporal_probe: Option<&JoinAnalysis>,
474 ) -> Result<StreamingPlan, PlanningError> {
475 if let Statement::Query(query) = stmt {
477 if let SetExpr::Select(select) = query.body.as_ref() {
478 if has_implicit_multi_source(select) {
479 return Err(PlanningError::InvalidQuery(
480 "implicit multi-source joins are unsupported; use one bounded INNER JOIN"
481 .to_string(),
482 ));
483 }
484 let window_function = Self::extract_window_from_select(select);
486
487 let mut join_analysis = analyze_joins(select).map_err(|e| {
489 PlanningError::InvalidQuery(format!("Join analysis failed: {e}"))
490 })?;
491
492 validate_temporal_version_shape(
493 stmt,
494 join_analysis.as_ref().map_or(0, |multi| {
495 multi
496 .joins
497 .iter()
498 .filter(|join| join.is_temporal_join())
499 .count()
500 }),
501 )?;
502
503 if let Some(ref mut multi) = join_analysis {
504 apply_temporal_probe_analysis(multi, temporal_probe)?;
505 self.resolve_temporal_source_contracts(multi)?;
506 validate_streaming_joins(multi, &self.lookup_tables, None)?;
507 }
508
509 let order_analysis = analyze_order_by(stmt);
511 let order_config = OrderOperatorConfig::from_analysis(&order_analysis)
512 .map_err(PlanningError::InvalidQuery)?;
513
514 let analytic_analysis = analyze_analytic_functions(stmt);
516 let analytic_config =
517 analytic_analysis.map(|a| AnalyticWindowConfig::from_analysis(&a));
518
519 let has_having = analyze_aggregates(stmt).has_having;
520
521 let frame_analysis = analyze_window_frames(stmt);
523 let frame_config = frame_analysis
524 .as_ref()
525 .map(WindowFrameConfig::from_analysis);
526
527 if let Some(fa) = &frame_analysis {
529 for f in &fa.functions {
530 if matches!(f.end_bound, FrameBound::UnboundedFollowing) {
531 return Err(PlanningError::InvalidQuery(
532 "UNBOUNDED FOLLOWING is not supported in streaming window frames"
533 .to_string(),
534 ));
535 }
536 }
537 }
538
539 let has_streaming_features = window_function.is_some()
540 || join_analysis.is_some()
541 || order_config.is_some()
542 || analytic_config.is_some()
543 || has_having
544 || frame_config.is_some();
545
546 if has_streaming_features {
547 let window_config = match window_function {
548 Some(w) => Some(
549 WindowOperatorConfig::from_window_function(&w)
550 .map_err(|e| PlanningError::InvalidQuery(e.to_string()))?,
551 ),
552 None => None,
553 };
554
555 let join_config = join_analysis
556 .map(|m| JoinOperatorConfig::from_multi_analysis(&m))
557 .transpose()
558 .map_err(PlanningError::InvalidQuery)?;
559
560 return Ok(StreamingPlan::Query(QueryPlan {
561 name: None,
562 window_config,
563 join_config,
564 order_config,
565 analytic_config,
566 frame_config,
567 emit_clause: None,
568 statement: Box::new(stmt.clone()),
569 }));
570 }
571 }
572 }
573
574 validate_temporal_version_shape(stmt, 0)?;
575
576 Ok(StreamingPlan::Standard(Box::new(stmt.clone())))
578 }
579
580 fn analyze_query(
582 &self,
583 stmt: &Statement,
584 emit_clause: Option<&EmitClause>,
585 changelog_enrich: Option<&ChangelogEnrichAdmission>,
586 temporal_probe: Option<&JoinAnalysis>,
587 ) -> Result<QueryAnalysis, PlanningError> {
588 let mut analysis = QueryAnalysis::default();
589 let mut recognized_temporal_versions = 0;
590
591 if let Statement::Query(query) = stmt {
592 if let SetExpr::Select(select) = query.body.as_ref() {
593 if has_implicit_multi_source(select) {
594 return Err(PlanningError::InvalidQuery(
595 "implicit multi-source joins are unsupported; use one bounded INNER JOIN"
596 .to_string(),
597 ));
598 }
599 if let Some(window) = Self::extract_window_from_select(select) {
601 let mut config = WindowOperatorConfig::from_window_function(&window)
602 .map_err(|e| PlanningError::InvalidQuery(e.to_string()))?;
603
604 if let Some(emit) = emit_clause {
606 config = config
607 .with_emit_clause(emit)
608 .map_err(|e| PlanningError::InvalidQuery(e.to_string()))?;
609 }
610
611 analysis.window_config = Some(config);
612 }
613
614 let join_analysis = analyze_joins(select).map_err(|e| {
616 PlanningError::InvalidQuery(format!("Join analysis failed: {e}"))
617 })?;
618 recognized_temporal_versions = join_analysis.as_ref().map_or(0, |multi| {
619 multi
620 .joins
621 .iter()
622 .filter(|join| join.is_temporal_join())
623 .count()
624 });
625 if let Some(mut multi) = join_analysis {
626 apply_temporal_probe_analysis(&mut multi, temporal_probe)?;
627 self.resolve_temporal_source_contracts(&mut multi)?;
628 validate_streaming_joins(&multi, &self.lookup_tables, changelog_enrich)?;
629 analysis.join_config = Some(
630 JoinOperatorConfig::from_multi_analysis(&multi)
631 .map_err(PlanningError::InvalidQuery)?,
632 );
633 }
634 }
635 }
636
637 validate_temporal_version_shape(stmt, recognized_temporal_versions)?;
638
639 let order_analysis = analyze_order_by(stmt);
641 analysis.order_config = OrderOperatorConfig::from_analysis(&order_analysis)
642 .map_err(PlanningError::InvalidQuery)?;
643
644 if let Some(analytic) = analyze_analytic_functions(stmt) {
646 analysis.analytic_config = Some(AnalyticWindowConfig::from_analysis(&analytic));
647 }
648
649 if let Some(frame_analysis) = analyze_window_frames(stmt) {
651 for f in &frame_analysis.functions {
653 if matches!(f.end_bound, FrameBound::UnboundedFollowing) {
654 return Err(PlanningError::InvalidQuery(
655 "UNBOUNDED FOLLOWING is not supported in streaming window frames"
656 .to_string(),
657 ));
658 }
659 }
660 analysis.frame_config = Some(WindowFrameConfig::from_analysis(&frame_analysis));
661 }
662
663 Ok(analysis)
664 }
665
666 fn extract_window_from_select(select: &sqlparser::ast::Select) -> Option<WindowFunction> {
668 use sqlparser::ast::GroupByExpr;
670 match &select.group_by {
671 GroupByExpr::Expressions(exprs, _modifiers) => {
672 for group_by_expr in exprs {
673 if let Ok(Some(window)) = WindowRewriter::extract_window_function(group_by_expr)
674 {
675 return Some(window);
676 }
677 }
678 }
679 GroupByExpr::All(_) => {}
680 }
681 None
682 }
683
684 fn plan_create_lookup_table(
686 &mut self,
687 lt: &CreateLookupTableStatement,
688 ) -> Result<StreamingPlan, PlanningError> {
689 let name = object_name_to_string(<.name);
690
691 if !lt.or_replace && !lt.if_not_exists && self.lookup_tables.contains_key(&name) {
692 return Err(PlanningError::InvalidQuery(format!(
693 "Lookup table '{}' already exists",
694 name
695 )));
696 }
697
698 let columns: Vec<(String, String)> = lt
699 .columns
700 .iter()
701 .map(|c| (c.name.value.clone(), c.data_type.to_string()))
702 .collect();
703
704 let properties = validate_properties(<.with_options).map_err(|e| {
705 PlanningError::InvalidQuery(format!("Invalid lookup table properties: {e}"))
706 })?;
707
708 let arrow_fields: Vec<Field> = lt
710 .columns
711 .iter()
712 .map(|c| {
713 let dt = crate::translator::streaming_ddl::sql_type_to_arrow(&c.data_type)
714 .map_err(|e| PlanningError::InvalidQuery(e.to_string()))?;
715 let nullable = !c
716 .options
717 .iter()
718 .any(|opt| matches!(opt.option, sqlparser::ast::ColumnOption::NotNull));
719 Ok(Field::new(&c.name.value, dt, nullable))
720 })
721 .collect::<Result<_, PlanningError>>()?;
722 let arrow_schema = Arc::new(Schema::new(arrow_fields));
723
724 let info = LookupTableInfo {
725 name: name.clone(),
726 columns,
727 primary_key: lt.primary_key.clone(),
728 properties,
729 arrow_schema,
730 raw_options: lt.with_options.clone(),
731 };
732
733 self.lookup_tables.insert(name, info.clone());
734
735 Ok(StreamingPlan::RegisterLookupTable(info))
736 }
737
738 fn plan_drop_lookup_table(
740 &mut self,
741 name: &ObjectName,
742 if_exists: bool,
743 ) -> Result<StreamingPlan, PlanningError> {
744 let name_str = object_name_to_string(name);
745
746 if !if_exists && !self.lookup_tables.contains_key(&name_str) {
747 return Err(PlanningError::InvalidQuery(format!(
748 "Lookup table '{}' does not exist",
749 name_str
750 )));
751 }
752
753 self.lookup_tables.remove(&name_str);
754
755 Ok(StreamingPlan::DropLookupTable { name: name_str })
756 }
757
758 #[must_use]
760 pub fn get_source(&self, name: &str) -> Option<&SourceInfo> {
761 self.sources.get(name)
762 }
763
764 #[must_use]
766 pub fn get_sink(&self, name: &str) -> Option<&SinkInfo> {
767 self.sinks.get(name)
768 }
769
770 #[must_use]
772 pub fn list_sources(&self) -> Vec<&SourceInfo> {
773 self.sources.values().collect()
774 }
775
776 fn resolve_temporal_source_contracts(
777 &self,
778 multi: &mut MultiJoinAnalysis,
779 ) -> Result<(), PlanningError> {
780 for step in &mut multi.joins {
781 if !step.is_temporal_join() {
782 continue;
783 }
784 let (_, right_key_columns) = temporal_key_columns(step)?;
785 let left = self.sources.get(&step.left_table).ok_or_else(|| {
786 PlanningError::SourceNotFound(format!(
787 "{} (temporal left input must be a registered event-time source)",
788 step.left_table
789 ))
790 })?;
791 let left_time = step.left_time_column.as_ref().ok_or_else(|| {
792 PlanningError::InvalidQuery(
793 "temporal join is missing its explicit left event-time column".into(),
794 )
795 })?;
796 let left_watermark = left.watermark_column.as_ref().ok_or_else(|| {
797 PlanningError::InvalidQuery(format!(
798 "temporal left source '{}' must declare WATERMARK FOR {}",
799 step.left_table, left_time
800 ))
801 })?;
802 if left_watermark != left_time {
803 return Err(PlanningError::InvalidQuery(format!(
804 "temporal left timestamp '{}' does not match WATERMARK FOR {} on source '{}'",
805 left_time, left_watermark, step.left_table
806 )));
807 }
808 let right = self.sources.get(&step.right_table).ok_or_else(|| {
809 PlanningError::SourceNotFound(format!(
810 "{} (temporal right input must be a registered source)",
811 step.right_table
812 ))
813 })?;
814 if !right
815 .primary_key
816 .iter()
817 .map(String::as_str)
818 .eq(right_key_columns.iter().copied())
819 {
820 return Err(PlanningError::InvalidQuery(format!(
821 "temporal right source '{}' must declare PRIMARY KEY ({}) matching the join key",
822 step.right_table,
823 right_key_columns.join(", ")
824 )));
825 }
826 let right_time = right.watermark_column.as_ref().ok_or_else(|| {
827 PlanningError::InvalidQuery(format!(
828 "temporal right source '{}' must declare WATERMARK FOR its version column",
829 step.right_table
830 ))
831 })?;
832 if step
833 .right_time_column
834 .as_ref()
835 .is_some_and(|column| column != right_time)
836 {
837 return Err(PlanningError::InvalidQuery(format!(
838 "temporal right timestamp '{}' does not match WATERMARK FOR {} on source '{}'",
839 step.right_time_column.as_deref().unwrap_or_default(),
840 right_time,
841 step.right_table
842 )));
843 }
844 step.right_time_column = Some(right_time.clone());
845 }
846 Ok(())
847 }
848
849 #[must_use]
851 pub fn list_sinks(&self) -> Vec<&SinkInfo> {
852 self.sinks.values().collect()
853 }
854
855 #[must_use]
857 pub fn get_lookup_table(&self, name: &str) -> Option<&LookupTableInfo> {
858 self.lookup_tables.get(name)
859 }
860
861 #[must_use]
863 pub fn list_lookup_tables(&self) -> Vec<&LookupTableInfo> {
864 self.lookup_tables.values().collect()
865 }
866
867 #[must_use]
869 pub fn lookup_tables_cloned(&self) -> HashMap<String, LookupTableInfo> {
870 self.lookup_tables.clone()
871 }
872
873 #[allow(clippy::unused_self)] pub async fn to_logical_plan(
884 &self,
885 plan: &QueryPlan,
886 ctx: &SessionContext,
887 ) -> Result<LogicalPlan, PlanningError> {
888 let sql = plan.statement.to_string();
892 ctx.state()
893 .create_logical_plan(&sql)
894 .await
895 .map_err(PlanningError::DataFusion)
896 }
897}
898
899impl Default for StreamingPlanner {
900 fn default() -> Self {
901 Self::new()
902 }
903}
904
905#[derive(Debug, Default)]
907#[allow(clippy::struct_field_names)]
908struct QueryAnalysis {
909 window_config: Option<WindowOperatorConfig>,
910 join_config: Option<Vec<JoinOperatorConfig>>,
911 order_config: Option<OrderOperatorConfig>,
912 analytic_config: Option<AnalyticWindowConfig>,
913 frame_config: Option<WindowFrameConfig>,
914}
915
916fn object_name_to_string(name: &ObjectName) -> String {
918 match name.0.as_slice() {
919 [sqlparser::ast::ObjectNamePart::Identifier(ident)] => ident.value.clone(),
920 _ => name.to_string(),
921 }
922}
923
924fn temporal_key_columns(step: &JoinAnalysis) -> Result<(Vec<&str>, Vec<&str>), PlanningError> {
925 let mut left = Vec::with_capacity(1 + step.additional_key_columns.len());
926 let mut right = Vec::with_capacity(1 + step.additional_key_columns.len());
927 left.push(step.left_key_column.as_str());
928 right.push(step.right_key_column.as_str());
929 for (left_column, right_column) in &step.additional_key_columns {
930 left.push(left_column);
931 right.push(right_column);
932 }
933 if left.is_empty()
934 || left.len() != right.len()
935 || left.iter().chain(&right).any(|column| column.is_empty())
936 {
937 return Err(PlanningError::InvalidQuery(
938 "temporal join equality keys must be non-empty and have matching cardinality".into(),
939 ));
940 }
941 Ok((left, right))
942}
943
944fn apply_temporal_probe_analysis(
945 multi: &mut MultiJoinAnalysis,
946 temporal_probe: Option<&JoinAnalysis>,
947) -> Result<(), PlanningError> {
948 let Some(temporal_probe) = temporal_probe else {
949 return Ok(());
950 };
951 let [normalized] = multi.joins.as_slice() else {
952 return Err(PlanningError::InvalidQuery(
953 "TEMPORAL PROBE JOIN requires one explicitly named two-way stage".into(),
954 ));
955 };
956 let (normalized_left_keys, normalized_right_keys) = temporal_key_columns(normalized)?;
957 let (probe_left_keys, probe_right_keys) = temporal_key_columns(temporal_probe)?;
958 if !normalized.is_temporal_join()
959 || normalized.left_table != temporal_probe.left_table
960 || normalized.right_table != temporal_probe.right_table
961 || normalized_left_keys != probe_left_keys
962 || normalized_right_keys != probe_right_keys
963 || normalized.left_time_column != temporal_probe.left_time_column
964 || normalized.join_type != temporal_probe.join_type
965 {
966 return Err(PlanningError::InvalidQuery(
967 "TEMPORAL PROBE JOIN metadata does not match its normalized AS-OF plan".into(),
968 ));
969 }
970 multi.joins[0] = temporal_probe.clone();
971 Ok(())
972}
973
974fn validate_temporal_version_shape(
975 statement: &Statement,
976 recognized_versions: usize,
977) -> Result<(), PlanningError> {
978 let ast_versions = temporal_table_version_count(statement);
979 if ast_versions != recognized_versions {
980 return Err(PlanningError::InvalidQuery(
981 "FOR SYSTEM_TIME AS OF is supported only on the right input of one direct two-input temporal join; nested and set-operation temporal joins are unsupported"
982 .into(),
983 ));
984 }
985 Ok(())
986}
987
988fn validate_streaming_joins(
990 multi: &MultiJoinAnalysis,
991 lookup_tables: &HashMap<String, LookupTableInfo>,
992 changelog_enrich: Option<&ChangelogEnrichAdmission>,
993) -> Result<(), PlanningError> {
994 if multi.joins.len() != 1 {
995 return Err(PlanningError::InvalidQuery(
996 "multi-way streaming joins require explicitly named two-way stages".to_string(),
997 ));
998 }
999 for step in &multi.joins {
1000 if step.time_bound.is_some_and(|bound| bound.is_zero()) {
1001 return Err(PlanningError::InvalidQuery(
1002 "streaming interval joins require a positive finite time bound".to_string(),
1003 ));
1004 }
1005 if step.is_bounded() {
1006 continue;
1007 }
1008 let left_lookup = lookup_tables.contains_key(&step.left_table);
1009 let right_lookup = lookup_tables.contains_key(&step.right_table);
1010 if !left_lookup && !right_lookup {
1011 if changelog_enrich.is_some_and(|admission| admission.matches(step)) {
1012 continue;
1013 }
1014 return Err(PlanningError::InvalidQuery(format!(
1015 "unbounded join between streaming sources '{}' and '{}'; \
1016 add a temporal predicate or use a lookup table",
1017 step.left_table, step.right_table,
1018 )));
1019 }
1020 }
1021 Ok(())
1022}
1023
1024#[derive(Debug, thiserror::Error)]
1026pub enum PlanningError {
1027 UnsupportedSql(String),
1029
1030 InvalidQuery(String),
1032
1033 SourceNotFound(String),
1035
1036 SinkNotFound(String),
1038
1039 DataFusion(#[from] datafusion_common::DataFusionError),
1041}
1042
1043impl std::fmt::Display for PlanningError {
1044 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1045 match self {
1046 Self::UnsupportedSql(msg) => write!(f, "Unsupported SQL: {msg}"),
1047 Self::InvalidQuery(msg) => write!(f, "Invalid query: {msg}"),
1048 Self::SourceNotFound(name) => write!(f, "Source not found: {name}"),
1049 Self::SinkNotFound(name) => write!(f, "Sink not found: {name}"),
1050 Self::DataFusion(e) => {
1051 let translated = crate::error::translate_datafusion_error(&e.to_string());
1052 write!(f, "{translated}")
1053 }
1054 }
1055 }
1056}
1057
1058#[cfg(test)]
1059mod tests;