Skip to main content

laminar_sql/planner/
mod.rs

1//! Query planner for streaming SQL
2//!
3//! This module translates parsed streaming SQL statements into execution plans.
4//! It integrates with the parser and translator modules to produce complete
5//! operator configurations for Ring 0 execution.
6
7pub mod channel_derivation;
8/// Optimizer rules for lookup join rewriting.
9pub mod lookup_join;
10/// Predicate splitting and pushdown for lookup joins.
11pub mod predicate_split;
12/// Physical optimizer rule for streaming plan validation.
13pub mod streaming_optimizer;
14
15#[allow(clippy::disallowed_types)] // cold path: query planning
16use 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/// Information about a registered lookup table.
42#[derive(Debug, Clone)]
43pub struct LookupTableInfo {
44    /// Table name.
45    pub name: String,
46    /// Column names and types.
47    pub columns: Vec<(String, String)>,
48    /// Primary key columns.
49    pub primary_key: Vec<String>,
50    /// Validated properties.
51    pub properties: LookupTableProperties,
52    /// Pre-computed Arrow schema from column definitions.
53    pub arrow_schema: SchemaRef,
54    /// Raw WITH options for connector configuration pass-through.
55    pub raw_options: HashMap<String, String>,
56}
57
58/// Streaming query planner
59pub struct StreamingPlanner {
60    /// Registered sources
61    sources: HashMap<String, SourceInfo>,
62    /// Registered sinks
63    sinks: HashMap<String, SinkInfo>,
64    /// Registered lookup tables
65    lookup_tables: HashMap<String, LookupTableInfo>,
66    /// Names of views/streams for which planning retains window classification.
67    windowed_views: std::collections::HashSet<String>,
68}
69
70/// Information about a registered source
71#[derive(Debug, Clone)]
72pub struct SourceInfo {
73    /// Source name
74    pub name: String,
75    /// Watermark column (if configured)
76    pub watermark_column: Option<String>,
77    /// Declared non-null primary-key columns.
78    pub primary_key: Vec<String>,
79}
80
81/// Information about a registered sink
82#[derive(Debug, Clone)]
83pub struct SinkInfo {
84    /// Sink name
85    pub name: String,
86    /// Source table or query name
87    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/// Result of planning a streaming statement
115#[derive(Debug)]
116#[allow(clippy::large_enum_variant)]
117pub enum StreamingPlan {
118    /// Source registration (DDL)
119    RegisterSource(SourceInfo),
120
121    /// Sink registration (DDL)
122    RegisterSink(SinkInfo),
123
124    /// Query plan with streaming configurations
125    Query(QueryPlan),
126
127    /// Standard SQL statement (pass-through to DataFusion)
128    Standard(Box<Statement>),
129
130    /// Lookup table registration (DDL)
131    RegisterLookupTable(LookupTableInfo),
132
133    /// Drop a lookup table
134    DropLookupTable {
135        /// Name of the lookup table to drop.
136        name: String,
137    },
138}
139
140/// A query plan with streaming operator configurations
141#[derive(Debug)]
142pub struct QueryPlan {
143    /// Optional name for the continuous query
144    pub name: Option<String>,
145    /// Window configuration if the query has windowed aggregation
146    pub window_config: Option<WindowOperatorConfig>,
147    /// Join configuration(s) if the query has joins (one per join step)
148    pub join_config: Option<Vec<JoinOperatorConfig>>,
149    /// ORDER BY configuration if the query has ordering
150    pub order_config: Option<OrderOperatorConfig>,
151    /// Analytic window function configuration (LAG/LEAD/etc.)
152    pub analytic_config: Option<AnalyticWindowConfig>,
153    /// Window frame configuration (ROWS BETWEEN / RANGE BETWEEN)
154    pub frame_config: Option<WindowFrameConfig>,
155    /// Emit strategy
156    pub emit_clause: Option<EmitClause>,
157    /// The underlying SQL statement
158    pub statement: Box<Statement>,
159}
160
161/// Exact one-call admission for database-certified changelog-to-static enrichment.
162///
163/// The certificate is matched against the planner's independent parse, so it
164/// cannot authorize another relation pair, key mapping, or join type.
165#[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    /// Construct an exact INNER or LEFT changelog enrichment admission.
176    ///
177    /// # Errors
178    /// Returns an error for empty relations/keys or a mismatched key arity.
179    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    /// Creates a new streaming planner
231    #[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    /// Plans a streaming statement.
242    ///
243    /// # Errors
244    ///
245    /// Returns `PlanningError` if the statement cannot be planned.
246    pub fn plan(&mut self, statement: &StreamingStatement) -> Result<StreamingPlan, PlanningError> {
247        self.plan_internal(statement, None)
248    }
249
250    /// Plan one query whose unbounded join shape has been certified as
251    /// changelog-to-static enrichment by the database.
252    ///
253    /// This does not relax multi-way, join-type, temporal, or implicit-join
254    /// validation. Raw stream callers must use [`Self::plan`].
255    ///
256    /// # Errors
257    /// Returns `PlanningError` if the statement fails the remaining streaming checks.
258    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                // These statements are handled directly by the database facade
319                // and don't need query planning. Return as Standard pass-through.
320                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    /// Remove query classification installed by a successful plan when the surrounding catalog
329    /// transaction is rolled back or the query is dropped.
330    pub fn unregister_query(&mut self, name: &str) {
331        self.windowed_views.remove(name);
332    }
333
334    /// Whether query classification state remains for a catalog object.
335    #[must_use]
336    pub fn has_query(&self, name: &str) -> bool {
337        self.windowed_views.contains(name)
338    }
339
340    /// Remove a source installed by a catalog transaction that was dropped or rolled back.
341    pub fn unregister_source(&mut self, name: &str) {
342        self.sources.remove(name);
343    }
344
345    /// Remove a sink installed by a catalog transaction that was dropped or rolled back.
346    pub fn unregister_sink(&mut self, name: &str) {
347        self.sinks.remove(name);
348    }
349
350    /// Remove a lookup table installed by a catalog transaction that was dropped or rolled back.
351    pub fn unregister_lookup_table(&mut self, name: &str) {
352        self.lookup_tables.remove(name);
353    }
354
355    /// Plans a CREATE SOURCE statement.
356    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        // Check for existing source
363        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        // Extract watermark column
371        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        // Register the source
384        self.sources.insert(name, info.clone());
385
386        Ok(StreamingPlan::RegisterSource(info))
387    }
388
389    /// Plans a CREATE SINK statement.
390    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        // Check for existing sink
397        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        // Determine the source
405        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        // Register the sink
416        self.sinks.insert(name, info.clone());
417
418        Ok(StreamingPlan::RegisterSink(info))
419    }
420
421    /// Plans a CREATE CONTINUOUS QUERY statement.
422    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        // The query inside should be a standard SELECT
430        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        // Analyze the query for streaming features
444        let query_plan =
445            self.analyze_query(&stmt, emit_clause, changelog_enrich, temporal_probe)?;
446
447        // Keep planner classification in sync with catalog rollback/drop. A windowed query is the
448        // only query shape for which the planner retains classification after planning.
449        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    /// Plans a standard SQL statement.
469    #[allow(clippy::unused_self)] // Will use planner state for plan optimization
470    fn plan_standard_statement(
471        &self,
472        stmt: &Statement,
473        temporal_probe: Option<&JoinAnalysis>,
474    ) -> Result<StreamingPlan, PlanningError> {
475        // Check if it's a query that might have streaming features
476        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                // Check for window functions in GROUP BY
485                let window_function = Self::extract_window_from_select(select);
486
487                // Check for joins (multi-way)
488                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                // Check for ORDER BY
510                let order_analysis = analyze_order_by(stmt);
511                let order_config = OrderOperatorConfig::from_analysis(&order_analysis)
512                    .map_err(PlanningError::InvalidQuery)?;
513
514                // Check for analytic functions (LAG/LEAD/etc.)
515                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                // Check for window frame functions (ROWS BETWEEN / RANGE BETWEEN)
522                let frame_analysis = analyze_window_frames(stmt);
523                let frame_config = frame_analysis
524                    .as_ref()
525                    .map(WindowFrameConfig::from_analysis);
526
527                // Validate: reject UNBOUNDED FOLLOWING (streaming can't buffer infinite future)
528                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        // Pass through standard SQL
577        Ok(StreamingPlan::Standard(Box::new(stmt.clone())))
578    }
579
580    /// Analyzes a query for streaming features.
581    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                // Extract window function
600                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                    // Apply emit clause if present
605                    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                // Extract join info (multi-way)
615                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        // Extract ORDER BY info
640        let order_analysis = analyze_order_by(stmt);
641        analysis.order_config = OrderOperatorConfig::from_analysis(&order_analysis)
642            .map_err(PlanningError::InvalidQuery)?;
643
644        // Extract analytic function info (LAG/LEAD/etc.)
645        if let Some(analytic) = analyze_analytic_functions(stmt) {
646            analysis.analytic_config = Some(AnalyticWindowConfig::from_analysis(&analytic));
647        }
648
649        // Extract window frame functions (ROWS BETWEEN / RANGE BETWEEN)
650        if let Some(frame_analysis) = analyze_window_frames(stmt) {
651            // Validate: reject UNBOUNDED FOLLOWING
652            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    /// Extracts window function from a SELECT.
667    fn extract_window_from_select(select: &sqlparser::ast::Select) -> Option<WindowFunction> {
668        // Check GROUP BY for window functions
669        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    /// Plans a CREATE LOOKUP TABLE statement.
685    fn plan_create_lookup_table(
686        &mut self,
687        lt: &CreateLookupTableStatement,
688    ) -> Result<StreamingPlan, PlanningError> {
689        let name = object_name_to_string(&lt.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(&lt.with_options).map_err(|e| {
705            PlanningError::InvalidQuery(format!("Invalid lookup table properties: {e}"))
706        })?;
707
708        // Compute Arrow schema from column definitions
709        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    /// Plans a DROP LOOKUP TABLE statement.
739    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    /// Gets a registered source by name.
759    #[must_use]
760    pub fn get_source(&self, name: &str) -> Option<&SourceInfo> {
761        self.sources.get(name)
762    }
763
764    /// Gets a registered sink by name.
765    #[must_use]
766    pub fn get_sink(&self, name: &str) -> Option<&SinkInfo> {
767        self.sinks.get(name)
768    }
769
770    /// Lists all registered sources.
771    #[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    /// Lists all registered sinks.
850    #[must_use]
851    pub fn list_sinks(&self) -> Vec<&SinkInfo> {
852        self.sinks.values().collect()
853    }
854
855    /// Gets a registered lookup table by name.
856    #[must_use]
857    pub fn get_lookup_table(&self, name: &str) -> Option<&LookupTableInfo> {
858        self.lookup_tables.get(name)
859    }
860
861    /// Lists all registered lookup tables.
862    #[must_use]
863    pub fn list_lookup_tables(&self) -> Vec<&LookupTableInfo> {
864        self.lookup_tables.values().collect()
865    }
866
867    /// Returns a clone of the lookup tables map for optimizer rule construction.
868    #[must_use]
869    pub fn lookup_tables_cloned(&self) -> HashMap<String, LookupTableInfo> {
870        self.lookup_tables.clone()
871    }
872
873    /// Converts a query plan's SQL statement into a `DataFusion`
874    /// `LogicalPlan`. Window UDFs (TUMBLE, HOP, SESSION) must be registered
875    /// on `ctx` via
876    /// [`register_streaming_functions`](crate::datafusion::register_streaming_functions)
877    /// for windowed queries to resolve correctly.
878    ///
879    /// # Errors
880    ///
881    /// Returns `PlanningError` if `DataFusion` cannot create the logical plan.
882    #[allow(clippy::unused_self)] // Method will use planner state for plan optimization
883    pub async fn to_logical_plan(
884        &self,
885        plan: &QueryPlan,
886        ctx: &SessionContext,
887    ) -> Result<LogicalPlan, PlanningError> {
888        // Convert the AST statement back to SQL and let DataFusion re-parse
889        // it with its own sqlparser version. This avoids version mismatches
890        // between our sqlparser (0.60) and DataFusion's (0.59).
891        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/// Intermediate query analysis result
906#[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
916/// Helper to convert `ObjectName` to String
917fn 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
988/// Fail closed on stream-join shapes whose state/output semantics are not implemented.
989fn 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/// Planning errors
1025#[derive(Debug, thiserror::Error)]
1026pub enum PlanningError {
1027    /// Unsupported SQL feature
1028    UnsupportedSql(String),
1029
1030    /// Invalid query
1031    InvalidQuery(String),
1032
1033    /// Source not found
1034    SourceNotFound(String),
1035
1036    /// Sink not found
1037    SinkNotFound(String),
1038
1039    /// `DataFusion` error during logical plan creation (translated on display)
1040    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;