Skip to main content

laminar_sql/parser/
mod.rs

1//! SQL parser with streaming extensions.
2//!
3//! Routes streaming DDL (CREATE SOURCE/SINK/CONTINUOUS QUERY) to custom
4//! parsers that use sqlparser primitives. Routes standard SQL to sqlparser
5//! with `GenericDialect`.
6
7pub mod aggregation_parser;
8pub mod analytic_parser;
9mod continuous_query_parser;
10mod declare_parser;
11pub(crate) mod dialect;
12mod emit_parser;
13pub mod join_parser;
14mod late_data_parser;
15/// Parser for CREATE/DROP LOOKUP TABLE DDL statements
16pub mod lookup_table;
17pub mod order_analyzer;
18mod sink_parser;
19mod source_parser;
20mod statements;
21mod subscribe_parser;
22mod tokenizer;
23mod window_rewriter;
24
25pub use lookup_table::CreateLookupTableStatement;
26pub use statements::{
27    AlterSourceOperation, CreateSinkStatement, CreateSourceStatement, EmitClause, EmitStrategy,
28    FormatSpec, LateDataClause, ShowCommand, SinkFrom, StreamingStatement, SubscribeStatement,
29    WatermarkDef, WindowFunction,
30};
31pub use window_rewriter::WindowRewriter;
32
33use dialect::LaminarDialect;
34use tokenizer::{detect_streaming_ddl, StreamingDdlKind};
35
36/// Parses SQL with streaming extensions.
37///
38/// Routes streaming DDL to custom parsers that use sqlparser's `Parser` API
39/// for structured parsing. Standard SQL is delegated to sqlparser directly.
40///
41/// # Errors
42///
43/// Returns `ParseError` if the SQL syntax is invalid.
44pub fn parse_streaming_sql(sql: &str) -> Result<Vec<StreamingStatement>, ParseError> {
45    StreamingParser::parse_sql(sql).map_err(ParseError::SqlParseError)
46}
47
48/// Parser for streaming SQL extensions.
49///
50/// Provides static methods for parsing streaming SQL statements.
51/// Uses sqlparser's `Parser` API internally for structured parsing
52/// of identifiers, data types, expressions, and queries.
53pub struct StreamingParser;
54
55impl StreamingParser {
56    /// Parse a SQL string with streaming extensions.
57    ///
58    /// Tokenizes the input to detect statement type, then routes to the
59    /// appropriate parser:
60    /// - CREATE SOURCE → `source_parser`
61    /// - CREATE SINK → `sink_parser`
62    /// - CREATE CONTINUOUS QUERY → `continuous_query_parser`
63    /// - Everything else → `sqlparser::parser::Parser`
64    ///
65    /// # Errors
66    ///
67    /// Returns `ParserError` if the SQL syntax is invalid.
68    pub fn parse_sql(sql: &str) -> Result<Vec<StreamingStatement>, sqlparser::parser::ParserError> {
69        let sql_trimmed = sql.trim();
70        if sql_trimmed.is_empty() {
71            return Err(sqlparser::parser::ParserError::ParserError(
72                "Empty SQL statement".to_string(),
73            ));
74        }
75
76        let dialect = LaminarDialect::default();
77
78        let tokens = sqlparser::tokenizer::Tokenizer::new(&dialect, sql_trimmed)
79            .tokenize_with_location()
80            .map_err(|e| {
81                sqlparser::parser::ParserError::ParserError(format!("Tokenization error: {e}"))
82            })?;
83        let kind = detect_streaming_ddl(&tokens);
84        if kind == StreamingDdlKind::None {
85            return parse_standard_or_temporal_sql(&dialect, sql_trimmed, &tokens);
86        }
87
88        let mut parser =
89            sqlparser::parser::Parser::new(&dialect).with_tokens_with_locations(tokens);
90        let statement = match kind {
91            StreamingDdlKind::CreateSource { .. } => StreamingStatement::CreateSource(Box::new(
92                source_parser::parse_create_source(&mut parser)
93                    .map_err(parse_error_to_parser_error)?,
94            )),
95            StreamingDdlKind::CreateSink { .. } => StreamingStatement::CreateSink(Box::new(
96                sink_parser::parse_create_sink(&mut parser).map_err(parse_error_to_parser_error)?,
97            )),
98            StreamingDdlKind::CreateContinuousQuery { .. } => {
99                continuous_query_parser::parse_continuous_query(&mut parser)
100                    .map_err(parse_error_to_parser_error)?
101            }
102            StreamingDdlKind::DropSource { .. } => {
103                parse_drop_source(&mut parser).map_err(parse_error_to_parser_error)?
104            }
105            StreamingDdlKind::DropSink { .. } => {
106                parse_drop_sink(&mut parser).map_err(parse_error_to_parser_error)?
107            }
108            StreamingDdlKind::DropMaterializedView { .. } => {
109                parse_drop_materialized_view(&mut parser).map_err(parse_error_to_parser_error)?
110            }
111            StreamingDdlKind::ShowSources => StreamingStatement::Show(ShowCommand::Sources),
112            StreamingDdlKind::ShowSinks => StreamingStatement::Show(ShowCommand::Sinks),
113            StreamingDdlKind::ShowQueries => StreamingStatement::Show(ShowCommand::Queries),
114            StreamingDdlKind::ShowMaterializedViews => {
115                StreamingStatement::Show(ShowCommand::MaterializedViews)
116            }
117            StreamingDdlKind::DescribeSource => {
118                parse_describe(&mut parser).map_err(parse_error_to_parser_error)?
119            }
120            StreamingDdlKind::ExplainStreaming => {
121                parse_explain(&mut parser, sql_trimmed).map_err(parse_error_to_parser_error)?
122            }
123            StreamingDdlKind::CreateMaterializedView { .. } => {
124                parse_create_materialized_view(&mut parser, sql_trimmed)
125                    .map_err(parse_error_to_parser_error)?
126            }
127            StreamingDdlKind::CreateStream { .. } => parse_create_stream(&mut parser, sql_trimmed)
128                .map_err(parse_error_to_parser_error)?,
129            StreamingDdlKind::DropStream { .. } => {
130                parse_drop_stream(&mut parser).map_err(parse_error_to_parser_error)?
131            }
132            StreamingDdlKind::ShowStreams => StreamingStatement::Show(ShowCommand::Streams),
133            StreamingDdlKind::ShowTables => StreamingStatement::Show(ShowCommand::Tables),
134            StreamingDdlKind::CreateLookupTable { .. } => {
135                StreamingStatement::CreateLookupTable(Box::new(
136                    lookup_table::parse_create_lookup_table(&mut parser)
137                        .map_err(parse_error_to_parser_error)?,
138                ))
139            }
140            StreamingDdlKind::DropLookupTable { .. } => {
141                let (name, if_exists) = lookup_table::parse_drop_lookup_table(&mut parser)
142                    .map_err(parse_error_to_parser_error)?;
143                StreamingStatement::DropLookupTable { name, if_exists }
144            }
145            StreamingDdlKind::AlterSource => {
146                parse_alter_source(&mut parser).map_err(parse_error_to_parser_error)?
147            }
148            StreamingDdlKind::ShowCheckpointStatus => {
149                StreamingStatement::Show(ShowCommand::CheckpointStatus)
150            }
151            StreamingDdlKind::ShowCreateSource => {
152                parse_show_create_source(&mut parser).map_err(parse_error_to_parser_error)?
153            }
154            StreamingDdlKind::ShowCreateSink => {
155                parse_show_create_sink(&mut parser).map_err(parse_error_to_parser_error)?
156            }
157            StreamingDdlKind::Checkpoint => StreamingStatement::Checkpoint,
158            StreamingDdlKind::Subscribe => StreamingStatement::Subscribe(Box::new(
159                subscribe_parser::parse_subscribe(&mut parser)
160                    .map_err(parse_error_to_parser_error)?,
161            )),
162            StreamingDdlKind::DeclareCursor => declare_parser::parse_declare_cursor(&mut parser)
163                .map_err(parse_error_to_parser_error)?,
164            StreamingDdlKind::RestoreCheckpoint => {
165                parse_restore_checkpoint(&mut parser).map_err(parse_error_to_parser_error)?
166            }
167            StreamingDdlKind::None => unreachable!("handled before parser construction"),
168        };
169        Ok(vec![statement])
170    }
171
172    /// Check if an expression contains a window function.
173    #[must_use]
174    pub fn has_window_function(expr: &sqlparser::ast::Expr) -> bool {
175        match expr {
176            sqlparser::ast::Expr::Function(func) => {
177                if let Some(name) = func.name.0.last() {
178                    let func_name = name.to_string().to_uppercase();
179                    matches!(func_name.as_str(), "TUMBLE" | "HOP" | "SESSION")
180                } else {
181                    false
182                }
183            }
184            _ => false,
185        }
186    }
187
188    /// Parse EMIT clause from SQL string.
189    ///
190    /// # Errors
191    ///
192    /// Returns `ParseError::StreamingError` if the EMIT clause syntax is invalid.
193    pub fn parse_emit_clause(sql: &str) -> Result<Option<EmitClause>, ParseError> {
194        emit_parser::parse_emit_clause_from_sql(sql)
195    }
196
197    /// Parse late data handling clause from SQL string.
198    ///
199    /// # Errors
200    ///
201    /// Returns `ParseError::StreamingError` if the clause syntax is invalid.
202    pub fn parse_late_data_clause(sql: &str) -> Result<Option<LateDataClause>, ParseError> {
203        late_data_parser::parse_late_data_clause_from_sql(sql)
204    }
205}
206
207fn parse_standard_or_temporal_sql(
208    dialect: &LaminarDialect,
209    sql: &str,
210    tokens: &[sqlparser::tokenizer::TokenWithSpan],
211) -> Result<Vec<StreamingStatement>, sqlparser::parser::ParserError> {
212    if let Some(parsed) =
213        join_parser::parse_temporal_probe_query(tokens).map_err(parse_error_to_parser_error)?
214    {
215        return Ok(vec![StreamingStatement::TemporalProbeQuery {
216            statement: Box::new(parsed.statement),
217            analysis: Box::new(parsed.analysis),
218        }]);
219    }
220
221    let statements = sqlparser::parser::Parser::parse_sql(dialect, sql)?;
222    Ok(statements
223        .into_iter()
224        .map(convert_standard_statement)
225        .collect())
226}
227
228/// Convert `ParseError` to `ParserError` for backward compatibility.
229fn parse_error_to_parser_error(e: ParseError) -> sqlparser::parser::ParserError {
230    match e {
231        ParseError::SqlParseError(pe) => pe,
232        ParseError::StreamingError(msg) => sqlparser::parser::ParserError::ParserError(msg),
233        ParseError::WindowError(msg) => {
234            sqlparser::parser::ParserError::ParserError(format!("Window error: {msg}"))
235        }
236        ParseError::ValidationError(msg) => {
237            sqlparser::parser::ParserError::ParserError(format!("Validation error: {msg}"))
238        }
239    }
240}
241
242/// Parse a RESTORE FROM CHECKPOINT statement.
243///
244/// Syntax: `RESTORE FROM CHECKPOINT <id>`
245///
246/// # Errors
247///
248/// Returns `ParseError` if the statement syntax is invalid.
249fn parse_restore_checkpoint(
250    parser: &mut sqlparser::parser::Parser,
251) -> Result<StreamingStatement, ParseError> {
252    // Consume RESTORE
253    tokenizer::expect_custom_keyword(parser, "RESTORE")?;
254    // Consume FROM
255    if !parser.parse_keyword(sqlparser::keywords::Keyword::FROM) {
256        return Err(ParseError::StreamingError(
257            "Expected FROM after RESTORE".to_string(),
258        ));
259    }
260    // Consume CHECKPOINT
261    tokenizer::expect_custom_keyword(parser, "CHECKPOINT")?;
262    // Parse checkpoint ID (numeric literal)
263    let token = parser.next_token();
264    match &token.token {
265        sqlparser::tokenizer::Token::Number(n, _) => {
266            let id: u64 = n
267                .parse()
268                .map_err(|_| ParseError::StreamingError(format!("Invalid checkpoint ID: {n}")))?;
269            Ok(StreamingStatement::RestoreCheckpoint { checkpoint_id: id })
270        }
271        other => Err(ParseError::StreamingError(format!(
272            "Expected checkpoint ID (number), found {other}"
273        ))),
274    }
275}
276
277/// Convert a standard sqlparser statement to a `StreamingStatement`.
278///
279/// Detects INSERT INTO statements and converts them to the streaming
280/// `InsertInto` variant. All other statements are wrapped as `Standard`.
281fn convert_standard_statement(stmt: sqlparser::ast::Statement) -> StreamingStatement {
282    if let sqlparser::ast::Statement::Insert(insert) = &stmt {
283        // Extract table name from TableObject
284        if let sqlparser::ast::TableObject::TableName(ref name) = insert.table {
285            let table_name = name.clone();
286            let columns = insert.columns.clone();
287
288            // Try to extract VALUES rows from source query
289            if let Some(ref source) = insert.source {
290                if let sqlparser::ast::SetExpr::Values(ref values) = *source.body {
291                    let rows = values.rows.clone();
292                    return StreamingStatement::InsertInto {
293                        table_name,
294                        columns,
295                        values: rows,
296                    };
297                }
298            }
299        }
300    }
301    StreamingStatement::Standard(Box::new(stmt))
302}
303
304/// Parse a DROP SOURCE statement.
305///
306/// Syntax: `DROP SOURCE [IF EXISTS] name [CASCADE]`
307///
308/// # Errors
309///
310/// Returns `ParseError` if the statement syntax is invalid.
311fn parse_drop_source(
312    parser: &mut sqlparser::parser::Parser,
313) -> Result<StreamingStatement, ParseError> {
314    parser
315        .expect_keyword(sqlparser::keywords::Keyword::DROP)
316        .map_err(ParseError::SqlParseError)?;
317    tokenizer::expect_custom_keyword(parser, "SOURCE")?;
318    let if_exists = parser.parse_keywords(&[
319        sqlparser::keywords::Keyword::IF,
320        sqlparser::keywords::Keyword::EXISTS,
321    ]);
322    let name = parser
323        .parse_object_name(false)
324        .map_err(ParseError::SqlParseError)?;
325    let cascade = parser.parse_keyword(sqlparser::keywords::Keyword::CASCADE);
326    Ok(StreamingStatement::DropSource {
327        name,
328        if_exists,
329        cascade,
330    })
331}
332
333/// Parse a DROP SINK statement.
334///
335/// Syntax: `DROP SINK [IF EXISTS] name [CASCADE]`
336///
337/// # Errors
338///
339/// Returns `ParseError` if the statement syntax is invalid.
340fn parse_drop_sink(
341    parser: &mut sqlparser::parser::Parser,
342) -> Result<StreamingStatement, ParseError> {
343    parser
344        .expect_keyword(sqlparser::keywords::Keyword::DROP)
345        .map_err(ParseError::SqlParseError)?;
346    tokenizer::expect_custom_keyword(parser, "SINK")?;
347    let if_exists = parser.parse_keywords(&[
348        sqlparser::keywords::Keyword::IF,
349        sqlparser::keywords::Keyword::EXISTS,
350    ]);
351    let name = parser
352        .parse_object_name(false)
353        .map_err(ParseError::SqlParseError)?;
354    let cascade = parser.parse_keyword(sqlparser::keywords::Keyword::CASCADE);
355    Ok(StreamingStatement::DropSink {
356        name,
357        if_exists,
358        cascade,
359    })
360}
361
362/// Parse a DROP MATERIALIZED VIEW statement.
363///
364/// Syntax: `DROP MATERIALIZED VIEW [IF EXISTS] name [CASCADE]`
365///
366/// # Errors
367///
368/// Returns `ParseError` if the statement syntax is invalid.
369fn parse_drop_materialized_view(
370    parser: &mut sqlparser::parser::Parser,
371) -> Result<StreamingStatement, ParseError> {
372    parser
373        .expect_keyword(sqlparser::keywords::Keyword::DROP)
374        .map_err(ParseError::SqlParseError)?;
375    parser
376        .expect_keyword(sqlparser::keywords::Keyword::MATERIALIZED)
377        .map_err(ParseError::SqlParseError)?;
378    parser
379        .expect_keyword(sqlparser::keywords::Keyword::VIEW)
380        .map_err(ParseError::SqlParseError)?;
381    let if_exists = parser.parse_keywords(&[
382        sqlparser::keywords::Keyword::IF,
383        sqlparser::keywords::Keyword::EXISTS,
384    ]);
385    let name = parser
386        .parse_object_name(false)
387        .map_err(ParseError::SqlParseError)?;
388    let cascade = parser.parse_keyword(sqlparser::keywords::Keyword::CASCADE);
389    Ok(StreamingStatement::DropMaterializedView {
390        name,
391        if_exists,
392        cascade,
393    })
394}
395
396/// Parse a CREATE STREAM statement.
397///
398/// Syntax: `CREATE [OR REPLACE] STREAM [IF NOT EXISTS] name AS <select_query> [EMIT <strategy>]`
399///
400/// # Errors
401///
402/// Returns `ParseError` if the statement syntax is invalid.
403fn parse_create_stream(
404    parser: &mut sqlparser::parser::Parser,
405    original_sql: &str,
406) -> Result<StreamingStatement, ParseError> {
407    parser
408        .expect_keyword(sqlparser::keywords::Keyword::CREATE)
409        .map_err(ParseError::SqlParseError)?;
410
411    let or_replace = parser.parse_keywords(&[
412        sqlparser::keywords::Keyword::OR,
413        sqlparser::keywords::Keyword::REPLACE,
414    ]);
415
416    tokenizer::expect_custom_keyword(parser, "STREAM")?;
417
418    let if_not_exists = parser.parse_keywords(&[
419        sqlparser::keywords::Keyword::IF,
420        sqlparser::keywords::Keyword::NOT,
421        sqlparser::keywords::Keyword::EXISTS,
422    ]);
423
424    let name = parser
425        .parse_object_name(false)
426        .map_err(ParseError::SqlParseError)?;
427
428    parser
429        .expect_keyword(sqlparser::keywords::Keyword::AS)
430        .map_err(ParseError::SqlParseError)?;
431
432    // Collect remaining tokens, then peel off the optional trailing
433    // `WITH (...)` before splitting query / EMIT.
434    let remaining = collect_remaining_tokens(parser);
435    let (head_tokens, with_tokens) = split_off_trailing_with(&remaining);
436    let (query_tokens, emit_tokens) = split_at_emit(&head_tokens);
437    let raw_query_sql = query_body_sql(
438        original_sql,
439        &query_tokens,
440        &emit_tokens,
441        with_tokens.as_deref(),
442    );
443
444    let stream_dialect = LaminarDialect::default();
445
446    let (query_stmt, normalized_temporal_sql) = if query_tokens.is_empty() {
447        return Err(ParseError::StreamingError(
448            "Expected SELECT query after AS".to_string(),
449        ));
450    } else if let Some(parsed) = join_parser::parse_temporal_probe_query(&query_tokens)? {
451        (
452            StreamingStatement::TemporalProbeQuery {
453                statement: Box::new(parsed.statement),
454                analysis: Box::new(parsed.analysis),
455            },
456            Some(parsed.normalized_sql),
457        )
458    } else {
459        let mut query_parser = sqlparser::parser::Parser::new(&stream_dialect)
460            .with_tokens_with_locations(query_tokens);
461        let query = query_parser
462            .parse_query()
463            .map_err(ParseError::SqlParseError)?;
464        (
465            StreamingStatement::Standard(Box::new(sqlparser::ast::Statement::Query(query))),
466            None,
467        )
468    };
469    let query_sql = normalized_temporal_sql.unwrap_or(raw_query_sql);
470
471    let emit_clause = if emit_tokens.is_empty() {
472        None
473    } else {
474        let mut emit_parser =
475            sqlparser::parser::Parser::new(&stream_dialect).with_tokens_with_locations(emit_tokens);
476        emit_parser::parse_emit_clause(&mut emit_parser)?
477    };
478
479    let retention_bytes = match with_tokens {
480        None => None,
481        Some(tokens) => {
482            let mut with_parser =
483                sqlparser::parser::Parser::new(&stream_dialect).with_tokens_with_locations(tokens);
484            let opts = tokenizer::parse_with_options(&mut with_parser)?;
485            extract_retention_bytes(&opts)?
486        }
487    };
488
489    Ok(StreamingStatement::CreateStream {
490        name,
491        query: Box::new(query_stmt),
492        emit_clause,
493        or_replace,
494        if_not_exists,
495        query_sql,
496        retention_bytes,
497    })
498}
499
500/// Split off a trailing `WITH (` at depth 0. CTE-style `WITH ident AS (...)`
501/// is ignored because it's followed by an identifier, not `(`.
502fn split_off_trailing_with(
503    tokens: &[sqlparser::tokenizer::TokenWithSpan],
504) -> (
505    Vec<sqlparser::tokenizer::TokenWithSpan>,
506    Option<Vec<sqlparser::tokenizer::TokenWithSpan>>,
507) {
508    use sqlparser::tokenizer::Token;
509    let mut depth: i32 = 0;
510    let mut last_with_idx: Option<usize> = None;
511    for (i, t) in tokens.iter().enumerate() {
512        match &t.token {
513            Token::LParen => depth += 1,
514            Token::RParen => depth -= 1,
515            Token::Word(w) if depth == 0 && w.value.eq_ignore_ascii_case("WITH") => {
516                if matches!(tokens.get(i + 1).map(|t| &t.token), Some(Token::LParen)) {
517                    last_with_idx = Some(i);
518                }
519            }
520            _ => {}
521        }
522    }
523    match last_with_idx {
524        None => (tokens.to_vec(), None),
525        Some(i) => {
526            let mut head = tokens[..i].to_vec();
527            head.push(sqlparser::tokenizer::TokenWithSpan {
528                token: Token::EOF,
529                span: sqlparser::tokenizer::Span::empty(),
530            });
531            let tail = tokens[i..].to_vec();
532            (head, Some(tail))
533        }
534    }
535}
536
537/// Unknown keys are rejected so typos surface immediately.
538fn extract_retention_bytes(
539    opts: &std::collections::HashMap<String, String>,
540) -> Result<Option<u64>, ParseError> {
541    for key in opts.keys() {
542        if !key.eq_ignore_ascii_case("retain_history") {
543            return Err(ParseError::StreamingError(format!(
544                "unknown CREATE STREAM option '{key}' (expected RETAIN_HISTORY)"
545            )));
546        }
547    }
548    match opts
549        .iter()
550        .find(|(k, _)| k.eq_ignore_ascii_case("retain_history"))
551    {
552        None => Ok(None),
553        Some((_, v)) => {
554            let bytes = lookup_table::ByteSize::parse(v)?;
555            Ok(Some(bytes.as_bytes()))
556        }
557    }
558}
559
560/// Slice `original_sql` from the first query token to whichever trailing
561/// clause comes next (`EMIT` or `WITH`), or end of input if neither is
562/// present. Preserves custom streaming syntax that sqlparser's AST would
563/// drop. Falls back to joining token text if spans are empty.
564pub(super) fn query_body_sql(
565    original_sql: &str,
566    query_tokens: &[sqlparser::tokenizer::TokenWithSpan],
567    emit_tokens: &[sqlparser::tokenizer::TokenWithSpan],
568    with_tokens: Option<&[sqlparser::tokenizer::TokenWithSpan]>,
569) -> String {
570    use sqlparser::tokenizer::Token;
571
572    let from_spans = || -> Option<String> {
573        let first = query_tokens
574            .iter()
575            .find(|t| !matches!(t.token, Token::EOF))?;
576        let start = location_to_byte_offset(original_sql, first.span.start)?;
577        let trailer_starts = [emit_tokens.first(), with_tokens.and_then(|t| t.first())]
578            .into_iter()
579            .flatten()
580            .filter_map(|t| location_to_byte_offset(original_sql, t.span.start));
581        let end = trailer_starts.min().unwrap_or(original_sql.len());
582        let slice = original_sql.get(start..end)?;
583        Some(
584            slice
585                .trim_end_matches(|c: char| c.is_whitespace() || c == ';')
586                .to_string(),
587        )
588    };
589
590    from_spans().unwrap_or_else(|| {
591        query_tokens
592            .iter()
593            .take_while(|t| !matches!(t.token, Token::EOF))
594            .map(|t| t.token.to_string())
595            .collect::<Vec<_>>()
596            .join(" ")
597    })
598}
599
600/// Parse a DROP STREAM statement.
601///
602/// Syntax: `DROP STREAM [IF EXISTS] name [CASCADE]`
603///
604/// # Errors
605///
606/// Returns `ParseError` if the statement syntax is invalid.
607fn parse_drop_stream(
608    parser: &mut sqlparser::parser::Parser,
609) -> Result<StreamingStatement, ParseError> {
610    parser
611        .expect_keyword(sqlparser::keywords::Keyword::DROP)
612        .map_err(ParseError::SqlParseError)?;
613    tokenizer::expect_custom_keyword(parser, "STREAM")?;
614    let if_exists = parser.parse_keywords(&[
615        sqlparser::keywords::Keyword::IF,
616        sqlparser::keywords::Keyword::EXISTS,
617    ]);
618    let name = parser
619        .parse_object_name(false)
620        .map_err(ParseError::SqlParseError)?;
621    let cascade = parser.parse_keyword(sqlparser::keywords::Keyword::CASCADE);
622    Ok(StreamingStatement::DropStream {
623        name,
624        if_exists,
625        cascade,
626    })
627}
628
629/// Parse a DESCRIBE statement.
630///
631/// Syntax: `DESCRIBE [EXTENDED] name`
632/// Parse an ALTER SOURCE statement.
633///
634/// Syntax:
635/// - `ALTER SOURCE name ADD COLUMN col_name data_type`
636/// - `ALTER SOURCE name SET ('key' = 'value', ...)`
637///
638/// # Errors
639///
640/// Returns `ParseError` if the statement syntax is invalid.
641fn parse_alter_source(
642    parser: &mut sqlparser::parser::Parser,
643) -> Result<StreamingStatement, ParseError> {
644    parser
645        .expect_keyword(sqlparser::keywords::Keyword::ALTER)
646        .map_err(ParseError::SqlParseError)?;
647    tokenizer::expect_custom_keyword(parser, "SOURCE")?;
648    let name = parser
649        .parse_object_name(false)
650        .map_err(ParseError::SqlParseError)?;
651
652    // Determine operation: ADD COLUMN or SET
653    if parser.parse_keywords(&[
654        sqlparser::keywords::Keyword::ADD,
655        sqlparser::keywords::Keyword::COLUMN,
656    ]) {
657        // ALTER SOURCE name ADD COLUMN col_name data_type
658        let col_name = parser
659            .parse_identifier()
660            .map_err(ParseError::SqlParseError)?;
661        let data_type = parser
662            .parse_data_type()
663            .map_err(ParseError::SqlParseError)?;
664        let column_def = sqlparser::ast::ColumnDef {
665            name: col_name,
666            data_type,
667            options: vec![],
668        };
669        Ok(StreamingStatement::AlterSource {
670            name,
671            operation: statements::AlterSourceOperation::AddColumn { column_def },
672        })
673    } else if parser.parse_keyword(sqlparser::keywords::Keyword::SET) {
674        // ALTER SOURCE name SET ('key' = 'value', ...)
675        parser
676            .expect_token(&sqlparser::tokenizer::Token::LParen)
677            .map_err(ParseError::SqlParseError)?;
678        #[allow(clippy::disallowed_types)] // cold path: SQL parsing
679        let mut properties = std::collections::HashMap::new();
680        loop {
681            let key = parser
682                .parse_literal_string()
683                .map_err(ParseError::SqlParseError)?;
684            parser
685                .expect_token(&sqlparser::tokenizer::Token::Eq)
686                .map_err(ParseError::SqlParseError)?;
687            let value = parser
688                .parse_literal_string()
689                .map_err(ParseError::SqlParseError)?;
690            properties.insert(key, value);
691            if !parser.consume_token(&sqlparser::tokenizer::Token::Comma) {
692                break;
693            }
694        }
695        parser
696            .expect_token(&sqlparser::tokenizer::Token::RParen)
697            .map_err(ParseError::SqlParseError)?;
698        Ok(StreamingStatement::AlterSource {
699            name,
700            operation: statements::AlterSourceOperation::SetProperties { properties },
701        })
702    } else {
703        Err(ParseError::StreamingError(
704            "Expected ADD COLUMN or SET after ALTER SOURCE <name>".to_string(),
705        ))
706    }
707}
708
709/// Parse a DESCRIBE statement.
710///
711/// # Errors
712///
713/// Returns `ParseError` if the statement syntax is invalid.
714fn parse_describe(
715    parser: &mut sqlparser::parser::Parser,
716) -> Result<StreamingStatement, ParseError> {
717    // Consume DESCRIBE or DESC
718    let token = parser.next_token();
719    match &token.token {
720        sqlparser::tokenizer::Token::Word(w)
721            if w.keyword == sqlparser::keywords::Keyword::DESCRIBE
722                || w.keyword == sqlparser::keywords::Keyword::DESC => {}
723        _ => {
724            return Err(ParseError::StreamingError(
725                "Expected DESCRIBE or DESC".to_string(),
726            ));
727        }
728    }
729    let extended = tokenizer::try_parse_custom_keyword(parser, "EXTENDED");
730    let name = parser
731        .parse_object_name(false)
732        .map_err(ParseError::SqlParseError)?;
733    Ok(StreamingStatement::Describe { name, extended })
734}
735
736/// Parse `SHOW CREATE SOURCE <name>`.
737fn parse_show_create_source(
738    parser: &mut sqlparser::parser::Parser,
739) -> Result<StreamingStatement, ParseError> {
740    // Consume SHOW CREATE SOURCE
741    parser
742        .expect_keyword(sqlparser::keywords::Keyword::SHOW)
743        .map_err(ParseError::SqlParseError)?;
744    parser
745        .expect_keyword(sqlparser::keywords::Keyword::CREATE)
746        .map_err(ParseError::SqlParseError)?;
747    tokenizer::expect_custom_keyword(parser, "SOURCE")?;
748    let name = parser
749        .parse_object_name(false)
750        .map_err(ParseError::SqlParseError)?;
751    Ok(StreamingStatement::Show(ShowCommand::CreateSource { name }))
752}
753
754/// Parse `SHOW CREATE SINK <name>`.
755fn parse_show_create_sink(
756    parser: &mut sqlparser::parser::Parser,
757) -> Result<StreamingStatement, ParseError> {
758    // Consume SHOW CREATE SINK
759    parser
760        .expect_keyword(sqlparser::keywords::Keyword::SHOW)
761        .map_err(ParseError::SqlParseError)?;
762    parser
763        .expect_keyword(sqlparser::keywords::Keyword::CREATE)
764        .map_err(ParseError::SqlParseError)?;
765    tokenizer::expect_custom_keyword(parser, "SINK")?;
766    let name = parser
767        .parse_object_name(false)
768        .map_err(ParseError::SqlParseError)?;
769    Ok(StreamingStatement::Show(ShowCommand::CreateSink { name }))
770}
771
772/// Parse an EXPLAIN [ANALYZE] statement wrapping a streaming query.
773///
774/// Syntax: `EXPLAIN [ANALYZE] <streaming_statement>`
775///
776/// # Errors
777///
778/// Returns `ParseError` if the statement syntax is invalid.
779fn parse_explain(
780    parser: &mut sqlparser::parser::Parser,
781    original_sql: &str,
782) -> Result<StreamingStatement, ParseError> {
783    parser
784        .expect_keyword(sqlparser::keywords::Keyword::EXPLAIN)
785        .map_err(ParseError::SqlParseError)?;
786
787    // Check for optional ANALYZE keyword
788    let analyze = tokenizer::try_parse_custom_keyword(parser, "ANALYZE");
789
790    // Find the position after EXPLAIN [ANALYZE] in the original SQL
791    let explain_prefix_upper = original_sql.to_uppercase();
792    let skip_keyword = if analyze { "ANALYZE" } else { "EXPLAIN" };
793    let inner_start = if analyze {
794        explain_prefix_upper
795            .find("ANALYZE")
796            .map_or(0, |pos| pos + "ANALYZE".len())
797    } else {
798        explain_prefix_upper
799            .find("EXPLAIN")
800            .map_or(0, |pos| pos + "EXPLAIN".len())
801    };
802    let inner_sql = original_sql[inner_start..].trim();
803    let _ = skip_keyword; // suppress unused warning
804
805    // Parse the inner statement recursively
806    let inner_stmts = StreamingParser::parse_sql(inner_sql)?;
807    let inner = inner_stmts.into_iter().next().ok_or_else(|| {
808        sqlparser::parser::ParserError::ParserError("Expected statement after EXPLAIN".to_string())
809    })?;
810    Ok(StreamingStatement::Explain {
811        statement: Box::new(inner),
812        analyze,
813    })
814}
815
816/// Parse a CREATE MATERIALIZED VIEW statement.
817///
818/// Syntax:
819/// ```sql
820/// CREATE [OR REPLACE] MATERIALIZED VIEW [IF NOT EXISTS] name
821/// AS <select_query>
822/// [EMIT <strategy>]
823/// ```
824///
825/// # Errors
826///
827/// Returns `ParseError` if the statement syntax is invalid.
828fn parse_create_materialized_view(
829    parser: &mut sqlparser::parser::Parser,
830    original_sql: &str,
831) -> Result<StreamingStatement, ParseError> {
832    parser
833        .expect_keyword(sqlparser::keywords::Keyword::CREATE)
834        .map_err(ParseError::SqlParseError)?;
835
836    let or_replace = parser.parse_keywords(&[
837        sqlparser::keywords::Keyword::OR,
838        sqlparser::keywords::Keyword::REPLACE,
839    ]);
840
841    parser
842        .expect_keyword(sqlparser::keywords::Keyword::MATERIALIZED)
843        .map_err(ParseError::SqlParseError)?;
844    parser
845        .expect_keyword(sqlparser::keywords::Keyword::VIEW)
846        .map_err(ParseError::SqlParseError)?;
847
848    let if_not_exists = parser.parse_keywords(&[
849        sqlparser::keywords::Keyword::IF,
850        sqlparser::keywords::Keyword::NOT,
851        sqlparser::keywords::Keyword::EXISTS,
852    ]);
853
854    let name = parser
855        .parse_object_name(false)
856        .map_err(ParseError::SqlParseError)?;
857
858    parser
859        .expect_keyword(sqlparser::keywords::Keyword::AS)
860        .map_err(ParseError::SqlParseError)?;
861
862    // Collect remaining tokens and split at EMIT boundary (same strategy as continuous query)
863    let remaining = collect_remaining_tokens(parser);
864    let (query_tokens, emit_tokens) = split_at_emit(&remaining);
865    let raw_query_sql = query_body_sql(original_sql, &query_tokens, &emit_tokens, None);
866
867    let mv_dialect = LaminarDialect::default();
868
869    let (query_stmt, normalized_temporal_sql) = if query_tokens.is_empty() {
870        return Err(ParseError::StreamingError(
871            "Expected SELECT query after AS".to_string(),
872        ));
873    } else if let Some(parsed) = join_parser::parse_temporal_probe_query(&query_tokens)? {
874        (
875            StreamingStatement::TemporalProbeQuery {
876                statement: Box::new(parsed.statement),
877                analysis: Box::new(parsed.analysis),
878            },
879            Some(parsed.normalized_sql),
880        )
881    } else {
882        let mut query_parser =
883            sqlparser::parser::Parser::new(&mv_dialect).with_tokens_with_locations(query_tokens);
884        let query = query_parser
885            .parse_query()
886            .map_err(ParseError::SqlParseError)?;
887        (
888            StreamingStatement::Standard(Box::new(sqlparser::ast::Statement::Query(query))),
889            None,
890        )
891    };
892    let query_sql = normalized_temporal_sql.unwrap_or(raw_query_sql);
893
894    let emit_clause = if emit_tokens.is_empty() {
895        None
896    } else {
897        let mut emit_parser =
898            sqlparser::parser::Parser::new(&mv_dialect).with_tokens_with_locations(emit_tokens);
899        emit_parser::parse_emit_clause(&mut emit_parser)?
900    };
901
902    Ok(StreamingStatement::CreateMaterializedView {
903        name,
904        query: Box::new(query_stmt),
905        emit_clause,
906        or_replace,
907        if_not_exists,
908        query_sql,
909    })
910}
911
912/// Byte offset in `sql` for a sqlparser `Location` (1-indexed line/column).
913fn location_to_byte_offset(sql: &str, loc: sqlparser::tokenizer::Location) -> Option<usize> {
914    if loc.line == 0 {
915        return None;
916    }
917    let (mut line, mut col) = (1u64, 1u64);
918    for (idx, ch) in sql.char_indices() {
919        if line == loc.line && col == loc.column {
920            return Some(idx);
921        }
922        if ch == '\n' {
923            line += 1;
924            col = 1;
925        } else {
926            col += 1;
927        }
928    }
929    (line == loc.line && col == loc.column).then_some(sql.len())
930}
931
932/// Collect all remaining tokens from the parser into a Vec.
933fn collect_remaining_tokens(
934    parser: &mut sqlparser::parser::Parser,
935) -> Vec<sqlparser::tokenizer::TokenWithSpan> {
936    let mut tokens = Vec::new();
937    loop {
938        let token = parser.next_token();
939        if token.token == sqlparser::tokenizer::Token::EOF {
940            tokens.push(token);
941            break;
942        }
943        tokens.push(token);
944    }
945    tokens
946}
947
948/// Split tokens at the first standalone EMIT keyword (not inside parentheses).
949///
950/// Returns (query_tokens, emit_tokens) where emit_tokens starts with EMIT
951/// (or is empty if no EMIT found).
952fn split_at_emit(
953    tokens: &[sqlparser::tokenizer::TokenWithSpan],
954) -> (
955    Vec<sqlparser::tokenizer::TokenWithSpan>,
956    Vec<sqlparser::tokenizer::TokenWithSpan>,
957) {
958    let mut depth: i32 = 0;
959    for (i, token) in tokens.iter().enumerate() {
960        match &token.token {
961            sqlparser::tokenizer::Token::LParen => depth += 1,
962            sqlparser::tokenizer::Token::RParen => {
963                depth -= 1;
964            }
965            sqlparser::tokenizer::Token::Word(w)
966                if depth == 0 && w.value.eq_ignore_ascii_case("EMIT") =>
967            {
968                let mut query_tokens = tokens[..i].to_vec();
969                query_tokens.push(sqlparser::tokenizer::TokenWithSpan {
970                    token: sqlparser::tokenizer::Token::EOF,
971                    span: sqlparser::tokenizer::Span::empty(),
972                });
973                let emit_tokens = tokens[i..].to_vec();
974                return (query_tokens, emit_tokens);
975            }
976            _ => {}
977        }
978    }
979    (tokens.to_vec(), vec![])
980}
981
982/// SQL parsing errors
983#[derive(Debug, thiserror::Error)]
984pub enum ParseError {
985    /// Standard SQL parse error
986    #[error("SQL parse error: {0}")]
987    SqlParseError(#[from] sqlparser::parser::ParserError),
988
989    /// Streaming extension parse error
990    #[error("Streaming SQL error: {0}")]
991    StreamingError(String),
992
993    /// Window function error
994    #[error("Window function error: {0}")]
995    WindowError(String),
996
997    /// Validation error (e.g., invalid option values)
998    #[error("Validation error: {0}")]
999    ValidationError(String),
1000}
1001
1002#[cfg(test)]
1003mod tests;