Skip to main content

risingwave_sqlparser/ast/
statement.rs

1// Copyright 2025 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::fmt;
16use std::fmt::{Formatter, Write};
17
18use itertools::Itertools;
19use winnow::ModalResult;
20
21use super::ddl::SourceWatermark;
22use super::legacy_source::{CompatibleFormatEncode, parse_format_encode};
23use super::value::escape_single_quote_string;
24use super::{EmitMode, Ident, ObjectType, Query, Value};
25use crate::ast::{
26    CdcTableInfo, ColumnDef, ObjectName, REDACT_SQL_OPTION_KEYWORDS, SqlOption, TableConstraint,
27    display_comma_separated, display_separated,
28};
29use crate::keywords::Keyword;
30use crate::parser::{IncludeOption, IsOptional, Parser};
31use crate::parser_err;
32use crate::parser_v2::literal_u32;
33use crate::tokenizer::Token;
34
35/// Consumes token from the parser into an AST node.
36pub trait ParseTo: Sized {
37    fn parse_to(parser: &mut Parser<'_>) -> ModalResult<Self>;
38}
39
40#[macro_export]
41macro_rules! impl_parse_to {
42    () => {};
43    ($field:ident : $field_type:ty, $parser:ident) => {
44        let $field = <$field_type>::parse_to($parser)?;
45    };
46    ($field:ident => [$($arr:tt)+], $parser:ident) => {
47        let $field = $parser.parse_keywords(&[$($arr)+]);
48    };
49    ([$($arr:tt)+], $parser:ident) => {
50        $parser.expect_keywords(&[$($arr)+])?;
51    };
52}
53
54#[macro_export]
55macro_rules! impl_fmt_display {
56    () => {};
57    ($field:ident, $v:ident, $self:ident) => {{
58        let s = format!("{}", $self.$field);
59        if !s.is_empty() {
60            $v.push(s);
61        }
62    }};
63    ($field:ident => [$($arr:tt)+], $v:ident, $self:ident) => {
64        if $self.$field {
65            $v.push(format!("{}", display_separated(&[$($arr)+], " ")));
66        }
67    };
68    ([$($arr:tt)+], $v:ident) => {
69        $v.push(format!("{}", display_separated(&[$($arr)+], " ")));
70    };
71}
72
73// sql_grammar!(CreateSourceStatement {
74//     if_not_exists => [Keyword::IF, Keyword::NOT, Keyword::EXISTS],
75//     source_name: Ident,
76//     with_properties: AstOption<WithProperties>,
77//     [Keyword::ROW, Keyword::FORMAT],
78//     format_encode: SourceSchema,
79//     [Keyword::WATERMARK, Keyword::FOR] column [Keyword::AS] <expr>
80// });
81#[derive(Debug, Clone, PartialEq, Eq, Hash)]
82pub struct CreateSourceStatement {
83    pub temporary: bool,
84    pub if_not_exists: bool,
85    pub columns: Vec<ColumnDef>,
86    // The wildchar position in columns defined in sql. Only exist when using external schema.
87    pub wildcard_idx: Option<usize>,
88    pub constraints: Vec<TableConstraint>,
89    pub source_name: ObjectName,
90    pub with_properties: WithProperties,
91    pub format_encode: CompatibleFormatEncode,
92    pub source_watermarks: Vec<SourceWatermark>,
93    pub include_column_options: IncludeOption,
94    /// `FROM cdc_source TABLE database_name.table_name`
95    pub cdc_table_info: Option<CdcTableInfo>,
96}
97
98/// FORMAT means how to get the operation(Insert/Delete) from the input.
99///
100/// Check `CONNECTORS_COMPATIBLE_FORMATS` for what `FORMAT ... ENCODE ...` combinations are allowed.
101#[derive(Debug, Clone, PartialEq, Eq, Hash)]
102pub enum Format {
103    /// The format is the same with RisingWave's internal representation.
104    /// Used internally for schema change
105    Native,
106    /// for self-explanatory sources like iceberg, they have their own format, and should not be specified by user.
107    None,
108    // Keyword::DEBEZIUM
109    Debezium,
110    // Keyword::DEBEZIUM_MONGO
111    DebeziumMongo,
112    // Keyword::MAXWELL
113    Maxwell,
114    // Keyword::CANAL
115    Canal,
116    // Keyword::UPSERT
117    Upsert,
118    // Keyword::PLAIN
119    Plain,
120}
121
122// TODO: unify with `from_keyword`
123impl fmt::Display for Format {
124    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
125        write!(
126            f,
127            "{}",
128            match self {
129                Format::Native => "NATIVE",
130                Format::Debezium => "DEBEZIUM",
131                Format::DebeziumMongo => "DEBEZIUM_MONGO",
132                Format::Maxwell => "MAXWELL",
133                Format::Canal => "CANAL",
134                Format::Upsert => "UPSERT",
135                Format::Plain => "PLAIN",
136                Format::None => "NONE",
137            }
138        )
139    }
140}
141
142impl Format {
143    pub fn from_keyword(s: &str) -> ModalResult<Self> {
144        Ok(match s {
145            "DEBEZIUM" => Format::Debezium,
146            "DEBEZIUM_MONGO" => Format::DebeziumMongo,
147            "MAXWELL" => Format::Maxwell,
148            "CANAL" => Format::Canal,
149            "PLAIN" => Format::Plain,
150            "UPSERT" => Format::Upsert,
151            "NATIVE" => Format::Native,
152            "NONE" => Format::None,
153            _ => parser_err!("expected PROTOBUF | DEBEZIUM | PLAIN | NATIVE | NONE after FORMAT"),
154        })
155    }
156}
157
158/// Check `CONNECTORS_COMPATIBLE_FORMATS` for what `FORMAT ... ENCODE ...` combinations are allowed.
159#[derive(Debug, Clone, PartialEq, Eq, Hash)]
160pub enum Encode {
161    Avro,     // Keyword::Avro
162    Csv,      // Keyword::CSV
163    Protobuf, // Keyword::PROTOBUF
164    Json,     // Keyword::JSON
165    Bytes,    // Keyword::BYTES
166    /// for self-explanatory sources like iceberg, they have their own format, and should not be specified by user.
167    None,
168    Text, // Keyword::TEXT
169    /// The encode is the same with RisingWave's internal representation.
170    /// Used internally for schema change
171    Native,
172    Template,
173    Parquet,
174}
175
176// TODO: unify with `from_keyword`
177impl fmt::Display for Encode {
178    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
179        write!(
180            f,
181            "{}",
182            match self {
183                Encode::Avro => "AVRO",
184                Encode::Csv => "CSV",
185                Encode::Protobuf => "PROTOBUF",
186                Encode::Json => "JSON",
187                Encode::Bytes => "BYTES",
188                Encode::Native => "NATIVE",
189                Encode::Template => "TEMPLATE",
190                Encode::None => "NONE",
191                Encode::Parquet => "PARQUET",
192                Encode::Text => "TEXT",
193            }
194        )
195    }
196}
197
198impl Encode {
199    pub fn from_keyword(s: &str) -> ModalResult<Self> {
200        Ok(match s {
201            "AVRO" => Encode::Avro,
202            "TEXT" => Encode::Text,
203            "BYTES" => Encode::Bytes,
204            "CSV" => Encode::Csv,
205            "PROTOBUF" => Encode::Protobuf,
206            "JSON" => Encode::Json,
207            "TEMPLATE" => Encode::Template,
208            "PARQUET" => Encode::Parquet,
209            "NATIVE" => Encode::Native,
210            "NONE" => Encode::None,
211            _ => parser_err!(
212                "expected AVRO | BYTES | CSV | PROTOBUF | JSON | NATIVE | TEMPLATE | PARQUET | NONE after Encode"
213            ),
214        })
215    }
216}
217
218/// `FORMAT ... ENCODE ... [(a=b, ...)] [KEY ENCODE ...]`
219#[derive(Debug, Clone, PartialEq, Eq, Hash)]
220pub struct FormatEncodeOptions {
221    pub format: Format,
222    pub row_encode: Encode,
223    pub row_options: Vec<SqlOption>,
224
225    pub key_encode: Option<Encode>,
226}
227
228impl Parser<'_> {
229    /// Peek the next tokens to see if it is `FORMAT` or `ROW FORMAT` (for compatibility).
230    fn peek_format_encode_format(&mut self) -> bool {
231        (self.peek_nth_any_of_keywords(0, &[Keyword::ROW])
232            && self.peek_nth_any_of_keywords(1, &[Keyword::FORMAT])) // ROW FORMAT
233            || self.peek_nth_any_of_keywords(0, &[Keyword::FORMAT]) // FORMAT
234    }
235
236    /// Parse the source schema. The behavior depends on the `connector` type.
237    pub fn parse_format_encode_with_connector(
238        &mut self,
239        connector: &str,
240        cdc_source_job: bool,
241    ) -> ModalResult<CompatibleFormatEncode> {
242        // row format for cdc source must be debezium json
243        // row format for nexmark source must be native
244        // default row format for datagen source is native
245        // FIXME: parse input `connector` to enum type instead using string here
246        if connector.contains("-cdc") {
247            let expected = if cdc_source_job {
248                FormatEncodeOptions::plain_json()
249            } else if connector.contains("mongodb") {
250                FormatEncodeOptions::debezium_mongo_json()
251            } else {
252                FormatEncodeOptions::debezium_json()
253            };
254
255            if self.peek_format_encode_format() {
256                let schema = parse_format_encode(self)?.into_v2();
257                if schema != expected {
258                    parser_err!(
259                        "Row format for CDC connectors should be \
260                         either omitted or set to `{expected}`",
261                    );
262                }
263            }
264            Ok(expected.into())
265        } else if connector.contains("nexmark") {
266            let expected = FormatEncodeOptions::native();
267            if self.peek_format_encode_format() {
268                let schema = parse_format_encode(self)?.into_v2();
269                if schema != expected {
270                    parser_err!(
271                        "Row format for nexmark connectors should be \
272                         either omitted or set to `{expected}`",
273                    );
274                }
275            }
276            Ok(expected.into())
277        } else if connector.contains("datagen") {
278            Ok(if self.peek_format_encode_format() {
279                parse_format_encode(self)?
280            } else {
281                FormatEncodeOptions::native().into()
282            })
283        } else if connector.contains("iceberg") || connector.contains("adbc_snowflake") {
284            let expected = FormatEncodeOptions::none();
285            if self.peek_format_encode_format() {
286                let schema = parse_format_encode(self)?.into_v2();
287                if schema != expected {
288                    parser_err!(
289                        "Row format for iceberg connectors should be \
290                         either omitted or set to `{expected}`",
291                    );
292                }
293            }
294            Ok(expected.into())
295        } else if connector.contains("webhook") {
296            parser_err!(
297                "Source with webhook connector is not supported. \
298                 Please use the `CREATE TABLE ... WITH ...` statement instead.",
299            );
300        } else {
301            Ok(parse_format_encode(self)?)
302        }
303    }
304
305    /// Parse `FORMAT ... ENCODE ... (...)`.
306    pub fn parse_schema(&mut self) -> ModalResult<Option<FormatEncodeOptions>> {
307        if !self.parse_keyword(Keyword::FORMAT) {
308            return Ok(None);
309        }
310
311        let id = self.parse_identifier()?;
312        let s = id.value.to_ascii_uppercase();
313        let format = Format::from_keyword(&s)?;
314        self.expect_keyword(Keyword::ENCODE)?;
315        let id = self.parse_identifier()?;
316        let s = id.value.to_ascii_uppercase();
317        let row_encode = Encode::from_keyword(&s)?;
318        let row_options = self.parse_options()?;
319
320        let key_encode = if self.parse_keywords(&[Keyword::KEY, Keyword::ENCODE]) {
321            Some(Encode::from_keyword(
322                self.parse_identifier()?.value.to_ascii_uppercase().as_str(),
323            )?)
324        } else {
325            None
326        };
327
328        Ok(Some(FormatEncodeOptions {
329            format,
330            row_encode,
331            row_options,
332            key_encode,
333        }))
334    }
335}
336
337impl FormatEncodeOptions {
338    pub const fn plain_json() -> Self {
339        FormatEncodeOptions {
340            format: Format::Plain,
341            row_encode: Encode::Json,
342            row_options: Vec::new(),
343            key_encode: None,
344        }
345    }
346
347    /// Create a new source schema with `Debezium` format and `Json` encoding.
348    pub const fn debezium_json() -> Self {
349        FormatEncodeOptions {
350            format: Format::Debezium,
351            row_encode: Encode::Json,
352            row_options: Vec::new(),
353            key_encode: None,
354        }
355    }
356
357    pub const fn debezium_mongo_json() -> Self {
358        FormatEncodeOptions {
359            format: Format::DebeziumMongo,
360            row_encode: Encode::Json,
361            row_options: Vec::new(),
362            key_encode: None,
363        }
364    }
365
366    /// Create a new source schema with `Native` format and encoding.
367    pub const fn native() -> Self {
368        FormatEncodeOptions {
369            format: Format::Native,
370            row_encode: Encode::Native,
371            row_options: Vec::new(),
372            key_encode: None,
373        }
374    }
375
376    /// Create a new source schema with `None` format and encoding.
377    /// Used for self-explanatory source like iceberg.
378    pub const fn none() -> Self {
379        FormatEncodeOptions {
380            format: Format::None,
381            row_encode: Encode::None,
382            row_options: Vec::new(),
383            key_encode: None,
384        }
385    }
386
387    pub fn row_options(&self) -> &[SqlOption] {
388        self.row_options.as_ref()
389    }
390}
391
392impl fmt::Display for FormatEncodeOptions {
393    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
394        write!(f, "FORMAT {} ENCODE {}", self.format, self.row_encode)?;
395
396        if !self.row_options().is_empty() {
397            write!(f, " ({})", display_comma_separated(self.row_options()))?;
398        }
399
400        if let Some(key_encode) = &self.key_encode {
401            write!(f, " KEY ENCODE {}", key_encode)?;
402        }
403
404        Ok(())
405    }
406}
407
408pub(super) fn fmt_create_items(
409    columns: &[ColumnDef],
410    constraints: &[TableConstraint],
411    watermarks: &[SourceWatermark],
412    wildcard_idx: Option<usize>,
413) -> std::result::Result<String, fmt::Error> {
414    let mut items = String::new();
415    let has_items = !columns.is_empty()
416        || !constraints.is_empty()
417        || !watermarks.is_empty()
418        || wildcard_idx.is_some();
419    has_items.then(|| write!(&mut items, "("));
420
421    if let Some(wildcard_idx) = wildcard_idx {
422        let (columns_l, columns_r) = columns.split_at(wildcard_idx);
423        write!(&mut items, "{}", display_comma_separated(columns_l))?;
424        if !columns_l.is_empty() {
425            write!(&mut items, ", ")?;
426        }
427        write!(&mut items, "{}", Token::Mul)?;
428        if !columns_r.is_empty() {
429            write!(&mut items, ", ")?;
430        }
431        write!(&mut items, "{}", display_comma_separated(columns_r))?;
432    } else {
433        write!(&mut items, "{}", display_comma_separated(columns))?;
434    }
435    let mut leading_items = !columns.is_empty() || wildcard_idx.is_some();
436
437    if leading_items && !constraints.is_empty() {
438        write!(&mut items, ", ")?;
439    }
440    write!(&mut items, "{}", display_comma_separated(constraints))?;
441    leading_items |= !constraints.is_empty();
442
443    if leading_items && !watermarks.is_empty() {
444        write!(&mut items, ", ")?;
445    }
446    write!(&mut items, "{}", display_comma_separated(watermarks))?;
447    // uncomment this when adding more sections below
448    // leading_items |= !watermarks.is_empty();
449
450    has_items.then(|| write!(&mut items, ")"));
451    Ok(items)
452}
453
454impl fmt::Display for CreateSourceStatement {
455    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
456        let mut v: Vec<String> = vec![];
457        impl_fmt_display!(if_not_exists => [Keyword::IF, Keyword::NOT, Keyword::EXISTS], v, self);
458        impl_fmt_display!(source_name, v, self);
459
460        let items = fmt_create_items(
461            &self.columns,
462            &self.constraints,
463            &self.source_watermarks,
464            self.wildcard_idx,
465        )?;
466        if !items.is_empty() {
467            v.push(items);
468        }
469
470        for item in &self.include_column_options {
471            v.push(format!("{}", item));
472        }
473
474        // skip format_encode for cdc source
475        let is_cdc_source = self.with_properties.0.iter().any(|option| {
476            option.name.real_value().eq_ignore_ascii_case("connector")
477                && option.value.to_string().contains("cdc")
478        });
479
480        impl_fmt_display!(with_properties, v, self);
481        if self.cdc_table_info.is_none() && !is_cdc_source {
482            impl_fmt_display!(format_encode, v, self);
483        }
484        if let Some(info) = &self.cdc_table_info {
485            v.push(format!(
486                "FROM {} TABLE '{}'",
487                info.source_name,
488                escape_single_quote_string(&info.external_table_name)
489            ));
490        }
491        v.iter().join(" ").fmt(f)
492    }
493}
494
495#[derive(Debug, Clone, PartialEq, Eq, Hash)]
496pub enum CreateSink {
497    From(ObjectName),
498    AsQuery(Box<Query>),
499}
500
501impl fmt::Display for CreateSink {
502    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
503        match self {
504            Self::From(mv) => write!(f, "FROM {}", mv),
505            Self::AsQuery(query) => write!(f, "AS {}", query),
506        }
507    }
508}
509// sql_grammar!(CreateSinkStatement {
510//     if_not_exists => [Keyword::IF, Keyword::NOT, Keyword::EXISTS],
511//     sink_name: Ident,
512//     [Keyword::FROM],
513//     materialized_view: Ident,
514//     with_properties: AstOption<WithProperties>,
515// });
516#[derive(Debug, Clone, PartialEq, Eq, Hash)]
517pub struct CreateSinkStatement {
518    pub or_replace: bool,
519    pub if_not_exists: bool,
520    pub sink_name: ObjectName,
521    pub with_properties: WithProperties,
522    pub sink_from: CreateSink,
523
524    // only used when creating sink into a table
525    // insert to specific columns of the target table
526    pub columns: Vec<Ident>,
527    pub emit_mode: Option<EmitMode>,
528    pub sink_schema: Option<FormatEncodeOptions>,
529    pub into_table_name: Option<ObjectName>,
530}
531
532impl ParseTo for CreateSinkStatement {
533    fn parse_to(p: &mut Parser<'_>) -> ModalResult<Self> {
534        Self::parse_to_with_or_replace(p, false)
535    }
536}
537
538impl CreateSinkStatement {
539    pub fn parse_to_with_or_replace(p: &mut Parser<'_>, or_replace: bool) -> ModalResult<Self> {
540        impl_parse_to!(if_not_exists => [Keyword::IF, Keyword::NOT, Keyword::EXISTS], p);
541        impl_parse_to!(sink_name: ObjectName, p);
542
543        let mut target_spec_columns = Vec::new();
544        let into_table_name = if p.parse_keyword(Keyword::INTO) {
545            impl_parse_to!(into_table_name: ObjectName, p);
546
547            // we only allow specify columns when creating sink into a table
548            target_spec_columns = p.parse_parenthesized_column_list(IsOptional::Optional)?;
549            Some(into_table_name)
550        } else {
551            None
552        };
553
554        let sink_from = if p.parse_keyword(Keyword::FROM) {
555            impl_parse_to!(from_name: ObjectName, p);
556            CreateSink::From(from_name)
557        } else if p.parse_keyword(Keyword::AS) {
558            let query = Box::new(p.parse_query()?);
559            CreateSink::AsQuery(query)
560        } else {
561            p.expected(if or_replace {
562                "FROM or AS after REPLACE SINK sink_name"
563            } else {
564                "FROM or AS after CREATE SINK sink_name"
565            })?
566        };
567
568        let emit_mode: Option<EmitMode> = p.parse_emit_mode()?;
569
570        // This check cannot be put into the `WithProperties::parse_to`, since other
571        // statements may not need the with properties.
572        if !p.peek_nth_any_of_keywords(0, &[Keyword::WITH]) && into_table_name.is_none() {
573            p.expected("WITH")?
574        }
575        impl_parse_to!(with_properties: WithProperties, p);
576
577        if with_properties.0.is_empty() && into_table_name.is_none() {
578            parser_err!("sink properties not provided");
579        }
580
581        let sink_schema = p.parse_schema()?;
582
583        Ok(Self {
584            or_replace,
585            if_not_exists,
586            sink_name,
587            with_properties,
588            sink_from,
589            columns: target_spec_columns,
590            emit_mode,
591            sink_schema,
592            into_table_name,
593        })
594    }
595}
596
597impl fmt::Display for CreateSinkStatement {
598    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
599        let mut v: Vec<String> = vec![];
600        impl_fmt_display!(if_not_exists => [Keyword::IF, Keyword::NOT, Keyword::EXISTS], v, self);
601        impl_fmt_display!(sink_name, v, self);
602        if let Some(into_table) = &self.into_table_name {
603            impl_fmt_display!([Keyword::INTO], v);
604            impl_fmt_display!([into_table], v);
605            if !self.columns.is_empty() {
606                v.push(format!("({})", display_comma_separated(&self.columns)));
607            }
608        }
609        impl_fmt_display!(sink_from, v, self);
610        if let Some(ref emit_mode) = self.emit_mode {
611            v.push(format!("EMIT {}", emit_mode));
612        }
613        impl_fmt_display!(with_properties, v, self);
614        if let Some(schema) = &self.sink_schema {
615            v.push(format!("{}", schema));
616        }
617        v.iter().join(" ").fmt(f)
618    }
619}
620
621// sql_grammar!(CreateSubscriptionStatement {
622//     if_not_exists => [Keyword::IF, Keyword::NOT, Keyword::EXISTS],
623//     subscription_name: Ident,
624//     [Keyword::FROM],
625//     materialized_view: Ident,
626//     with_properties: AstOption<WithProperties>,
627// });
628#[derive(Debug, Clone, PartialEq, Eq, Hash)]
629pub struct CreateSubscriptionStatement {
630    pub if_not_exists: bool,
631    pub subscription_name: ObjectName,
632    pub with_properties: WithProperties,
633    pub subscription_from: ObjectName,
634    // pub emit_mode: Option<EmitMode>,
635}
636
637impl ParseTo for CreateSubscriptionStatement {
638    fn parse_to(p: &mut Parser<'_>) -> ModalResult<Self> {
639        impl_parse_to!(if_not_exists => [Keyword::IF, Keyword::NOT, Keyword::EXISTS], p);
640        impl_parse_to!(subscription_name: ObjectName, p);
641
642        let subscription_from = if p.parse_keyword(Keyword::FROM) {
643            impl_parse_to!(from_name: ObjectName, p);
644            from_name
645        } else {
646            p.expected("FROM after CREATE SUBSCRIPTION subscription_name")?
647        };
648
649        // let emit_mode = p.parse_emit_mode()?;
650
651        // This check cannot be put into the `WithProperties::parse_to`, since other
652        // statements may not need the with properties.
653        if !p.peek_nth_any_of_keywords(0, &[Keyword::WITH]) {
654            p.expected("WITH")?
655        }
656        impl_parse_to!(with_properties: WithProperties, p);
657
658        if with_properties.0.is_empty() {
659            parser_err!("subscription properties not provided");
660        }
661
662        Ok(Self {
663            if_not_exists,
664            subscription_name,
665            with_properties,
666            subscription_from,
667        })
668    }
669}
670
671impl fmt::Display for CreateSubscriptionStatement {
672    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
673        let mut v: Vec<String> = vec![];
674        impl_fmt_display!(if_not_exists => [Keyword::IF, Keyword::NOT, Keyword::EXISTS], v, self);
675        impl_fmt_display!(subscription_name, v, self);
676        v.push(format!("FROM {}", self.subscription_from));
677        impl_fmt_display!(with_properties, v, self);
678        v.iter().join(" ").fmt(f)
679    }
680}
681
682#[derive(Debug, Clone, PartialEq, Eq, Hash)]
683pub enum DeclareCursor {
684    Query(Box<Query>),
685    Subscription(ObjectName, Since),
686}
687
688impl fmt::Display for DeclareCursor {
689    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
690        let mut v: Vec<String> = vec![];
691        match self {
692            DeclareCursor::Query(query) => v.push(format!("{}", query.as_ref())),
693            DeclareCursor::Subscription(name, since) => {
694                v.push(format!("{}", name));
695                v.push(format!("{:?}", since));
696            }
697        }
698        v.iter().join(" ").fmt(f)
699    }
700}
701
702// sql_grammar!(DeclareCursorStatement {
703//     cursor_name: Ident,
704//     [Keyword::SUBSCRIPTION]
705//     [Keyword::CURSOR],
706//     [Keyword::FOR],
707//     subscription: Ident or query: Query,
708//     [Keyword::SINCE],
709//     rw_timestamp: Ident,
710// });
711#[derive(Debug, Clone, PartialEq, Eq, Hash)]
712pub struct DeclareCursorStatement {
713    pub cursor_name: Ident,
714    pub declare_cursor: DeclareCursor,
715}
716
717impl ParseTo for DeclareCursorStatement {
718    fn parse_to(p: &mut Parser<'_>) -> ModalResult<Self> {
719        let cursor_name = p.parse_identifier_non_reserved()?;
720
721        let declare_cursor = if !p.parse_keyword(Keyword::SUBSCRIPTION) {
722            p.expect_keyword(Keyword::CURSOR)?;
723            p.expect_keyword(Keyword::FOR)?;
724            DeclareCursor::Query(Box::new(p.parse_query()?))
725        } else {
726            p.expect_keyword(Keyword::CURSOR)?;
727            p.expect_keyword(Keyword::FOR)?;
728            let cursor_for_name = p.parse_object_name()?;
729            let rw_timestamp = p.parse_since()?;
730            DeclareCursor::Subscription(cursor_for_name, rw_timestamp)
731        };
732
733        Ok(Self {
734            cursor_name,
735            declare_cursor,
736        })
737    }
738}
739
740impl fmt::Display for DeclareCursorStatement {
741    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
742        let mut v: Vec<String> = vec![];
743        impl_fmt_display!(cursor_name, v, self);
744        match &self.declare_cursor {
745            DeclareCursor::Query(_) => {
746                v.push("CURSOR FOR ".to_owned());
747            }
748            DeclareCursor::Subscription { .. } => {
749                v.push("SUBSCRIPTION CURSOR FOR ".to_owned());
750            }
751        }
752        impl_fmt_display!(declare_cursor, v, self);
753        v.iter().join(" ").fmt(f)
754    }
755}
756
757// sql_grammar!(FetchCursorStatement {
758//     cursor_name: Ident,
759// });
760#[derive(Debug, Clone, PartialEq, Eq, Hash)]
761pub struct FetchCursorStatement {
762    pub cursor_name: Ident,
763    pub count: u32,
764    pub with_properties: WithProperties,
765}
766
767impl ParseTo for FetchCursorStatement {
768    fn parse_to(p: &mut Parser<'_>) -> ModalResult<Self> {
769        let count = if p.parse_keyword(Keyword::NEXT) {
770            1
771        } else {
772            literal_u32(p)?
773        };
774        p.expect_keyword(Keyword::FROM)?;
775        let cursor_name = p.parse_identifier_non_reserved()?;
776        impl_parse_to!(with_properties: WithProperties, p);
777
778        Ok(Self {
779            cursor_name,
780            count,
781            with_properties,
782        })
783    }
784}
785
786impl fmt::Display for FetchCursorStatement {
787    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
788        let mut v: Vec<String> = vec![];
789        if self.count == 1 {
790            v.push("NEXT ".to_owned());
791        } else {
792            impl_fmt_display!(count, v, self);
793        }
794        v.push("FROM ".to_owned());
795        impl_fmt_display!(cursor_name, v, self);
796        v.iter().join(" ").fmt(f)
797    }
798}
799
800// sql_grammar!(CloseCursorStatement {
801//     cursor_name: Ident,
802// });
803#[derive(Debug, Clone, PartialEq, Eq, Hash)]
804pub struct CloseCursorStatement {
805    pub cursor_name: Option<Ident>,
806}
807
808impl ParseTo for CloseCursorStatement {
809    fn parse_to(p: &mut Parser<'_>) -> ModalResult<Self> {
810        let cursor_name = if p.parse_keyword(Keyword::ALL) {
811            None
812        } else {
813            Some(p.parse_identifier_non_reserved()?)
814        };
815
816        Ok(Self { cursor_name })
817    }
818}
819
820impl fmt::Display for CloseCursorStatement {
821    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
822        let mut v: Vec<String> = vec![];
823        if let Some(cursor_name) = &self.cursor_name {
824            v.push(format!("{}", cursor_name));
825        } else {
826            v.push("ALL".to_owned());
827        }
828        v.iter().join(" ").fmt(f)
829    }
830}
831
832// sql_grammar!(CreateConnectionStatement {
833//     if_not_exists => [Keyword::IF, Keyword::NOT, Keyword::EXISTS],
834//     connection_name: Ident,
835//     with_properties: AstOption<WithProperties>,
836// });
837#[derive(Debug, Clone, PartialEq, Eq, Hash)]
838pub struct CreateConnectionStatement {
839    pub if_not_exists: bool,
840    pub connection_name: ObjectName,
841    pub with_properties: WithProperties,
842}
843
844impl ParseTo for CreateConnectionStatement {
845    fn parse_to(p: &mut Parser<'_>) -> ModalResult<Self> {
846        impl_parse_to!(if_not_exists => [Keyword::IF, Keyword::NOT, Keyword::EXISTS], p);
847        impl_parse_to!(connection_name: ObjectName, p);
848        impl_parse_to!(with_properties: WithProperties, p);
849        if with_properties.0.is_empty() {
850            parser_err!("connection properties not provided");
851        }
852
853        Ok(Self {
854            if_not_exists,
855            connection_name,
856            with_properties,
857        })
858    }
859}
860
861impl fmt::Display for CreateConnectionStatement {
862    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
863        let mut v: Vec<String> = vec![];
864        impl_fmt_display!(if_not_exists => [Keyword::IF, Keyword::NOT, Keyword::EXISTS], v, self);
865        impl_fmt_display!(connection_name, v, self);
866        impl_fmt_display!(with_properties, v, self);
867        v.iter().join(" ").fmt(f)
868    }
869}
870
871#[derive(Debug, Clone, PartialEq, Eq, Hash)]
872pub struct CreateSecretStatement {
873    pub if_not_exists: bool,
874    pub secret_name: ObjectName,
875    pub credential: Value,
876    pub with_properties: WithProperties,
877}
878
879impl ParseTo for CreateSecretStatement {
880    fn parse_to(parser: &mut Parser<'_>) -> ModalResult<Self> {
881        impl_parse_to!(if_not_exists => [Keyword::IF, Keyword::NOT, Keyword::EXISTS], parser);
882        impl_parse_to!(secret_name: ObjectName, parser);
883        impl_parse_to!(with_properties: WithProperties, parser);
884        let mut credential = Value::Null;
885        if parser.parse_keyword(Keyword::AS) {
886            credential = parser.ensure_parse_value()?;
887        }
888        Ok(Self {
889            if_not_exists,
890            secret_name,
891            credential,
892            with_properties,
893        })
894    }
895}
896
897impl fmt::Display for CreateSecretStatement {
898    fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
899        let mut v: Vec<String> = vec![];
900        impl_fmt_display!(if_not_exists => [Keyword::IF, Keyword::NOT, Keyword::EXISTS], v, self);
901        impl_fmt_display!(secret_name, v, self);
902        impl_fmt_display!(with_properties, v, self);
903        if self.credential != Value::Null {
904            v.push("AS".to_owned());
905            impl_fmt_display!(credential, v, self);
906        }
907        v.iter().join(" ").fmt(f)
908    }
909}
910
911#[derive(Debug, Clone, PartialEq, Eq, Hash)]
912pub struct WithProperties(pub Vec<SqlOption>);
913
914impl ParseTo for WithProperties {
915    fn parse_to(parser: &mut Parser<'_>) -> ModalResult<Self> {
916        Ok(Self(
917            parser.parse_options_with_preceding_keyword(Keyword::WITH)?,
918        ))
919    }
920}
921
922impl fmt::Display for WithProperties {
923    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
924        if !self.0.is_empty() {
925            write!(f, "WITH ({})", display_comma_separated(self.0.as_slice()))
926        } else {
927            Ok(())
928        }
929    }
930}
931
932#[derive(Debug, Clone, PartialEq, Eq, Hash)]
933pub enum Since {
934    TimestampMsNum(u64),
935    ProcessTime,
936    Begin,
937    Full,
938}
939
940impl fmt::Display for Since {
941    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
942        use Since::*;
943        match self {
944            TimestampMsNum(ts) => write!(f, " SINCE {}", ts),
945            ProcessTime => write!(f, " SINCE PROCTIME()"),
946            Begin => write!(f, " SINCE BEGIN()"),
947            Full => write!(f, " FULL"),
948        }
949    }
950}
951
952/// String literal. The difference with String is that it is displayed with
953/// single-quotes.
954#[derive(Debug, Clone, PartialEq, Eq, Hash)]
955pub struct AstString(pub String);
956
957impl ParseTo for AstString {
958    fn parse_to(parser: &mut Parser<'_>) -> ModalResult<Self> {
959        Ok(Self(parser.parse_literal_string()?))
960    }
961}
962
963impl fmt::Display for AstString {
964    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
965        write!(f, "'{}'", self.0)
966    }
967}
968
969/// This trait is used to replace `Option` because `fmt::Display` can not be implemented for
970/// `Option<T>`.
971#[derive(Debug, Clone, PartialEq, Eq, Hash)]
972pub enum AstOption<T> {
973    /// No value
974    None,
975    /// Some value `T`
976    Some(T),
977}
978
979impl<T: ParseTo> ParseTo for AstOption<T> {
980    fn parse_to(parser: &mut Parser<'_>) -> ModalResult<Self> {
981        match T::parse_to(parser) {
982            Ok(t) => Ok(AstOption::Some(t)),
983            Err(_) => Ok(AstOption::None),
984        }
985    }
986}
987
988impl<T: fmt::Display> fmt::Display for AstOption<T> {
989    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
990        match &self {
991            AstOption::Some(t) => t.fmt(f),
992            AstOption::None => Ok(()),
993        }
994    }
995}
996
997impl<T> From<AstOption<T>> for Option<T> {
998    fn from(val: AstOption<T>) -> Self {
999        match val {
1000            AstOption::Some(t) => Some(t),
1001            AstOption::None => None,
1002        }
1003    }
1004}
1005
1006#[derive(Debug, Clone, PartialEq, Eq, Hash)]
1007pub struct CreateUserStatement {
1008    pub user_name: ObjectName,
1009    pub with_options: UserOptions,
1010}
1011
1012#[derive(Debug, Clone, PartialEq, Eq, Hash)]
1013pub struct AlterUserStatement {
1014    pub user_name: ObjectName,
1015    pub mode: AlterUserMode,
1016}
1017
1018#[derive(Debug, Clone, PartialEq, Eq, Hash)]
1019pub enum AlterUserMode {
1020    Options(UserOptions),
1021    Rename(ObjectName),
1022}
1023
1024#[derive(Debug, Clone, PartialEq, Eq, Hash)]
1025pub enum UserOption {
1026    SuperUser,
1027    NoSuperUser,
1028    CreateDB,
1029    NoCreateDB,
1030    CreateUser,
1031    NoCreateUser,
1032    Login,
1033    NoLogin,
1034    Admin,
1035    NoAdmin,
1036    EncryptedPassword(AstString),
1037    Password(Option<AstString>),
1038    OAuth(Vec<SqlOption>),
1039}
1040
1041impl fmt::Display for UserOption {
1042    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1043        match self {
1044            UserOption::SuperUser => write!(f, "SUPERUSER"),
1045            UserOption::NoSuperUser => write!(f, "NOSUPERUSER"),
1046            UserOption::CreateDB => write!(f, "CREATEDB"),
1047            UserOption::NoCreateDB => write!(f, "NOCREATEDB"),
1048            UserOption::CreateUser => write!(f, "CREATEUSER"),
1049            UserOption::NoCreateUser => write!(f, "NOCREATEUSER"),
1050            UserOption::Login => write!(f, "LOGIN"),
1051            UserOption::NoLogin => write!(f, "NOLOGIN"),
1052            UserOption::Admin => write!(f, "ADMIN"),
1053            UserOption::NoAdmin => write!(f, "NOADMIN"),
1054            UserOption::EncryptedPassword(p) => {
1055                if should_redact_user_password() {
1056                    write!(f, "ENCRYPTED PASSWORD [REDACTED]")
1057                } else {
1058                    write!(f, "ENCRYPTED PASSWORD {}", p)
1059                }
1060            }
1061            UserOption::Password(None) => write!(f, "PASSWORD NULL"),
1062            UserOption::Password(Some(p)) => {
1063                if should_redact_user_password() {
1064                    write!(f, "PASSWORD [REDACTED]")
1065                } else {
1066                    write!(f, "PASSWORD {}", p)
1067                }
1068            }
1069            UserOption::OAuth(options) => {
1070                write!(f, "({})", display_comma_separated(options.as_slice()))
1071            }
1072        }
1073    }
1074}
1075
1076#[derive(Debug, Clone, PartialEq, Eq, Hash)]
1077pub struct UserOptions(pub Vec<UserOption>, pub bool);
1078
1079fn should_redact_user_password() -> bool {
1080    REDACT_SQL_OPTION_KEYWORDS.try_with(|_| ()).is_ok()
1081}
1082
1083#[derive(Default)]
1084struct UserOptionsBuilder {
1085    with_prefix: bool,
1086    super_user: Option<UserOption>,
1087    create_db: Option<UserOption>,
1088    create_user: Option<UserOption>,
1089    login: Option<UserOption>,
1090    admin: Option<UserOption>,
1091    password: Option<UserOption>,
1092}
1093
1094impl UserOptionsBuilder {
1095    fn build(self) -> UserOptions {
1096        let mut options = vec![];
1097        if let Some(option) = self.super_user {
1098            options.push(option);
1099        }
1100        if let Some(option) = self.create_db {
1101            options.push(option);
1102        }
1103        if let Some(option) = self.create_user {
1104            options.push(option);
1105        }
1106        if let Some(option) = self.login {
1107            options.push(option);
1108        }
1109        if let Some(option) = self.admin {
1110            options.push(option);
1111        }
1112        if let Some(option) = self.password {
1113            options.push(option);
1114        }
1115        UserOptions(options, self.with_prefix)
1116    }
1117}
1118
1119impl ParseTo for UserOptions {
1120    fn parse_to(parser: &mut Parser<'_>) -> ModalResult<Self> {
1121        let mut builder = UserOptionsBuilder {
1122            with_prefix: parser.parse_keyword(Keyword::WITH),
1123            ..Default::default()
1124        };
1125        let add_option = |item: &mut Option<UserOption>, user_option| {
1126            let old_value = item.replace(user_option);
1127            if old_value.is_some() {
1128                parser_err!("conflicting or redundant options");
1129            }
1130            Ok(())
1131        };
1132        loop {
1133            let token = parser.peek_token();
1134            if token == Token::EOF || token == Token::SemiColon {
1135                break;
1136            }
1137
1138            let option = parser.parse_identifier()?;
1139            let s = option.real_value();
1140            let (item_mut_ref, user_option) = match &s[..] {
1141                "superuser" => (&mut builder.super_user, UserOption::SuperUser),
1142                "nosuperuser" => (&mut builder.super_user, UserOption::NoSuperUser),
1143                "createdb" => (&mut builder.create_db, UserOption::CreateDB),
1144                "nocreatedb" => (&mut builder.create_db, UserOption::NoCreateDB),
1145                "createuser" => (&mut builder.create_user, UserOption::CreateUser),
1146                "nocreateuser" => (&mut builder.create_user, UserOption::NoCreateUser),
1147                "login" => (&mut builder.login, UserOption::Login),
1148                "nologin" => (&mut builder.login, UserOption::NoLogin),
1149                "admin" => (&mut builder.admin, UserOption::Admin),
1150                "noadmin" => (&mut builder.admin, UserOption::NoAdmin),
1151                "password" => {
1152                    if parser.parse_keyword(Keyword::NULL) {
1153                        (&mut builder.password, UserOption::Password(None))
1154                    } else {
1155                        (
1156                            &mut builder.password,
1157                            UserOption::Password(Some(AstString::parse_to(parser)?)),
1158                        )
1159                    }
1160                }
1161                "encrypted" => {
1162                    let option = parser.parse_identifier()?;
1163                    let s = option.real_value();
1164                    if s != "password" {
1165                        parser_err!("expected PASSWORD after ENCRYPTED, found {}", option);
1166                    }
1167                    (
1168                        &mut builder.password,
1169                        UserOption::EncryptedPassword(AstString::parse_to(parser)?),
1170                    )
1171                }
1172                "oauth" => {
1173                    let options = parser.parse_options()?;
1174                    (&mut builder.password, UserOption::OAuth(options))
1175                }
1176                _ => {
1177                    parser_err!("unexpected user option: {}", option);
1178                }
1179            };
1180            add_option(item_mut_ref, user_option)?;
1181        }
1182        Ok(builder.build())
1183    }
1184}
1185
1186impl fmt::Display for UserOptions {
1187    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1188        if !self.0.is_empty() {
1189            if self.1 {
1190                write!(f, "WITH ")?;
1191            }
1192            write!(f, "{}", display_separated(self.0.as_slice(), " "))
1193        } else {
1194            Ok(())
1195        }
1196    }
1197}
1198
1199impl ParseTo for CreateUserStatement {
1200    fn parse_to(p: &mut Parser<'_>) -> ModalResult<Self> {
1201        impl_parse_to!(user_name: ObjectName, p);
1202        impl_parse_to!(with_options: UserOptions, p);
1203
1204        Ok(CreateUserStatement {
1205            user_name,
1206            with_options,
1207        })
1208    }
1209}
1210
1211impl fmt::Display for CreateUserStatement {
1212    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1213        let mut v: Vec<String> = vec![];
1214        impl_fmt_display!(user_name, v, self);
1215        impl_fmt_display!(with_options, v, self);
1216        v.iter().join(" ").fmt(f)
1217    }
1218}
1219
1220impl fmt::Display for AlterUserMode {
1221    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1222        match self {
1223            AlterUserMode::Options(options) => {
1224                write!(f, "{}", options)
1225            }
1226            AlterUserMode::Rename(new_name) => {
1227                write!(f, "RENAME TO {}", new_name)
1228            }
1229        }
1230    }
1231}
1232
1233impl fmt::Display for AlterUserStatement {
1234    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1235        let mut v: Vec<String> = vec![];
1236        impl_fmt_display!(user_name, v, self);
1237        impl_fmt_display!(mode, v, self);
1238        v.iter().join(" ").fmt(f)
1239    }
1240}
1241
1242impl ParseTo for AlterUserStatement {
1243    fn parse_to(p: &mut Parser<'_>) -> ModalResult<Self> {
1244        impl_parse_to!(user_name: ObjectName, p);
1245        impl_parse_to!(mode: AlterUserMode, p);
1246
1247        Ok(AlterUserStatement { user_name, mode })
1248    }
1249}
1250
1251impl ParseTo for AlterUserMode {
1252    fn parse_to(p: &mut Parser<'_>) -> ModalResult<Self> {
1253        if p.parse_keyword(Keyword::RENAME) {
1254            p.expect_keyword(Keyword::TO)?;
1255            impl_parse_to!(new_name: ObjectName, p);
1256            Ok(AlterUserMode::Rename(new_name))
1257        } else {
1258            impl_parse_to!(with_options: UserOptions, p);
1259            Ok(AlterUserMode::Options(with_options))
1260        }
1261    }
1262}
1263
1264#[derive(Debug, Clone, PartialEq, Eq, Hash)]
1265pub struct DropStatement {
1266    /// The type of the object to drop: TABLE, VIEW, etc.
1267    pub object_type: ObjectType,
1268    /// An optional `IF EXISTS` clause. (Non-standard.)
1269    pub if_exists: bool,
1270    /// Object to drop.
1271    pub object_name: ObjectName,
1272    /// Whether `CASCADE` was specified. This will be `false` when
1273    /// `RESTRICT` or no drop behavior at all was specified.
1274    pub drop_mode: AstOption<DropMode>,
1275}
1276
1277// sql_grammar!(DropStatement {
1278//     object_type: ObjectType,
1279//     if_exists => [Keyword::IF, Keyword::EXISTS],
1280//     name: ObjectName,
1281//     drop_mode: AstOption<DropMode>,
1282// });
1283impl ParseTo for DropStatement {
1284    fn parse_to(p: &mut Parser<'_>) -> ModalResult<Self> {
1285        impl_parse_to!(object_type: ObjectType, p);
1286        impl_parse_to!(if_exists => [Keyword::IF, Keyword::EXISTS], p);
1287        let object_name = p.parse_object_name()?;
1288        impl_parse_to!(drop_mode: AstOption<DropMode>, p);
1289        Ok(Self {
1290            object_type,
1291            if_exists,
1292            object_name,
1293            drop_mode,
1294        })
1295    }
1296}
1297
1298impl fmt::Display for DropStatement {
1299    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1300        let mut v: Vec<String> = vec![];
1301        impl_fmt_display!(object_type, v, self);
1302        impl_fmt_display!(if_exists => [Keyword::IF, Keyword::EXISTS], v, self);
1303        impl_fmt_display!(object_name, v, self);
1304        impl_fmt_display!(drop_mode, v, self);
1305        v.iter().join(" ").fmt(f)
1306    }
1307}
1308
1309#[derive(Debug, Clone, PartialEq, Eq, Hash)]
1310pub enum DropMode {
1311    Cascade,
1312    Restrict,
1313}
1314
1315impl ParseTo for DropMode {
1316    fn parse_to(parser: &mut Parser<'_>) -> ModalResult<Self> {
1317        let drop_mode = if parser.parse_keyword(Keyword::CASCADE) {
1318            DropMode::Cascade
1319        } else if parser.parse_keyword(Keyword::RESTRICT) {
1320            DropMode::Restrict
1321        } else {
1322            return parser.expected("CASCADE | RESTRICT");
1323        };
1324        Ok(drop_mode)
1325    }
1326}
1327
1328impl fmt::Display for DropMode {
1329    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1330        f.write_str(match self {
1331            DropMode::Cascade => "CASCADE",
1332            DropMode::Restrict => "RESTRICT",
1333        })
1334    }
1335}