1#[cfg(feature = "enterprise")]
16pub mod trigger;
17
18use std::collections::{HashMap, HashSet};
19
20use api::helper::ColumnDataTypeWrapper;
21use api::v1::alter_database_expr::Kind as AlterDatabaseKind;
22use api::v1::alter_table_expr::Kind as AlterTableKind;
23use api::v1::column_def::{options_from_column_schema, try_as_column_schema};
24use api::v1::{
25 AddColumn, AddColumns, AlterDatabaseExpr, AlterTableExpr, Analyzer, ColumnDataType,
26 ColumnDataTypeExtension, CreateFlowExpr, CreateTableExpr, CreateViewExpr, DropColumn,
27 DropColumns, DropDefaults, ExpireAfter, FulltextBackend as PbFulltextBackend,
28 JsonSettings as PbJsonSettings, JsonTypeHint as PbJsonTypeHint, ModifyColumnType,
29 ModifyColumnTypes, RenameTable, SemanticType, SetDatabaseOptions, SetDefaults, SetFulltext,
30 SetIndex, SetIndexes, SetInverted, SetJsonSettings, SetSkipping, SetTableOptions,
31 SkippingIndexType as PbSkippingIndexType, TableName, UnsetDatabaseOptions, UnsetFulltext,
32 UnsetIndex, UnsetIndexes, UnsetInverted, UnsetSkipping, UnsetTableOptions, set_index,
33 unset_index,
34};
35use common_datasource::object_store::LocalFileAccess;
36use common_error::ext::BoxedError;
37use common_grpc_expr::util::ColumnExpr;
38use common_time::Timezone;
39use datafusion::sql::planner::object_name_to_table_reference;
40use datatypes::json::JsonSettings;
41use datatypes::prelude::ConcreteDataType;
42use datatypes::schema::{
43 COLUMN_FULLTEXT_OPT_KEY_ANALYZER, COLUMN_FULLTEXT_OPT_KEY_BACKEND,
44 COLUMN_FULLTEXT_OPT_KEY_CASE_SENSITIVE, COLUMN_FULLTEXT_OPT_KEY_FALSE_POSITIVE_RATE,
45 COLUMN_FULLTEXT_OPT_KEY_GRANULARITY, COLUMN_SKIPPING_INDEX_OPT_KEY_FALSE_POSITIVE_RATE,
46 COLUMN_SKIPPING_INDEX_OPT_KEY_GRANULARITY, COLUMN_SKIPPING_INDEX_OPT_KEY_TYPE, COMMENT_KEY,
47 ColumnDefaultConstraint, ColumnSchema, FulltextAnalyzer, FulltextBackend, Schema,
48 SkippingIndexType,
49};
50use file_engine::FileOptions;
51use query::sql::{
52 check_file_to_table_schema_compatibility, file_column_schemas_to_table,
53 infer_file_table_schema, prepare_file_table_files,
54};
55use session::context::QueryContextRef;
56use session::table_name::table_idents_to_full_name;
57use snafu::{OptionExt, ResultExt, ensure};
58use sql::ast::{
59 ColumnDef, ColumnOption, ColumnOptionDef, Expr, Ident, ObjectName, ObjectNamePartExt,
60};
61use sql::dialect::GreptimeDbDialect;
62use sql::parser::ParserContext;
63use sql::statements::alter::{
64 AlterDatabase, AlterDatabaseOperation, AlterTable, AlterTableOperation,
65};
66use sql::statements::create::{
67 Column as SqlColumn, ColumnExtensions, CreateExternalTable, CreateFlow, CreateTable,
68 CreateView, TableConstraint,
69};
70use sql::statements::{
71 OptionMap, column_to_schema, concrete_data_type_to_sql_data_type,
72 sql_column_def_to_grpc_column_def, sql_data_type_to_concrete_data_type, value_to_sql_value,
73};
74use sql::util::extract_tables_from_query;
75use store_api::mito_engine_options::{COMPACTION_OVERRIDE, COMPACTION_TYPE};
76use table::requests::{FILE_TABLE_META_KEY, TableOptions};
77use table::table_reference::TableReference;
78#[cfg(feature = "enterprise")]
79pub use trigger::to_create_trigger_task_expr;
80
81use crate::error::{
82 BuildCreateExprOnInsertionSnafu, ColumnDataTypeSnafu, ConvertColumnDefaultConstraintSnafu,
83 ConvertIdentifierSnafu, EncodeJsonSnafu, ExternalSnafu, FindNewColumnsOnInsertionSnafu,
84 IllegalPrimaryKeysDefSnafu, InferFileTableSchemaSnafu, InvalidColumnDefSnafu,
85 InvalidFlowNameSnafu, InvalidSqlSnafu, NotSupportedSnafu, ParseSqlSnafu, ParseSqlValueSnafu,
86 PrepareFileTableSnafu, Result, SchemaIncompatibleSnafu, UnrecognizedTableOptionSnafu,
87};
88
89pub fn create_table_expr_by_column_schemas(
90 table_name: &TableReference<'_>,
91 column_schemas: &[api::v1::ColumnSchema],
92 engine: &str,
93 desc: Option<&str>,
94) -> Result<CreateTableExpr> {
95 let column_exprs = ColumnExpr::from_column_schemas(column_schemas);
96 let expr = common_grpc_expr::util::build_create_table_expr(
97 None,
98 table_name,
99 column_exprs,
100 engine,
101 desc.unwrap_or("Created on insertion"),
102 )
103 .context(BuildCreateExprOnInsertionSnafu)?;
104
105 validate_create_expr(&expr)?;
106 Ok(expr)
107}
108
109pub fn extract_add_columns_expr(
110 schema: &Schema,
111 column_exprs: Vec<ColumnExpr>,
112) -> Result<Option<AddColumns>> {
113 let add_columns = common_grpc_expr::util::extract_new_columns(schema, column_exprs)
114 .context(FindNewColumnsOnInsertionSnafu)?;
115 if let Some(add_columns) = &add_columns {
116 validate_add_columns_expr(add_columns)?;
117 }
118 Ok(add_columns)
119}
120
121pub(crate) async fn create_external_expr(
145 create: CreateExternalTable,
146 query_ctx: &QueryContextRef,
147 local_file_access: &LocalFileAccess,
148) -> Result<CreateTableExpr> {
149 let (catalog_name, schema_name, table_name) =
150 table_idents_to_full_name(&create.name, query_ctx)
151 .map_err(BoxedError::new)
152 .context(ExternalSnafu)?;
153
154 let mut table_options = create.options.into_map();
155
156 let (object_store, files) = prepare_file_table_files(&table_options, local_file_access)
157 .await
158 .context(PrepareFileTableSnafu)?;
159
160 let file_column_schemas = infer_file_table_schema(&object_store, &files, &table_options)
161 .await
162 .context(InferFileTableSchemaSnafu)?
163 .column_schemas()
164 .to_vec();
165
166 let (time_index, primary_keys, table_column_schemas) = if !create.columns.is_empty() {
167 let time_index = find_time_index(&create.constraints)?;
169 let primary_keys = find_primary_keys(&create.columns, &create.constraints)?;
170 let column_schemas =
171 columns_to_column_schemas(&create.columns, &time_index, Some(&query_ctx.timezone()))?;
172 (time_index, primary_keys, column_schemas)
173 } else {
174 let (column_schemas, time_index) = file_column_schemas_to_table(&file_column_schemas);
176 let primary_keys = vec![];
177 (time_index, primary_keys, column_schemas)
178 };
179
180 check_file_to_table_schema_compatibility(&file_column_schemas, &table_column_schemas)
181 .context(SchemaIncompatibleSnafu)?;
182
183 let meta = FileOptions {
184 files,
185 file_column_schemas,
186 };
187 table_options.insert(
188 FILE_TABLE_META_KEY.to_string(),
189 serde_json::to_string(&meta).context(EncodeJsonSnafu)?,
190 );
191
192 let column_defs = column_schemas_to_defs(table_column_schemas, &primary_keys)?;
193 let expr = CreateTableExpr {
194 catalog_name,
195 schema_name,
196 table_name,
197 desc: String::default(),
198 column_defs,
199 time_index,
200 primary_keys,
201 create_if_not_exists: create.if_not_exists,
202 table_options,
203 table_id: None,
204 engine: create.engine.clone(),
205 };
206
207 Ok(expr)
208}
209
210pub fn create_to_expr(
212 create: &CreateTable,
213 query_ctx: &QueryContextRef,
214) -> Result<CreateTableExpr> {
215 let (catalog_name, schema_name, table_name) =
216 table_idents_to_full_name(&create.name, query_ctx)
217 .map_err(BoxedError::new)
218 .context(ExternalSnafu)?;
219
220 let time_index = find_time_index(&create.constraints)?;
221 let mut table_options = HashMap::from(
222 &TableOptions::try_from_iter(create.options.to_str_map())
223 .context(UnrecognizedTableOptionSnafu)?,
224 );
225
226 if table_options.contains_key(COMPACTION_TYPE) {
227 table_options.insert(COMPACTION_OVERRIDE.to_string(), "true".to_string());
228 }
229
230 let primary_keys = find_primary_keys(&create.columns, &create.constraints)?;
231
232 let expr = CreateTableExpr {
233 catalog_name,
234 schema_name,
235 table_name,
236 desc: String::default(),
237 column_defs: columns_to_expr(
238 &create.columns,
239 &time_index,
240 &primary_keys,
241 Some(&query_ctx.timezone()),
242 )?,
243 time_index,
244 primary_keys,
245 create_if_not_exists: create.if_not_exists,
246 table_options,
247 table_id: None,
248 engine: create.engine.clone(),
249 };
250
251 validate_create_expr(&expr)?;
252 Ok(expr)
253}
254
255pub fn expr_to_create(expr: &CreateTableExpr, quote_style: Option<char>) -> Result<CreateTable> {
263 let quote_style = quote_style.unwrap_or('`');
264
265 let table_name = ObjectName(vec![sql::ast::ObjectNamePart::Identifier(
267 sql::ast::Ident::with_quote(quote_style, &expr.table_name),
268 )]);
269
270 let mut columns = Vec::with_capacity(expr.column_defs.len());
272 for column_def in &expr.column_defs {
273 let column_schema = try_as_column_schema(column_def).context(InvalidColumnDefSnafu {
274 column: &column_def.name,
275 })?;
276
277 let mut options = Vec::new();
278
279 if column_def.is_nullable {
281 options.push(ColumnOptionDef {
282 name: None,
283 option: ColumnOption::Null,
284 });
285 } else {
286 options.push(ColumnOptionDef {
287 name: None,
288 option: ColumnOption::NotNull,
289 });
290 }
291
292 if let Some(default_constraint) = column_schema.default_constraint() {
294 let expr = match default_constraint {
295 ColumnDefaultConstraint::Value(v) => {
296 Expr::Value(value_to_sql_value(v).context(ParseSqlValueSnafu)?.into())
297 }
298 ColumnDefaultConstraint::Function(func_expr) => {
299 ParserContext::parse_function(func_expr, &GreptimeDbDialect {})
300 .context(ParseSqlSnafu)?
301 }
302 };
303 options.push(ColumnOptionDef {
304 name: None,
305 option: ColumnOption::Default(expr),
306 });
307 }
308
309 if !column_def.comment.is_empty() {
311 options.push(ColumnOptionDef {
312 name: None,
313 option: ColumnOption::Comment(column_def.comment.clone()),
314 });
315 }
316
317 let mut extensions = ColumnExtensions::default();
322
323 if let Ok(Some(opt)) = column_schema.fulltext_options()
325 && opt.enable
326 {
327 let mut map = HashMap::from([
328 (
329 COLUMN_FULLTEXT_OPT_KEY_ANALYZER.to_string(),
330 opt.analyzer.to_string(),
331 ),
332 (
333 COLUMN_FULLTEXT_OPT_KEY_CASE_SENSITIVE.to_string(),
334 opt.case_sensitive.to_string(),
335 ),
336 (
337 COLUMN_FULLTEXT_OPT_KEY_BACKEND.to_string(),
338 opt.backend.to_string(),
339 ),
340 ]);
341 if opt.backend == FulltextBackend::Bloom {
342 map.insert(
343 COLUMN_FULLTEXT_OPT_KEY_GRANULARITY.to_string(),
344 opt.granularity.to_string(),
345 );
346 map.insert(
347 COLUMN_FULLTEXT_OPT_KEY_FALSE_POSITIVE_RATE.to_string(),
348 opt.false_positive_rate().to_string(),
349 );
350 }
351 extensions.fulltext_index_options = Some(map.into());
352 }
353
354 if let Ok(Some(opt)) = column_schema.skipping_index_options() {
356 let map = HashMap::from([
357 (
358 COLUMN_SKIPPING_INDEX_OPT_KEY_GRANULARITY.to_string(),
359 opt.granularity.to_string(),
360 ),
361 (
362 COLUMN_SKIPPING_INDEX_OPT_KEY_FALSE_POSITIVE_RATE.to_string(),
363 opt.false_positive_rate().to_string(),
364 ),
365 (
366 COLUMN_SKIPPING_INDEX_OPT_KEY_TYPE.to_string(),
367 opt.index_type.to_string(),
368 ),
369 ]);
370 extensions.skipping_index_options = Some(map.into());
371 }
372
373 if column_schema.is_inverted_indexed() {
375 extensions.inverted_index_options = Some(HashMap::new().into());
376 }
377
378 let sql_column = SqlColumn {
379 column_def: ColumnDef {
380 name: Ident::with_quote(quote_style, &column_def.name),
381 data_type: concrete_data_type_to_sql_data_type(&column_schema.data_type)
382 .context(ParseSqlSnafu)?,
383 options,
384 },
385 extensions,
386 };
387
388 columns.push(sql_column);
389 }
390
391 let mut constraints = Vec::new();
393
394 constraints.push(TableConstraint::TimeIndex {
396 column: Ident::with_quote(quote_style, &expr.time_index),
397 });
398
399 if !expr.primary_keys.is_empty() {
401 let primary_key_columns: Vec<Ident> = expr
402 .primary_keys
403 .iter()
404 .map(|pk| Ident::with_quote(quote_style, pk))
405 .collect();
406
407 constraints.push(TableConstraint::PrimaryKey {
408 columns: primary_key_columns,
409 });
410 }
411
412 let mut options = OptionMap::default();
414 for (key, value) in &expr.table_options {
415 options.insert(key.clone(), value.clone());
416 }
417
418 Ok(CreateTable {
419 if_not_exists: expr.create_if_not_exists,
420 table_id: expr.table_id.as_ref().map(|tid| tid.id).unwrap_or(0),
421 name: table_name,
422 columns,
423 engine: expr.engine.clone(),
424 constraints,
425 options,
426 partitions: None,
427 })
428}
429
430pub fn validate_create_expr(create: &CreateTableExpr) -> Result<()> {
432 let mut column_to_indices = HashMap::with_capacity(create.column_defs.len());
434 for (idx, column) in create.column_defs.iter().enumerate() {
435 if let Some(indices) = column_to_indices.get(&column.name) {
436 return InvalidSqlSnafu {
437 err_msg: format!(
438 "column name `{}` is duplicated at index {} and {}",
439 column.name, indices, idx
440 ),
441 }
442 .fail();
443 }
444 column_to_indices.insert(&column.name, idx);
445 }
446
447 let time_index_idx =
449 column_to_indices
450 .get(&create.time_index)
451 .with_context(|| InvalidSqlSnafu {
452 err_msg: format!(
453 "column name `{}` is not found in column list",
454 create.time_index
455 ),
456 })?;
457
458 let time_index_column = &create.column_defs[*time_index_idx];
460 let data_type = ConcreteDataType::from(
461 ColumnDataTypeWrapper::try_new(
462 time_index_column.data_type,
463 time_index_column.datatype_extension.clone(),
464 )
465 .context(ColumnDataTypeSnafu)?,
466 );
467 ensure!(
468 data_type.is_timestamp(),
469 InvalidSqlSnafu {
470 err_msg: format!(
471 "column `{}` is not a timestamp type, it can't be used as time index",
472 create.time_index
473 ),
474 }
475 );
476
477 for pk in &create.primary_keys {
479 let _ = column_to_indices
480 .get(&pk)
481 .with_context(|| InvalidSqlSnafu {
482 err_msg: format!("column name `{}` is not found in column list", pk),
483 })?;
484 }
485
486 let mut pk_set = HashSet::new();
488 for pk in &create.primary_keys {
489 if !pk_set.insert(pk) {
490 return InvalidSqlSnafu {
491 err_msg: format!("column name `{}` is duplicated in primary keys", pk),
492 }
493 .fail();
494 }
495 }
496
497 if pk_set.contains(&create.time_index) {
499 return InvalidSqlSnafu {
500 err_msg: format!(
501 "column name `{}` is both primary key and time index",
502 create.time_index
503 ),
504 }
505 .fail();
506 }
507
508 for column in &create.column_defs {
509 if is_interval_type(&column.data_type()) {
511 return InvalidSqlSnafu {
512 err_msg: format!(
513 "column name `{}` is interval type, which is not supported",
514 column.name
515 ),
516 }
517 .fail();
518 }
519 if is_date_time_type(&column.data_type()) {
521 return InvalidSqlSnafu {
522 err_msg: format!(
523 "column name `{}` is datetime type, which is not supported, please use `timestamp` type instead",
524 column.name
525 ),
526 }
527 .fail();
528 }
529 }
530 Ok(())
531}
532
533fn validate_add_columns_expr(add_columns: &AddColumns) -> Result<()> {
534 for add_column in &add_columns.add_columns {
535 let Some(column_def) = &add_column.column_def else {
536 continue;
537 };
538 if is_date_time_type(&column_def.data_type()) {
539 return InvalidSqlSnafu {
540 err_msg: format!("column name `{}` is datetime type, which is not supported, please use `timestamp` type instead", column_def.name),
541 }
542 .fail();
543 }
544 if is_interval_type(&column_def.data_type()) {
545 return InvalidSqlSnafu {
546 err_msg: format!(
547 "column name `{}` is interval type, which is not supported",
548 column_def.name
549 ),
550 }
551 .fail();
552 }
553 }
554 Ok(())
555}
556
557fn is_date_time_type(data_type: &ColumnDataType) -> bool {
558 matches!(data_type, ColumnDataType::Datetime)
559}
560
561fn is_interval_type(data_type: &ColumnDataType) -> bool {
562 matches!(
563 data_type,
564 ColumnDataType::IntervalYearMonth
565 | ColumnDataType::IntervalDayTime
566 | ColumnDataType::IntervalMonthDayNano
567 )
568}
569
570fn find_primary_keys(
571 columns: &[SqlColumn],
572 constraints: &[TableConstraint],
573) -> Result<Vec<String>> {
574 let columns_pk = columns
575 .iter()
576 .filter_map(|x| {
577 if x.options()
578 .iter()
579 .any(|o| matches!(o.option, ColumnOption::PrimaryKey(_)))
580 {
581 Some(x.name().value.clone())
582 } else {
583 None
584 }
585 })
586 .collect::<Vec<String>>();
587
588 ensure!(
589 columns_pk.len() <= 1,
590 IllegalPrimaryKeysDefSnafu {
591 msg: "not allowed to inline multiple primary keys in columns options"
592 }
593 );
594
595 let constraints_pk = constraints
596 .iter()
597 .filter_map(|constraint| match constraint {
598 TableConstraint::PrimaryKey { columns, .. } => {
599 Some(columns.iter().map(|ident| ident.value.clone()))
600 }
601 _ => None,
602 })
603 .flatten()
604 .collect::<Vec<String>>();
605
606 ensure!(
607 columns_pk.is_empty() || constraints_pk.is_empty(),
608 IllegalPrimaryKeysDefSnafu {
609 msg: "found definitions of primary keys in multiple places"
610 }
611 );
612
613 let mut primary_keys = Vec::with_capacity(columns_pk.len() + constraints_pk.len());
614 primary_keys.extend(columns_pk);
615 primary_keys.extend(constraints_pk);
616 Ok(primary_keys)
617}
618
619pub fn find_time_index(constraints: &[TableConstraint]) -> Result<String> {
620 let time_index = constraints
621 .iter()
622 .filter_map(|constraint| match constraint {
623 TableConstraint::TimeIndex { column, .. } => Some(&column.value),
624 _ => None,
625 })
626 .collect::<Vec<&String>>();
627 ensure!(
628 time_index.len() == 1,
629 InvalidSqlSnafu {
630 err_msg: "must have one and only one TimeIndex columns",
631 }
632 );
633 Ok(time_index[0].clone())
634}
635
636fn columns_to_expr(
637 column_defs: &[SqlColumn],
638 time_index: &str,
639 primary_keys: &[String],
640 timezone: Option<&Timezone>,
641) -> Result<Vec<api::v1::ColumnDef>> {
642 let column_schemas = columns_to_column_schemas(column_defs, time_index, timezone)?;
643 column_schemas_to_defs(column_schemas, primary_keys)
644}
645
646fn columns_to_column_schemas(
647 columns: &[SqlColumn],
648 time_index: &str,
649 timezone: Option<&Timezone>,
650) -> Result<Vec<ColumnSchema>> {
651 columns
652 .iter()
653 .map(|c| column_to_schema(c, time_index, timezone).context(ParseSqlSnafu))
654 .collect::<Result<Vec<ColumnSchema>>>()
655}
656
657pub fn column_schemas_to_defs(
659 column_schemas: Vec<ColumnSchema>,
660 primary_keys: &[String],
661) -> Result<Vec<api::v1::ColumnDef>> {
662 let column_datatypes: Vec<(ColumnDataType, Option<ColumnDataTypeExtension>)> = column_schemas
663 .iter()
664 .map(|c| {
665 ColumnDataTypeWrapper::try_from(c.data_type.clone())
666 .map(|w| w.to_parts())
667 .context(ColumnDataTypeSnafu)
668 })
669 .collect::<Result<Vec<_>>>()?;
670
671 column_schemas
672 .iter()
673 .zip(column_datatypes)
674 .map(|(schema, datatype)| {
675 let semantic_type = if schema.is_time_index() {
676 SemanticType::Timestamp
677 } else if primary_keys.contains(&schema.name) {
678 SemanticType::Tag
679 } else {
680 SemanticType::Field
681 } as i32;
682 let comment = schema
683 .metadata()
684 .get(COMMENT_KEY)
685 .cloned()
686 .unwrap_or_default();
687
688 Ok(api::v1::ColumnDef {
689 name: schema.name.clone(),
690 data_type: datatype.0 as i32,
691 is_nullable: schema.is_nullable(),
692 default_constraint: match schema.default_constraint() {
693 None => vec![],
694 Some(v) => {
695 v.clone()
696 .try_into()
697 .context(ConvertColumnDefaultConstraintSnafu {
698 column_name: &schema.name,
699 })?
700 }
701 },
702 semantic_type,
703 comment,
704 datatype_extension: datatype.1,
705 options: options_from_column_schema(schema),
706 })
707 })
708 .collect()
709}
710
711#[derive(Debug, Clone, PartialEq, Eq)]
712pub struct RepartitionRequest {
713 pub catalog_name: String,
714 pub schema_name: String,
715 pub table_name: String,
716 pub source: RepartitionSource,
717 pub into_exprs: Vec<Expr>,
718 pub options: OptionMap,
719}
720
721#[derive(Debug, Clone, PartialEq, Eq)]
722pub enum RepartitionSource {
723 Partitions {
724 from_exprs: Vec<Expr>,
725 target_partition_columns: Option<Vec<String>>,
726 },
727 Unpartitioned {
728 partition_columns: Vec<String>,
729 },
730}
731
732pub(crate) fn to_repartition_request(
733 alter_table: AlterTable,
734 query_ctx: &QueryContextRef,
735) -> Result<RepartitionRequest> {
736 let AlterTable {
737 table_name,
738 alter_operation,
739 options,
740 } = alter_table;
741
742 let (catalog_name, schema_name, table_name) = table_idents_to_full_name(&table_name, query_ctx)
743 .map_err(BoxedError::new)
744 .context(ExternalSnafu)?;
745
746 let (source, into_exprs) = match alter_operation {
747 AlterTableOperation::Repartition { operation } => (
748 RepartitionSource::Partitions {
749 from_exprs: operation.from_exprs,
750 target_partition_columns: operation.partition_columns.map(|columns| {
751 columns
752 .into_iter()
753 .map(|ident| ident.value)
754 .collect::<Vec<_>>()
755 }),
756 },
757 operation.into_exprs,
758 ),
759 AlterTableOperation::Partition { partitions } => (
760 RepartitionSource::Unpartitioned {
761 partition_columns: partitions
762 .column_list
763 .into_iter()
764 .map(|ident| ident.value)
765 .collect(),
766 },
767 partitions.exprs,
768 ),
769 _ => {
770 return InvalidSqlSnafu {
771 err_msg: "expected REPARTITION or PARTITION operation",
772 }
773 .fail();
774 }
775 };
776
777 Ok(RepartitionRequest {
778 catalog_name,
779 schema_name,
780 table_name,
781 source,
782 into_exprs,
783 options,
784 })
785}
786
787fn json_settings_to_proto(settings: JsonSettings) -> Result<PbJsonSettings> {
788 let (type_hints, max_auto_expanded_paths) = settings.into_parts();
789 let type_hints = type_hints
790 .into_iter()
791 .map(|hint| {
792 let (data_type, datatype_extension) = ColumnDataTypeWrapper::try_from(hint.data_type)
793 .map(|w| w.to_parts())
794 .context(ColumnDataTypeSnafu)?;
795
796 Ok(PbJsonTypeHint {
797 path: hint.path,
798 data_type: data_type as i32,
799 datatype_extension,
800 })
801 })
802 .collect::<Result<Vec<_>>>()?;
803
804 Ok(PbJsonSettings {
805 type_hints,
806 max_auto_expanded_paths,
807 })
808}
809
810pub(crate) fn to_alter_table_expr(
812 alter_table: AlterTable,
813 query_ctx: &QueryContextRef,
814) -> Result<AlterTableExpr> {
815 let (catalog_name, schema_name, table_name) =
816 table_idents_to_full_name(alter_table.table_name(), query_ctx)
817 .map_err(BoxedError::new)
818 .context(ExternalSnafu)?;
819
820 let kind = match alter_table.alter_operation {
821 AlterTableOperation::AddConstraint(_) => {
822 return NotSupportedSnafu {
823 feat: "ADD CONSTRAINT",
824 }
825 .fail();
826 }
827 AlterTableOperation::AddColumns { add_columns } => AlterTableKind::AddColumns(AddColumns {
828 add_columns: add_columns
829 .into_iter()
830 .map(|add_column| {
831 let column_def = sql_column_def_to_grpc_column_def(
832 &add_column.column_def,
833 Some(&query_ctx.timezone()),
834 )
835 .map_err(BoxedError::new)
836 .context(ExternalSnafu)?;
837 if is_interval_type(&column_def.data_type()) {
838 return NotSupportedSnafu {
839 feat: "Add column with interval type",
840 }
841 .fail();
842 }
843 Ok(AddColumn {
844 column_def: Some(column_def),
845 location: add_column.location.as_ref().map(From::from),
846 add_if_not_exists: add_column.add_if_not_exists,
847 })
848 })
849 .collect::<Result<Vec<AddColumn>>>()?,
850 }),
851 AlterTableOperation::ModifyColumnType {
852 column_name,
853 target_type,
854 ..
855 } => {
856 let target_type =
857 sql_data_type_to_concrete_data_type(&target_type).context(ParseSqlSnafu)?;
858
859 if target_type.is_json2() {
861 return NotSupportedSnafu {
862 feat: "ALTER TABLE MODIFY COLUMN to JSON2 type",
863 }
864 .fail();
865 }
866
867 let (target_type, target_type_extension) = ColumnDataTypeWrapper::try_from(target_type)
868 .map(|w| w.to_parts())
869 .context(ColumnDataTypeSnafu)?;
870 if is_interval_type(&target_type) {
871 return NotSupportedSnafu {
872 feat: "Modify column type to interval type",
873 }
874 .fail();
875 }
876 AlterTableKind::ModifyColumnTypes(ModifyColumnTypes {
877 modify_column_types: vec![ModifyColumnType {
878 column_name: column_name.value,
879 target_type: target_type as i32,
880 target_type_extension,
881 }],
882 })
883 }
884 AlterTableOperation::SetJsonSettings {
885 column_name,
886 json2_options,
887 } => {
888 let settings = match json2_options {
889 Some(options) => options.build_json_settings().context(ParseSqlSnafu)?,
890 None => datatypes::json::JsonSettings::new_v2(),
891 };
892
893 AlterTableKind::SetJsonSettings(SetJsonSettings {
894 column_name: column_name.value,
895 settings: Some(json_settings_to_proto(settings)?),
896 })
897 }
898 AlterTableOperation::DropColumn { name } => AlterTableKind::DropColumns(DropColumns {
899 drop_columns: vec![DropColumn {
900 name: name.value.clone(),
901 }],
902 }),
903 AlterTableOperation::RenameTable { new_table_name } => {
904 AlterTableKind::RenameTable(RenameTable {
905 new_table_name: new_table_name.clone(),
906 })
907 }
908 AlterTableOperation::SetTableOptions { options } => {
909 AlterTableKind::SetTableOptions(SetTableOptions {
910 table_options: options.into_iter().map(Into::into).collect(),
911 })
912 }
913 AlterTableOperation::UnsetTableOptions { keys } => {
914 AlterTableKind::UnsetTableOptions(UnsetTableOptions { keys })
915 }
916 AlterTableOperation::Repartition { .. } => {
917 return NotSupportedSnafu {
918 feat: "ALTER TABLE ... REPARTITION",
919 }
920 .fail();
921 }
922 AlterTableOperation::Partition { .. } => {
923 return NotSupportedSnafu {
924 feat: "ALTER TABLE ... PARTITION ON COLUMNS",
925 }
926 .fail();
927 }
928 AlterTableOperation::SetIndex { options } => {
929 let option = match options {
930 sql::statements::alter::SetIndexOperation::Fulltext {
931 column_name,
932 options,
933 } => SetIndex {
934 options: Some(set_index::Options::Fulltext(SetFulltext {
935 column_name: column_name.value,
936 enable: options.enable,
937 analyzer: match options.analyzer {
938 FulltextAnalyzer::English => Analyzer::English.into(),
939 FulltextAnalyzer::Chinese => Analyzer::Chinese.into(),
940 },
941 case_sensitive: options.case_sensitive,
942 backend: match options.backend {
943 FulltextBackend::Bloom => PbFulltextBackend::Bloom.into(),
944 FulltextBackend::Tantivy => PbFulltextBackend::Tantivy.into(),
945 },
946 granularity: options.granularity as u64,
947 false_positive_rate: options.false_positive_rate(),
948 })),
949 },
950 sql::statements::alter::SetIndexOperation::Inverted { column_name } => SetIndex {
951 options: Some(set_index::Options::Inverted(SetInverted {
952 column_name: column_name.value,
953 })),
954 },
955 sql::statements::alter::SetIndexOperation::Skipping {
956 column_name,
957 options,
958 } => SetIndex {
959 options: Some(set_index::Options::Skipping(SetSkipping {
960 column_name: column_name.value,
961 enable: true,
962 granularity: options.granularity as u64,
963 false_positive_rate: options.false_positive_rate(),
964 skipping_index_type: match options.index_type {
965 SkippingIndexType::BloomFilter => {
966 PbSkippingIndexType::BloomFilter.into()
967 }
968 },
969 })),
970 },
971 };
972 AlterTableKind::SetIndexes(SetIndexes {
973 set_indexes: vec![option],
974 })
975 }
976 AlterTableOperation::UnsetIndex { options } => {
977 let option = match options {
978 sql::statements::alter::UnsetIndexOperation::Fulltext { column_name } => {
979 UnsetIndex {
980 options: Some(unset_index::Options::Fulltext(UnsetFulltext {
981 column_name: column_name.value,
982 })),
983 }
984 }
985 sql::statements::alter::UnsetIndexOperation::Inverted { column_name } => {
986 UnsetIndex {
987 options: Some(unset_index::Options::Inverted(UnsetInverted {
988 column_name: column_name.value,
989 })),
990 }
991 }
992 sql::statements::alter::UnsetIndexOperation::Skipping { column_name } => {
993 UnsetIndex {
994 options: Some(unset_index::Options::Skipping(UnsetSkipping {
995 column_name: column_name.value,
996 })),
997 }
998 }
999 };
1000
1001 AlterTableKind::UnsetIndexes(UnsetIndexes {
1002 unset_indexes: vec![option],
1003 })
1004 }
1005 AlterTableOperation::DropDefaults { columns } => {
1006 AlterTableKind::DropDefaults(DropDefaults {
1007 drop_defaults: columns
1008 .into_iter()
1009 .map(|col| {
1010 let column_name = col.0.to_string();
1011 Ok(api::v1::DropDefault { column_name })
1012 })
1013 .collect::<Result<Vec<_>>>()?,
1014 })
1015 }
1016 AlterTableOperation::SetDefaults { defaults } => AlterTableKind::SetDefaults(SetDefaults {
1017 set_defaults: defaults
1018 .into_iter()
1019 .map(|col| {
1020 let column_name = col.column_name.to_string();
1021 let default_constraint = serde_json::to_string(&col.default_constraint)
1022 .context(EncodeJsonSnafu)?
1023 .into_bytes();
1024 Ok(api::v1::SetDefault {
1025 column_name,
1026 default_constraint,
1027 })
1028 })
1029 .collect::<Result<Vec<_>>>()?,
1030 }),
1031 };
1032
1033 Ok(AlterTableExpr {
1034 catalog_name,
1035 schema_name,
1036 table_name,
1037 kind: Some(kind),
1038 })
1039}
1040
1041pub fn to_alter_database_expr(
1043 alter_database: AlterDatabase,
1044 query_ctx: &QueryContextRef,
1045) -> Result<AlterDatabaseExpr> {
1046 let catalog = query_ctx.current_catalog();
1047 let schema = alter_database.database_name;
1048
1049 let kind = match alter_database.alter_operation {
1050 AlterDatabaseOperation::SetDatabaseOption { options } => {
1051 let options = options.into_iter().map(Into::into).collect();
1052 AlterDatabaseKind::SetDatabaseOptions(SetDatabaseOptions {
1053 set_database_options: options,
1054 })
1055 }
1056 AlterDatabaseOperation::UnsetDatabaseOption { keys } => {
1057 AlterDatabaseKind::UnsetDatabaseOptions(UnsetDatabaseOptions { keys })
1058 }
1059 };
1060
1061 Ok(AlterDatabaseExpr {
1062 catalog_name: catalog.to_string(),
1063 schema_name: schema.to_string(),
1064 kind: Some(kind),
1065 })
1066}
1067
1068pub fn to_create_view_expr(
1070 stmt: CreateView,
1071 logical_plan: Vec<u8>,
1072 table_names: Vec<TableName>,
1073 columns: Vec<String>,
1074 plan_columns: Vec<String>,
1075 definition: String,
1076 query_ctx: QueryContextRef,
1077) -> Result<CreateViewExpr> {
1078 let (catalog_name, schema_name, view_name) = table_idents_to_full_name(&stmt.name, &query_ctx)
1079 .map_err(BoxedError::new)
1080 .context(ExternalSnafu)?;
1081
1082 let expr = CreateViewExpr {
1083 catalog_name,
1084 schema_name,
1085 view_name,
1086 logical_plan,
1087 create_if_not_exists: stmt.if_not_exists,
1088 or_replace: stmt.or_replace,
1089 table_names,
1090 columns,
1091 plan_columns,
1092 definition,
1093 };
1094
1095 Ok(expr)
1096}
1097
1098pub fn to_create_flow_task_expr(
1099 create_flow: CreateFlow,
1100 query_ctx: &QueryContextRef,
1101) -> Result<CreateFlowExpr> {
1102 let sink_table_ref = object_name_to_table_reference(create_flow.sink_table_name.clone(), true)
1104 .with_context(|_| ConvertIdentifierSnafu {
1105 ident: create_flow.sink_table_name.to_string(),
1106 })?;
1107 let catalog = sink_table_ref
1108 .catalog()
1109 .unwrap_or(query_ctx.current_catalog())
1110 .to_string();
1111 let schema = sink_table_ref
1112 .schema()
1113 .map(|s| s.to_owned())
1114 .unwrap_or(query_ctx.current_schema());
1115
1116 let sink_table_name = TableName {
1117 catalog_name: catalog,
1118 schema_name: schema,
1119 table_name: sink_table_ref.table().to_string(),
1120 };
1121
1122 let source_table_names = extract_tables_from_query(&create_flow.query)
1123 .map(|name| {
1124 let reference =
1125 object_name_to_table_reference(name.clone(), true).with_context(|_| {
1126 ConvertIdentifierSnafu {
1127 ident: name.to_string(),
1128 }
1129 })?;
1130 let catalog = reference
1131 .catalog()
1132 .unwrap_or(query_ctx.current_catalog())
1133 .to_string();
1134 let schema = reference
1135 .schema()
1136 .map(|s| s.to_string())
1137 .unwrap_or(query_ctx.current_schema());
1138
1139 let table_name = TableName {
1140 catalog_name: catalog,
1141 schema_name: schema,
1142 table_name: reference.table().to_string(),
1143 };
1144 Ok(table_name)
1145 })
1146 .collect::<Result<Vec<_>>>()?;
1147
1148 let eval_interval = create_flow.eval_interval;
1149
1150 let flow_options = stringify_flow_options(create_flow.flow_options)?;
1151 Ok(CreateFlowExpr {
1152 catalog_name: query_ctx.current_catalog().to_string(),
1153 flow_name: sanitize_flow_name(create_flow.flow_name)?,
1154 source_table_names,
1155 sink_table_name: Some(sink_table_name),
1156 or_replace: create_flow.or_replace,
1157 create_if_not_exists: create_flow.if_not_exists,
1158 expire_after: create_flow.expire_after.map(|value| ExpireAfter { value }),
1159 eval_interval: eval_interval.map(|seconds| api::v1::EvalInterval { seconds }),
1160 comment: create_flow.comment.unwrap_or_default(),
1161 sql: create_flow.query.to_string(),
1162 flow_options,
1163 })
1164}
1165
1166fn stringify_flow_options(flow_options: OptionMap) -> Result<HashMap<String, String>> {
1167 let options_len = flow_options.len();
1168 let flow_options = flow_options.into_map();
1169 ensure!(
1170 flow_options.len() == options_len,
1171 InvalidSqlSnafu {
1172 err_msg: "flow options only support scalar string-compatible values".to_string(),
1173 }
1174 );
1175 Ok(flow_options)
1176}
1177
1178fn sanitize_flow_name(mut flow_name: ObjectName) -> Result<String> {
1180 ensure!(
1181 flow_name.0.len() == 1,
1182 InvalidFlowNameSnafu {
1183 name: flow_name.to_string(),
1184 }
1185 );
1186 Ok(flow_name.0.swap_remove(0).to_string_unquoted())
1188}
1189
1190#[cfg(test)]
1191mod tests {
1192 use std::collections::HashMap;
1193
1194 use api::v1::{SetDatabaseOptions, UnsetDatabaseOptions};
1195 use datatypes::value::Value;
1196 use session::context::{QueryContext, QueryContextBuilder};
1197 use sql::dialect::GreptimeDbDialect;
1198 use sql::parser::{ParseOptions, ParserContext};
1199 use sql::statements::statement::Statement;
1200 use store_api::storage::ColumnDefaultConstraint;
1201
1202 use super::*;
1203
1204 #[test]
1205 fn test_create_flow_tql_expr() {
1206 let sql = r#"
1207CREATE FLOW calc_reqs SINK TO cnt_reqs AS
1208TQL EVAL (0, 15, '5s') count_values("status_code", http_requests);"#;
1209 let stmt =
1210 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default());
1211
1212 assert!(
1213 stmt.is_err(),
1214 "Expected error for invalid TQL EVAL parameters: {:#?}",
1215 stmt
1216 );
1217
1218 let sql = r#"
1219CREATE FLOW calc_reqs SINK TO cnt_reqs AS
1220TQL EVAL (now() - '15s'::interval, now(), '5s') count_values("status_code", http_requests);"#;
1221 let stmt =
1222 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default())
1223 .unwrap()
1224 .pop()
1225 .unwrap();
1226
1227 let Statement::CreateFlow(create_flow) = stmt else {
1228 unreachable!()
1229 };
1230 let expr = to_create_flow_task_expr(create_flow, &QueryContext::arc()).unwrap();
1231
1232 let to_dot_sep =
1233 |c: TableName| format!("{}.{}.{}", c.catalog_name, c.schema_name, c.table_name);
1234 assert_eq!("calc_reqs", expr.flow_name);
1235 assert_eq!("greptime", expr.catalog_name);
1236 assert_eq!(
1237 "greptime.public.cnt_reqs",
1238 expr.sink_table_name.map(to_dot_sep).unwrap()
1239 );
1240 assert_eq!(1, expr.source_table_names.len());
1241 assert_eq!(
1242 "greptime.public.http_requests",
1243 to_dot_sep(expr.source_table_names[0].clone())
1244 );
1245 assert_eq!(
1246 r#"TQL EVAL (now() - '15s'::interval, now(), '5s') count_values("status_code", http_requests)"#,
1247 expr.sql
1248 );
1249
1250 let sql = r#"
1251CREATE FLOW calc_reqs SINK TO cnt_reqs AS
1252TQL EVAL (now() - '15s'::interval, now(), '5s') count_values("status_code", http_requests{__schema__="greptime_private"});"#;
1253 let stmt =
1254 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default())
1255 .unwrap()
1256 .pop()
1257 .unwrap();
1258 let Statement::CreateFlow(create_flow) = stmt else {
1259 unreachable!()
1260 };
1261 let expr = to_create_flow_task_expr(create_flow, &QueryContext::arc()).unwrap();
1262 assert_eq!(1, expr.source_table_names.len());
1263 assert_eq!(
1264 "greptime.greptime_private.http_requests",
1265 to_dot_sep(expr.source_table_names[0].clone())
1266 );
1267
1268 let sql = r#"
1269CREATE FLOW calc_reqs SINK TO cnt_reqs AS
1270TQL EVAL (now() - '15s'::interval, now(), '5s') count_values("status_code", http_requests{__database__="greptime_private"});"#;
1271 let stmt =
1272 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default())
1273 .unwrap()
1274 .pop()
1275 .unwrap();
1276 let Statement::CreateFlow(create_flow) = stmt else {
1277 unreachable!()
1278 };
1279 let expr = to_create_flow_task_expr(create_flow, &QueryContext::arc()).unwrap();
1280 assert_eq!(1, expr.source_table_names.len());
1281 assert_eq!(
1282 "greptime.greptime_private.http_requests",
1283 to_dot_sep(expr.source_table_names[0].clone())
1284 );
1285 }
1286
1287 #[test]
1288 fn test_create_flow_tql_cte_source_tables() {
1289 let sql = r#"
1290CREATE FLOW calc_cte
1291SINK TO metric_cte_sink
1292EVAL INTERVAL '1m'
1293AS
1294WITH tql(ts, the_value) AS (
1295 TQL EVAL (now() - '1m'::interval, now(), '5s') metric_cte
1296)
1297SELECT * FROM tql;
1298"#;
1299
1300 let stmt =
1301 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default())
1302 .unwrap()
1303 .pop()
1304 .unwrap();
1305
1306 let Statement::CreateFlow(create_flow) = stmt else {
1307 unreachable!()
1308 };
1309 let expr = to_create_flow_task_expr(create_flow, &QueryContext::arc()).unwrap();
1310
1311 let to_dot_sep =
1312 |c: TableName| format!("{}.{}.{}", c.catalog_name, c.schema_name, c.table_name);
1313 assert_eq!(1, expr.source_table_names.len());
1314 assert_eq!(
1315 "greptime.public.metric_cte",
1316 to_dot_sep(expr.source_table_names[0].clone())
1317 );
1318 }
1319
1320 #[test]
1321 fn test_create_flow_tql_cte_source_tables_quoted_cte_name() {
1322 let sql = r#"
1323CREATE FLOW calc_cte
1324SINK TO metric_cte_sink
1325EVAL INTERVAL '1m'
1326AS
1327WITH "TQL"(ts, the_value) AS (
1328 TQL EVAL (now() - '1m'::interval, now(), '5s') metric_cte
1329)
1330SELECT * FROM "TQL";
1331"#;
1332
1333 let stmt =
1334 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default())
1335 .unwrap()
1336 .pop()
1337 .unwrap();
1338
1339 let Statement::CreateFlow(create_flow) = stmt else {
1340 unreachable!()
1341 };
1342 let expr = to_create_flow_task_expr(create_flow, &QueryContext::arc()).unwrap();
1343
1344 let to_dot_sep =
1345 |c: TableName| format!("{}.{}.{}", c.catalog_name, c.schema_name, c.table_name);
1346 assert_eq!(1, expr.source_table_names.len());
1347 assert_eq!(
1348 "greptime.public.metric_cte",
1349 to_dot_sep(expr.source_table_names[0].clone())
1350 );
1351 }
1352
1353 #[test]
1354 fn test_create_flow_tql_cte_source_tables_same_name() {
1355 let sql = r#"
1356CREATE FLOW calc_cte
1357SINK TO metric_cte_sink
1358EVAL INTERVAL '1m'
1359AS
1360WITH tql(ts, the_value) AS (
1361 TQL EVAL (now() - '1m'::interval, now(), '5s') tql
1362)
1363SELECT * FROM tql;
1364"#;
1365
1366 let stmt =
1367 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default())
1368 .unwrap()
1369 .pop()
1370 .unwrap();
1371
1372 let Statement::CreateFlow(create_flow) = stmt else {
1373 unreachable!()
1374 };
1375 let expr = to_create_flow_task_expr(create_flow, &QueryContext::arc()).unwrap();
1376
1377 let to_dot_sep =
1378 |c: TableName| format!("{}.{}.{}", c.catalog_name, c.schema_name, c.table_name);
1379 assert_eq!(1, expr.source_table_names.len());
1380 assert_eq!(
1381 "greptime.public.tql",
1382 to_dot_sep(expr.source_table_names[0].clone())
1383 );
1384 }
1385
1386 #[test]
1387 fn test_create_flow_expr() {
1388 let sql = r"
1389CREATE FLOW test_distinct_basic SINK TO out_distinct_basic AS
1390SELECT
1391 DISTINCT number as dis
1392FROM
1393 distinct_basic;";
1394 let stmt =
1395 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default())
1396 .unwrap()
1397 .pop()
1398 .unwrap();
1399
1400 let Statement::CreateFlow(create_flow) = stmt else {
1401 unreachable!()
1402 };
1403 let expr = to_create_flow_task_expr(create_flow, &QueryContext::arc()).unwrap();
1404
1405 let to_dot_sep =
1406 |c: TableName| format!("{}.{}.{}", c.catalog_name, c.schema_name, c.table_name);
1407 assert_eq!("test_distinct_basic", expr.flow_name);
1408 assert_eq!("greptime", expr.catalog_name);
1409 assert_eq!(
1410 "greptime.public.out_distinct_basic",
1411 expr.sink_table_name.map(to_dot_sep).unwrap()
1412 );
1413 assert_eq!(1, expr.source_table_names.len());
1414 assert_eq!(
1415 "greptime.public.distinct_basic",
1416 to_dot_sep(expr.source_table_names[0].clone())
1417 );
1418 assert_eq!(
1419 r"SELECT
1420 DISTINCT number as dis
1421FROM
1422 distinct_basic",
1423 expr.sql
1424 );
1425
1426 let sql = r"
1427CREATE FLOW `task_2`
1428SINK TO schema_1.table_1
1429AS
1430SELECT max(c1), min(c2) FROM schema_2.table_2;";
1431 let stmt =
1432 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default())
1433 .unwrap()
1434 .pop()
1435 .unwrap();
1436
1437 let Statement::CreateFlow(create_flow) = stmt else {
1438 unreachable!()
1439 };
1440 let expr = to_create_flow_task_expr(create_flow, &QueryContext::arc()).unwrap();
1441
1442 let to_dot_sep =
1443 |c: TableName| format!("{}.{}.{}", c.catalog_name, c.schema_name, c.table_name);
1444 assert_eq!("task_2", expr.flow_name);
1445 assert_eq!("greptime", expr.catalog_name);
1446 assert_eq!(
1447 "greptime.schema_1.table_1",
1448 expr.sink_table_name.map(to_dot_sep).unwrap()
1449 );
1450 assert_eq!(1, expr.source_table_names.len());
1451 assert_eq!(
1452 "greptime.schema_2.table_2",
1453 to_dot_sep(expr.source_table_names[0].clone())
1454 );
1455 assert_eq!("SELECT max(c1), min(c2) FROM schema_2.table_2", expr.sql);
1456 assert!(expr.flow_options.is_empty());
1457
1458 let sql = r"
1459CREATE FLOW task_3
1460SINK TO schema_1.table_1
1461WITH (defer_on_missing_source = 'true', foo = 'bar')
1462AS
1463SELECT max(c1), min(c2) FROM schema_2.table_2;";
1464 let stmt =
1465 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default())
1466 .unwrap()
1467 .pop()
1468 .unwrap();
1469
1470 let Statement::CreateFlow(create_flow) = stmt else {
1471 unreachable!()
1472 };
1473 let expr = to_create_flow_task_expr(create_flow, &QueryContext::arc()).unwrap();
1474 assert_eq!(
1475 expr.flow_options,
1476 HashMap::from([
1477 ("defer_on_missing_source".to_string(), "true".to_string()),
1478 ("foo".to_string(), "bar".to_string()),
1479 ])
1480 );
1481
1482 let sql = r"
1483CREATE FLOW task_4
1484SINK TO schema_1.table_1
1485WITH (defer_on_missing_source = true)
1486AS
1487SELECT max(c1), min(c2) FROM schema_2.table_2;";
1488 let stmt =
1489 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default())
1490 .unwrap()
1491 .pop()
1492 .unwrap();
1493
1494 let Statement::CreateFlow(create_flow) = stmt else {
1495 unreachable!()
1496 };
1497 let expr = to_create_flow_task_expr(create_flow, &QueryContext::arc()).unwrap();
1498 assert_eq!(
1499 expr.flow_options,
1500 HashMap::from([("defer_on_missing_source".to_string(), "true".to_string(),)])
1501 );
1502
1503 let sql = r"
1504CREATE FLOW task_5
1505SINK TO schema_1.table_1
1506WITH (defer_on_missing_source = [true])
1507AS
1508SELECT max(c1), min(c2) FROM schema_2.table_2;";
1509 let stmt =
1510 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default())
1511 .unwrap()
1512 .pop()
1513 .unwrap();
1514
1515 let Statement::CreateFlow(create_flow) = stmt else {
1516 unreachable!()
1517 };
1518 let res = to_create_flow_task_expr(create_flow, &QueryContext::arc());
1519 assert!(res.is_err());
1520 assert!(
1521 res.unwrap_err()
1522 .to_string()
1523 .contains("flow options only support scalar string-compatible values")
1524 );
1525
1526 let sql = r"
1527CREATE FLOW abc.`task_2`
1528SINK TO schema_1.table_1
1529AS
1530SELECT max(c1), min(c2) FROM schema_2.table_2;";
1531 let stmt =
1532 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default())
1533 .unwrap()
1534 .pop()
1535 .unwrap();
1536
1537 let Statement::CreateFlow(create_flow) = stmt else {
1538 unreachable!()
1539 };
1540 let res = to_create_flow_task_expr(create_flow, &QueryContext::arc());
1541
1542 assert!(res.is_err());
1543 assert!(
1544 res.unwrap_err()
1545 .to_string()
1546 .contains("Invalid flow name: abc.`task_2`")
1547 );
1548 }
1549
1550 #[test]
1551 fn test_create_to_expr() {
1552 let sql = "CREATE TABLE monitor (host STRING,ts TIMESTAMP,TIME INDEX (ts),PRIMARY KEY(host)) ENGINE=mito WITH(ttl='3days', write_buffer_size='1024KB');";
1553 let stmt =
1554 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default())
1555 .unwrap()
1556 .pop()
1557 .unwrap();
1558
1559 let Statement::CreateTable(create_table) = stmt else {
1560 unreachable!()
1561 };
1562 let expr = create_to_expr(&create_table, &QueryContext::arc()).unwrap();
1563 assert_eq!("3days", expr.table_options.get("ttl").unwrap());
1564 assert_eq!(
1565 "1.0MiB",
1566 expr.table_options.get("write_buffer_size").unwrap()
1567 );
1568
1569 let sql = "CREATE TABLE monitor (ts TIMESTAMP TIME INDEX) WITH(skip_wal='false');";
1570 let stmt =
1571 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default())
1572 .unwrap()
1573 .pop()
1574 .unwrap();
1575 let Statement::CreateTable(create_table) = stmt else {
1576 unreachable!()
1577 };
1578 let expr = create_to_expr(&create_table, &QueryContext::arc()).unwrap();
1579 assert_eq!(
1580 Some("false"),
1581 expr.table_options
1582 .get(store_api::mito_engine_options::SKIP_WAL_KEY)
1583 .map(String::as_str)
1584 );
1585 }
1586
1587 #[test]
1588 fn test_invalid_create_to_expr() {
1589 let cases = [
1590 "CREATE TABLE monitor (host STRING primary key, ts TIMESTAMP TIME INDEX, some_column text, some_column string);",
1592 "CREATE TABLE monitor (host STRING, ts TIMESTAMP TIME INDEX, some_column STRING, PRIMARY KEY (some_column, host, some_column));",
1594 "CREATE TABLE monitor (host STRING, ts TIMESTAMP TIME INDEX, PRIMARY KEY (host, ts));",
1596 ];
1597
1598 for sql in cases {
1599 let stmt = ParserContext::create_with_dialect(
1600 sql,
1601 &GreptimeDbDialect {},
1602 ParseOptions::default(),
1603 )
1604 .unwrap()
1605 .pop()
1606 .unwrap();
1607 let Statement::CreateTable(create_table) = stmt else {
1608 unreachable!()
1609 };
1610 create_to_expr(&create_table, &QueryContext::arc()).unwrap_err();
1611 }
1612 }
1613
1614 #[test]
1615 fn test_create_to_expr_with_default_timestamp_value() {
1616 let sql = "CREATE TABLE monitor (v double,ts TIMESTAMP default '2024-01-30T00:01:01',TIME INDEX (ts)) engine=mito;";
1617 let stmt =
1618 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default())
1619 .unwrap()
1620 .pop()
1621 .unwrap();
1622
1623 let Statement::CreateTable(create_table) = stmt else {
1624 unreachable!()
1625 };
1626
1627 let expr = create_to_expr(&create_table, &QueryContext::arc()).unwrap();
1629 let ts_column = &expr.column_defs[1];
1630 let constraint = assert_ts_column(ts_column);
1631 assert!(
1632 matches!(constraint, ColumnDefaultConstraint::Value(Value::Timestamp(ts))
1633 if ts.to_iso8601_string() == "2024-01-30 00:01:01+0000")
1634 );
1635
1636 let ctx = QueryContextBuilder::default()
1638 .timezone(Timezone::from_tz_string("+08:00").unwrap())
1639 .build()
1640 .into();
1641 let expr = create_to_expr(&create_table, &ctx).unwrap();
1642 let ts_column = &expr.column_defs[1];
1643 let constraint = assert_ts_column(ts_column);
1644 assert!(
1645 matches!(constraint, ColumnDefaultConstraint::Value(Value::Timestamp(ts))
1646 if ts.to_iso8601_string() == "2024-01-29 16:01:01+0000")
1647 );
1648 }
1649
1650 fn assert_ts_column(ts_column: &api::v1::ColumnDef) -> ColumnDefaultConstraint {
1651 assert_eq!("ts", ts_column.name);
1652 assert_eq!(
1653 ColumnDataType::TimestampMillisecond as i32,
1654 ts_column.data_type
1655 );
1656 assert!(!ts_column.default_constraint.is_empty());
1657
1658 ColumnDefaultConstraint::try_from(&ts_column.default_constraint[..]).unwrap()
1659 }
1660
1661 #[test]
1662 fn test_to_alter_expr() {
1663 let sql = "ALTER DATABASE greptime SET key1='value1', key2='value2';";
1664 let stmt =
1665 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default())
1666 .unwrap()
1667 .pop()
1668 .unwrap();
1669
1670 let Statement::AlterDatabase(alter_database) = stmt else {
1671 unreachable!()
1672 };
1673
1674 let expr = to_alter_database_expr(alter_database, &QueryContext::arc()).unwrap();
1675 let kind = expr.kind.unwrap();
1676
1677 let AlterDatabaseKind::SetDatabaseOptions(SetDatabaseOptions {
1678 set_database_options,
1679 }) = kind
1680 else {
1681 unreachable!()
1682 };
1683
1684 assert_eq!(2, set_database_options.len());
1685 assert_eq!("key1", set_database_options[0].key);
1686 assert_eq!("value1", set_database_options[0].value);
1687 assert_eq!("key2", set_database_options[1].key);
1688 assert_eq!("value2", set_database_options[1].value);
1689
1690 let sql = "ALTER DATABASE greptime UNSET key1, key2;";
1691 let stmt =
1692 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default())
1693 .unwrap()
1694 .pop()
1695 .unwrap();
1696
1697 let Statement::AlterDatabase(alter_database) = stmt else {
1698 unreachable!()
1699 };
1700
1701 let expr = to_alter_database_expr(alter_database, &QueryContext::arc()).unwrap();
1702 let kind = expr.kind.unwrap();
1703
1704 let AlterDatabaseKind::UnsetDatabaseOptions(UnsetDatabaseOptions { keys }) = kind else {
1705 unreachable!()
1706 };
1707
1708 assert_eq!(2, keys.len());
1709 assert!(keys.contains(&"key1".to_string()));
1710 assert!(keys.contains(&"key2".to_string()));
1711
1712 let sql = "ALTER TABLE monitor add column ts TIMESTAMP default '2024-01-30T00:01:01';";
1713 let stmt =
1714 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default())
1715 .unwrap()
1716 .pop()
1717 .unwrap();
1718
1719 let Statement::AlterTable(alter_table) = stmt else {
1720 unreachable!()
1721 };
1722
1723 let expr = to_alter_table_expr(alter_table.clone(), &QueryContext::arc()).unwrap();
1725 let kind = expr.kind.unwrap();
1726
1727 let AlterTableKind::AddColumns(AddColumns { add_columns, .. }) = kind else {
1728 unreachable!()
1729 };
1730
1731 assert_eq!(1, add_columns.len());
1732 let ts_column = add_columns[0].column_def.clone().unwrap();
1733 let constraint = assert_ts_column(&ts_column);
1734 assert!(
1735 matches!(constraint, ColumnDefaultConstraint::Value(Value::Timestamp(ts))
1736 if ts.to_iso8601_string() == "2024-01-30 00:01:01+0000")
1737 );
1738
1739 let ctx = QueryContextBuilder::default()
1742 .timezone(Timezone::from_tz_string("+08:00").unwrap())
1743 .build()
1744 .into();
1745 let expr = to_alter_table_expr(alter_table, &ctx).unwrap();
1746 let kind = expr.kind.unwrap();
1747
1748 let AlterTableKind::AddColumns(AddColumns { add_columns, .. }) = kind else {
1749 unreachable!()
1750 };
1751
1752 assert_eq!(1, add_columns.len());
1753 let ts_column = add_columns[0].column_def.clone().unwrap();
1754 let constraint = assert_ts_column(&ts_column);
1755 assert!(
1756 matches!(constraint, ColumnDefaultConstraint::Value(Value::Timestamp(ts))
1757 if ts.to_iso8601_string() == "2024-01-29 16:01:01+0000")
1758 );
1759 }
1760
1761 #[test]
1762 fn test_to_alter_modify_column_type_expr() {
1763 let sql = "ALTER TABLE monitor MODIFY COLUMN mem_usage STRING;";
1764 let stmt =
1765 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default())
1766 .unwrap()
1767 .pop()
1768 .unwrap();
1769
1770 let Statement::AlterTable(alter_table) = stmt else {
1771 unreachable!()
1772 };
1773
1774 let expr = to_alter_table_expr(alter_table.clone(), &QueryContext::arc()).unwrap();
1776 let kind = expr.kind.unwrap();
1777
1778 let AlterTableKind::ModifyColumnTypes(ModifyColumnTypes {
1779 modify_column_types,
1780 }) = kind
1781 else {
1782 unreachable!()
1783 };
1784
1785 assert_eq!(1, modify_column_types.len());
1786 let modify_column_type = &modify_column_types[0];
1787
1788 assert_eq!("mem_usage", modify_column_type.column_name);
1789 assert_eq!(
1790 ColumnDataType::String as i32,
1791 modify_column_type.target_type
1792 );
1793 assert!(modify_column_type.target_type_extension.is_none());
1794 }
1795
1796 #[test]
1797 fn test_to_alter_set_json_settings_expr() {
1798 let sql = "ALTER TABLE monitor MODIFY COLUMN payload JSON2 (service STRING);";
1799 let stmt =
1800 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default())
1801 .unwrap()
1802 .pop()
1803 .unwrap();
1804
1805 let Statement::AlterTable(alter_table) = stmt else {
1806 unreachable!()
1807 };
1808 let expr = to_alter_table_expr(alter_table, &QueryContext::arc()).unwrap();
1809 let kind = expr.kind.unwrap();
1810
1811 let AlterTableKind::SetJsonSettings(modify) = kind else {
1812 unreachable!()
1813 };
1814 assert_eq!("payload", modify.column_name);
1815 let settings = modify.settings.as_ref().unwrap();
1816 assert_eq!(Some(100), settings.max_auto_expanded_paths);
1817 assert_eq!(1, settings.type_hints.len());
1818 assert_eq!(["service"], &settings.type_hints[0].path[..]);
1819
1820 let sql = "ALTER TABLE monitor MODIFY COLUMN payload JSON2;";
1821 let stmt =
1822 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default())
1823 .unwrap()
1824 .pop()
1825 .unwrap();
1826
1827 let Statement::AlterTable(alter_table) = stmt else {
1828 unreachable!()
1829 };
1830 let expr = to_alter_table_expr(alter_table, &QueryContext::arc()).unwrap();
1831 let kind = expr.kind.unwrap();
1832
1833 let AlterTableKind::SetJsonSettings(modify) = kind else {
1834 unreachable!()
1835 };
1836 let settings = modify.settings.as_ref().unwrap();
1837 assert_eq!(Some(100), settings.max_auto_expanded_paths);
1838 assert!(settings.type_hints.is_empty());
1839 }
1840
1841 #[test]
1842 fn test_to_repartition_request() {
1843 let sql = r#"
1844ALTER TABLE metrics REPARTITION (
1845 device_id < 100
1846) INTO (
1847 device_id < 100 AND area < 'South',
1848 device_id < 100 AND area >= 'South'
1849);"#;
1850 let stmt =
1851 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default())
1852 .unwrap()
1853 .pop()
1854 .unwrap();
1855
1856 let Statement::AlterTable(alter_table) = stmt else {
1857 unreachable!()
1858 };
1859
1860 let request = to_repartition_request(alter_table, &QueryContext::arc()).unwrap();
1861 assert_eq!("greptime", request.catalog_name);
1862 assert_eq!("public", request.schema_name);
1863 assert_eq!("metrics", request.table_name);
1864 let RepartitionSource::Partitions {
1865 from_exprs,
1866 target_partition_columns,
1867 } = request.source
1868 else {
1869 unreachable!()
1870 };
1871 assert!(target_partition_columns.is_none());
1872 assert_eq!(
1873 from_exprs
1874 .into_iter()
1875 .map(|x| x.to_string())
1876 .collect::<Vec<_>>(),
1877 vec!["device_id < 100".to_string()]
1878 );
1879 assert_eq!(
1880 request
1881 .into_exprs
1882 .into_iter()
1883 .map(|x| x.to_string())
1884 .collect::<Vec<_>>(),
1885 vec![
1886 "device_id < 100 AND area < 'South'".to_string(),
1887 "device_id < 100 AND area >= 'South'".to_string()
1888 ]
1889 );
1890 }
1891
1892 #[test]
1893 fn test_to_repartition_request_with_target_partition_columns() {
1894 let sql = r#"
1895ALTER TABLE metrics REPARTITION (
1896 device_id < 100
1897) ON COLUMNS (device_id, area) INTO (
1898 device_id < 100 AND area < 'South',
1899 device_id < 100 AND area >= 'South'
1900);"#;
1901 let stmt =
1902 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default())
1903 .unwrap()
1904 .pop()
1905 .unwrap();
1906
1907 let Statement::AlterTable(alter_table) = stmt else {
1908 unreachable!()
1909 };
1910
1911 let request = to_repartition_request(alter_table, &QueryContext::arc()).unwrap();
1912 let RepartitionSource::Partitions {
1913 target_partition_columns,
1914 ..
1915 } = request.source
1916 else {
1917 unreachable!()
1918 };
1919
1920 assert_eq!(
1921 target_partition_columns,
1922 Some(vec!["device_id".to_string(), "area".to_string()])
1923 );
1924 }
1925
1926 #[test]
1927 fn test_to_repartition_request_with_unpartitioned_source() {
1928 let sql = r#"
1929ALTER TABLE metrics PARTITION ON COLUMNS (device_id, area) (
1930 device_id < 100 AND area < 'South',
1931 device_id < 100 AND area >= 'South'
1932);"#;
1933 let stmt =
1934 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default())
1935 .unwrap()
1936 .pop()
1937 .unwrap();
1938
1939 let Statement::AlterTable(alter_table) = stmt else {
1940 unreachable!()
1941 };
1942
1943 let request = to_repartition_request(alter_table, &QueryContext::arc()).unwrap();
1944 assert_eq!("greptime", request.catalog_name);
1945 assert_eq!("public", request.schema_name);
1946 assert_eq!("metrics", request.table_name);
1947 let RepartitionSource::Unpartitioned { partition_columns } = request.source else {
1948 unreachable!()
1949 };
1950 assert_eq!(partition_columns, vec!["device_id", "area"]);
1951 assert_eq!(
1952 request
1953 .into_exprs
1954 .into_iter()
1955 .map(|x| x.to_string())
1956 .collect::<Vec<_>>(),
1957 vec![
1958 "device_id < 100 AND area < 'South'".to_string(),
1959 "device_id < 100 AND area >= 'South'".to_string()
1960 ]
1961 );
1962 }
1963
1964 fn new_test_table_names() -> Vec<TableName> {
1965 vec![
1966 TableName {
1967 catalog_name: "greptime".to_string(),
1968 schema_name: "public".to_string(),
1969 table_name: "a_table".to_string(),
1970 },
1971 TableName {
1972 catalog_name: "greptime".to_string(),
1973 schema_name: "public".to_string(),
1974 table_name: "b_table".to_string(),
1975 },
1976 ]
1977 }
1978
1979 #[test]
1980 fn test_to_create_view_expr() {
1981 let sql = "CREATE VIEW test AS SELECT * FROM NUMBERS";
1982 let stmt =
1983 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default())
1984 .unwrap()
1985 .pop()
1986 .unwrap();
1987
1988 let Statement::CreateView(stmt) = stmt else {
1989 unreachable!()
1990 };
1991
1992 let logical_plan = vec![1, 2, 3];
1993 let table_names = new_test_table_names();
1994 let columns = vec!["a".to_string()];
1995 let plan_columns = vec!["number".to_string()];
1996
1997 let expr = to_create_view_expr(
1998 stmt,
1999 logical_plan.clone(),
2000 table_names.clone(),
2001 columns.clone(),
2002 plan_columns.clone(),
2003 sql.to_string(),
2004 QueryContext::arc(),
2005 )
2006 .unwrap();
2007
2008 assert_eq!("greptime", expr.catalog_name);
2009 assert_eq!("public", expr.schema_name);
2010 assert_eq!("test", expr.view_name);
2011 assert!(!expr.create_if_not_exists);
2012 assert!(!expr.or_replace);
2013 assert_eq!(logical_plan, expr.logical_plan);
2014 assert_eq!(table_names, expr.table_names);
2015 assert_eq!(sql, expr.definition);
2016 assert_eq!(columns, expr.columns);
2017 assert_eq!(plan_columns, expr.plan_columns);
2018 }
2019
2020 #[test]
2021 fn test_to_create_view_expr_complex() {
2022 let sql = "CREATE OR REPLACE VIEW IF NOT EXISTS test.test_view AS SELECT * FROM NUMBERS";
2023 let stmt =
2024 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default())
2025 .unwrap()
2026 .pop()
2027 .unwrap();
2028
2029 let Statement::CreateView(stmt) = stmt else {
2030 unreachable!()
2031 };
2032
2033 let logical_plan = vec![1, 2, 3];
2034 let table_names = new_test_table_names();
2035 let columns = vec!["a".to_string()];
2036 let plan_columns = vec!["number".to_string()];
2037
2038 let expr = to_create_view_expr(
2039 stmt,
2040 logical_plan.clone(),
2041 table_names.clone(),
2042 columns.clone(),
2043 plan_columns.clone(),
2044 sql.to_string(),
2045 QueryContext::arc(),
2046 )
2047 .unwrap();
2048
2049 assert_eq!("greptime", expr.catalog_name);
2050 assert_eq!("test", expr.schema_name);
2051 assert_eq!("test_view", expr.view_name);
2052 assert!(expr.create_if_not_exists);
2053 assert!(expr.or_replace);
2054 assert_eq!(logical_plan, expr.logical_plan);
2055 assert_eq!(table_names, expr.table_names);
2056 assert_eq!(sql, expr.definition);
2057 assert_eq!(columns, expr.columns);
2058 assert_eq!(plan_columns, expr.plan_columns);
2059 }
2060
2061 #[test]
2062 fn test_expr_to_create() {
2063 let sql = r#"CREATE TABLE IF NOT EXISTS `tt` (
2064 `timestamp` TIMESTAMP(9) NOT NULL,
2065 `ip_address` STRING NULL SKIPPING INDEX WITH(false_positive_rate = '0.01', granularity = '10240', type = 'BLOOM'),
2066 `username` STRING NULL,
2067 `http_method` STRING NULL INVERTED INDEX,
2068 `request_line` STRING NULL FULLTEXT INDEX WITH(analyzer = 'English', backend = 'bloom', case_sensitive = 'false', false_positive_rate = '0.01', granularity = '10240'),
2069 `protocol` STRING NULL,
2070 `status_code` INT NULL INVERTED INDEX,
2071 `response_size` BIGINT NULL,
2072 `message` STRING NULL,
2073 TIME INDEX (`timestamp`),
2074 PRIMARY KEY (`username`, `status_code`)
2075)
2076ENGINE=mito
2077WITH(
2078 append_mode = 'true'
2079)"#;
2080 let stmt =
2081 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default())
2082 .unwrap()
2083 .pop()
2084 .unwrap();
2085
2086 let Statement::CreateTable(original_create) = stmt else {
2087 unreachable!()
2088 };
2089
2090 let expr = create_to_expr(&original_create, &QueryContext::arc()).unwrap();
2092
2093 let create_table = expr_to_create(&expr, Some('`')).unwrap();
2094 let new_sql = format!("{:#}", create_table);
2095 assert_eq!(sql, new_sql);
2096 }
2097}