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::translator::{
36    AnalyticWindowConfig, HavingFilterConfig, JoinOperatorConfig, OrderOperatorConfig,
37    WindowFrameConfig, WindowOperatorConfig,
38};
39
40/// Information about a registered lookup table.
41#[derive(Debug, Clone)]
42pub struct LookupTableInfo {
43    /// Table name.
44    pub name: String,
45    /// Column names and types.
46    pub columns: Vec<(String, String)>,
47    /// Primary key columns.
48    pub primary_key: Vec<String>,
49    /// Validated properties.
50    pub properties: LookupTableProperties,
51    /// Pre-computed Arrow schema from column definitions.
52    pub arrow_schema: SchemaRef,
53    /// Raw WITH options for connector configuration pass-through.
54    pub raw_options: HashMap<String, String>,
55}
56
57/// Streaming query planner
58pub struct StreamingPlanner {
59    /// Registered sources
60    sources: HashMap<String, SourceInfo>,
61    /// Registered sinks
62    sinks: HashMap<String, SinkInfo>,
63    /// Registered lookup tables
64    lookup_tables: HashMap<String, LookupTableInfo>,
65    /// Names of views/streams for which planning retains window classification.
66    windowed_views: std::collections::HashSet<String>,
67}
68
69/// Information about a registered source
70#[derive(Debug, Clone)]
71pub struct SourceInfo {
72    /// Source name
73    pub name: String,
74    /// Watermark column (if configured)
75    pub watermark_column: Option<String>,
76    /// Connector options
77    pub options: HashMap<String, String>,
78}
79
80/// Information about a registered sink
81#[derive(Debug, Clone)]
82pub struct SinkInfo {
83    /// Sink name
84    pub name: String,
85    /// Source table or query name
86    pub from: String,
87    /// Connector options
88    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/// Result of planning a streaming statement
116#[derive(Debug)]
117#[allow(clippy::large_enum_variant)]
118pub enum StreamingPlan {
119    /// Source registration (DDL)
120    RegisterSource(SourceInfo),
121
122    /// Sink registration (DDL)
123    RegisterSink(SinkInfo),
124
125    /// Query plan with streaming configurations
126    Query(QueryPlan),
127
128    /// Standard SQL statement (pass-through to DataFusion)
129    Standard(Box<Statement>),
130
131    /// Lookup table registration (DDL)
132    RegisterLookupTable(LookupTableInfo),
133
134    /// Drop a lookup table
135    DropLookupTable {
136        /// Name of the lookup table to drop.
137        name: String,
138    },
139}
140
141/// A query plan with streaming operator configurations
142#[derive(Debug)]
143pub struct QueryPlan {
144    /// Optional name for the continuous query
145    pub name: Option<String>,
146    /// Window configuration if the query has windowed aggregation
147    pub window_config: Option<WindowOperatorConfig>,
148    /// Join configuration(s) if the query has joins (one per join step)
149    pub join_config: Option<Vec<JoinOperatorConfig>>,
150    /// ORDER BY configuration if the query has ordering
151    pub order_config: Option<OrderOperatorConfig>,
152    /// Analytic window function configuration (LAG/LEAD/etc.)
153    pub analytic_config: Option<AnalyticWindowConfig>,
154    /// HAVING clause filter configuration
155    pub having_config: Option<HavingFilterConfig>,
156    /// Window frame configuration (ROWS BETWEEN / RANGE BETWEEN)
157    pub frame_config: Option<WindowFrameConfig>,
158    /// Emit strategy
159    pub emit_clause: Option<EmitClause>,
160    /// The underlying SQL statement
161    pub statement: Box<Statement>,
162}
163
164/// Exact one-call admission for a database-certified state-backed join.
165///
166/// The certificate is matched against the planner's independent parse, so it
167/// cannot authorize another relation pair, key mapping, or join type.
168#[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    /// Construct an exact INNER or LEFT state-backed join admission.
179    ///
180    /// # Errors
181    /// Returns an error for empty relations/keys or a mismatched key arity.
182    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    /// Creates a new streaming planner
235    #[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    /// Plans a streaming statement.
246    ///
247    /// # Errors
248    ///
249    /// Returns `PlanningError` if the statement cannot be planned.
250    pub fn plan(&mut self, statement: &StreamingStatement) -> Result<StreamingPlan, PlanningError> {
251        self.plan_internal(statement, None)
252    }
253
254    /// Plan one query whose unbounded join shape has already been certified by
255    /// the database as a durable state-backed incremental join.
256    ///
257    /// This does not relax multi-way, join-type, temporal, or implicit-join
258    /// validation. Raw stream callers must use [`Self::plan`].
259    ///
260    /// # Errors
261    /// Returns `PlanningError` if the statement fails the remaining streaming checks.
262    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                // 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            options: source.with_options.clone(),
377        };
378
379        // Register the source
380        self.sources.insert(name, info.clone());
381
382        Ok(StreamingPlan::RegisterSource(info))
383    }
384
385    /// Plans a CREATE SINK statement.
386    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        // Check for existing sink
393        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        // Determine the source
401        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        // Register the sink
413        self.sinks.insert(name, info.clone());
414
415        Ok(StreamingPlan::RegisterSink(info))
416    }
417
418    /// Plans a CREATE CONTINUOUS QUERY statement.
419    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        // The query inside should be a standard SELECT
427        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        // Analyze the query for streaming features
437        let query_plan = self.analyze_query(&stmt, emit_clause, state_backed_join)?;
438
439        // Keep planner classification in sync with catalog rollback/drop. A windowed query is the
440        // only query shape for which the planner retains classification after planning.
441        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    /// Plans a standard SQL statement.
462    #[allow(clippy::unused_self)] // Will use planner state for plan optimization
463    fn plan_standard_statement(&self, stmt: &Statement) -> Result<StreamingPlan, PlanningError> {
464        // Check if it's a query that might have streaming features
465        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                // Check for window functions in GROUP BY
474                let window_function = Self::extract_window_from_select(select);
475
476                // Check for joins (multi-way)
477                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                // Check for ORDER BY
486                let order_analysis = analyze_order_by(stmt);
487                let order_config = OrderOperatorConfig::from_analysis(&order_analysis)
488                    .map_err(PlanningError::InvalidQuery)?;
489
490                // Check for analytic functions (LAG/LEAD/etc.)
491                let analytic_analysis = analyze_analytic_functions(stmt);
492                let analytic_config =
493                    analytic_analysis.map(|a| AnalyticWindowConfig::from_analysis(&a));
494
495                // Check for HAVING clause
496                let agg_analysis = analyze_aggregates(stmt);
497                let having_config = agg_analysis.having_expr.map(HavingFilterConfig::new);
498
499                // Check for window frame functions (ROWS BETWEEN / RANGE BETWEEN)
500                let frame_analysis = analyze_window_frames(stmt);
501                let frame_config = frame_analysis
502                    .as_ref()
503                    .map(WindowFrameConfig::from_analysis);
504
505                // Validate: reject UNBOUNDED FOLLOWING (streaming can't buffer infinite future)
506                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        // Pass through standard SQL
552        Ok(StreamingPlan::Standard(Box::new(stmt.clone())))
553    }
554
555    /// Analyzes a query for streaming features.
556    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                // Extract window function
573                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                    // Apply emit clause if present
578                    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                // Extract join info (multi-way)
588                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        // Extract ORDER BY info
598        let order_analysis = analyze_order_by(stmt);
599        analysis.order_config = OrderOperatorConfig::from_analysis(&order_analysis)
600            .map_err(PlanningError::InvalidQuery)?;
601
602        // Extract analytic function info (LAG/LEAD/etc.)
603        if let Some(analytic) = analyze_analytic_functions(stmt) {
604            analysis.analytic_config = Some(AnalyticWindowConfig::from_analysis(&analytic));
605        }
606
607        // Extract HAVING clause
608        let agg_analysis = analyze_aggregates(stmt);
609        analysis.having_config = agg_analysis.having_expr.map(HavingFilterConfig::new);
610
611        // Extract window frame functions (ROWS BETWEEN / RANGE BETWEEN)
612        if let Some(frame_analysis) = analyze_window_frames(stmt) {
613            // Validate: reject UNBOUNDED FOLLOWING
614            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    /// Extracts window function from a SELECT.
629    fn extract_window_from_select(select: &sqlparser::ast::Select) -> Option<WindowFunction> {
630        // Check GROUP BY for window functions
631        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    /// Plans a CREATE LOOKUP TABLE statement.
647    fn plan_create_lookup_table(
648        &mut self,
649        lt: &CreateLookupTableStatement,
650    ) -> Result<StreamingPlan, PlanningError> {
651        let name = object_name_to_string(&lt.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(&lt.with_options).map_err(|e| {
667            PlanningError::InvalidQuery(format!("Invalid lookup table properties: {e}"))
668        })?;
669
670        // Compute Arrow schema from column definitions
671        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    /// Plans a DROP LOOKUP TABLE statement.
701    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    /// Gets a registered source by name.
721    #[must_use]
722    pub fn get_source(&self, name: &str) -> Option<&SourceInfo> {
723        self.sources.get(name)
724    }
725
726    /// Gets a registered sink by name.
727    #[must_use]
728    pub fn get_sink(&self, name: &str) -> Option<&SinkInfo> {
729        self.sinks.get(name)
730    }
731
732    /// Lists all registered sources.
733    #[must_use]
734    pub fn list_sources(&self) -> Vec<&SourceInfo> {
735        self.sources.values().collect()
736    }
737
738    /// Lists all registered sinks.
739    #[must_use]
740    pub fn list_sinks(&self) -> Vec<&SinkInfo> {
741        self.sinks.values().collect()
742    }
743
744    /// Gets a registered lookup table by name.
745    #[must_use]
746    pub fn get_lookup_table(&self, name: &str) -> Option<&LookupTableInfo> {
747        self.lookup_tables.get(name)
748    }
749
750    /// Lists all registered lookup tables.
751    #[must_use]
752    pub fn list_lookup_tables(&self) -> Vec<&LookupTableInfo> {
753        self.lookup_tables.values().collect()
754    }
755
756    /// Returns a clone of the lookup tables map for optimizer rule construction.
757    #[must_use]
758    pub fn lookup_tables_cloned(&self) -> HashMap<String, LookupTableInfo> {
759        self.lookup_tables.clone()
760    }
761
762    /// Converts a query plan's SQL statement into a `DataFusion`
763    /// `LogicalPlan`. Window UDFs (TUMBLE, HOP, SESSION) must be registered
764    /// on `ctx` via
765    /// [`register_streaming_functions`](crate::datafusion::register_streaming_functions)
766    /// for windowed queries to resolve correctly.
767    ///
768    /// # Errors
769    ///
770    /// Returns `PlanningError` if `DataFusion` cannot create the logical plan.
771    #[allow(clippy::unused_self)] // Method will use planner state for plan optimization
772    pub async fn to_logical_plan(
773        &self,
774        plan: &QueryPlan,
775        ctx: &SessionContext,
776    ) -> Result<LogicalPlan, PlanningError> {
777        // Convert the AST statement back to SQL and let DataFusion re-parse
778        // it with its own sqlparser version. This avoids version mismatches
779        // between our sqlparser (0.60) and DataFusion's (0.59).
780        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/// Intermediate query analysis result
795#[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
806/// Helper to convert `ObjectName` to String
807fn 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
814/// Fail closed on stream-join shapes whose state/output semantics are not implemented.
815fn 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/// Planning errors
861#[derive(Debug, thiserror::Error)]
862pub enum PlanningError {
863    /// Unsupported SQL feature
864    UnsupportedSql(String),
865
866    /// Invalid query
867    InvalidQuery(String),
868
869    /// Source not found
870    SourceNotFound(String),
871
872    /// Sink not found
873    SinkNotFound(String),
874
875    /// `DataFusion` error during logical plan creation (translated on display)
876    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        // First source
934        let statements =
935            StreamingParser::parse_sql("CREATE SOURCE events (id INT, name VARCHAR)").unwrap();
936        planner.plan(&statements[0]).unwrap();
937
938        // Duplicate should fail
939        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        // First source
948        let statements =
949            StreamingParser::parse_sql("CREATE SOURCE events (id INT, name VARCHAR)").unwrap();
950        planner.plan(&statements[0]).unwrap();
951
952        // IF NOT EXISTS should succeed
953        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        // First source
965        let statements =
966            StreamingParser::parse_sql("CREATE SOURCE events (id INT, name VARCHAR)").unwrap();
967        planner.plan(&statements[0]).unwrap();
968
969        // OR REPLACE should succeed
970        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        // Create sources
1016        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        // Create sinks
1022        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        // Two windowed views register as bounded inputs.
1250        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        // Closed-window batches are not an ordering/alignment contract between two independent
1260        // streams. A per-cycle join can silently miss rows when the batches arrive in different
1261        // cycles, so this shape remains unbounded and must fail closed.
1262        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        // HAVING without window function still produces a Query plan
1374        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    // -- Multi-way join planner tests --
1432
1433    #[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        // A lookup-table final step does not make a multi-way operator atomic; users must expose
1473        // the interval result as a named stage before enriching it.
1474        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(_) => {} // No join → pass-through
1503            _ => panic!("Expected Standard plan for simple SELECT"),
1504        }
1505    }
1506
1507    // -- Window Frame planner tests --
1508
1509    #[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(_) => {} // No frame → pass-through
1558            _ => 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}