1pub 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;
15pub 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
36pub fn parse_streaming_sql(sql: &str) -> Result<Vec<StreamingStatement>, ParseError> {
45 StreamingParser::parse_sql(sql).map_err(ParseError::SqlParseError)
46}
47
48pub struct StreamingParser;
54
55impl StreamingParser {
56 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 #[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 pub fn parse_emit_clause(sql: &str) -> Result<Option<EmitClause>, ParseError> {
194 emit_parser::parse_emit_clause_from_sql(sql)
195 }
196
197 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
228fn 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
242fn parse_restore_checkpoint(
250 parser: &mut sqlparser::parser::Parser,
251) -> Result<StreamingStatement, ParseError> {
252 tokenizer::expect_custom_keyword(parser, "RESTORE")?;
254 if !parser.parse_keyword(sqlparser::keywords::Keyword::FROM) {
256 return Err(ParseError::StreamingError(
257 "Expected FROM after RESTORE".to_string(),
258 ));
259 }
260 tokenizer::expect_custom_keyword(parser, "CHECKPOINT")?;
262 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
277fn convert_standard_statement(stmt: sqlparser::ast::Statement) -> StreamingStatement {
282 if let sqlparser::ast::Statement::Insert(insert) = &stmt {
283 if let sqlparser::ast::TableObject::TableName(ref name) = insert.table {
285 let table_name = name.clone();
286 let columns = insert.columns.clone();
287
288 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
304fn 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
333fn 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
362fn 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
396fn 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 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
500fn 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
537fn 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
560pub(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
600fn 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
629fn 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 if parser.parse_keywords(&[
654 sqlparser::keywords::Keyword::ADD,
655 sqlparser::keywords::Keyword::COLUMN,
656 ]) {
657 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 parser
676 .expect_token(&sqlparser::tokenizer::Token::LParen)
677 .map_err(ParseError::SqlParseError)?;
678 #[allow(clippy::disallowed_types)] 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
709fn parse_describe(
715 parser: &mut sqlparser::parser::Parser,
716) -> Result<StreamingStatement, ParseError> {
717 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
736fn parse_show_create_source(
738 parser: &mut sqlparser::parser::Parser,
739) -> Result<StreamingStatement, ParseError> {
740 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
754fn parse_show_create_sink(
756 parser: &mut sqlparser::parser::Parser,
757) -> Result<StreamingStatement, ParseError> {
758 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
772fn 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 let analyze = tokenizer::try_parse_custom_keyword(parser, "ANALYZE");
789
790 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; 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
816fn 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 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
912fn 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
932fn 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
948fn 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#[derive(Debug, thiserror::Error)]
984pub enum ParseError {
985 #[error("SQL parse error: {0}")]
987 SqlParseError(#[from] sqlparser::parser::ParserError),
988
989 #[error("Streaming SQL error: {0}")]
991 StreamingError(String),
992
993 #[error("Window function error: {0}")]
995 WindowError(String),
996
997 #[error("Validation error: {0}")]
999 ValidationError(String),
1000}
1001
1002#[cfg(test)]
1003mod tests;