1mod show_create_table;
16
17use std::collections::HashMap;
18use std::ops::ControlFlow;
19use std::sync::Arc;
20
21use catalog::CatalogManagerRef;
22use catalog::information_schema::{
23 CHARACTER_SETS, COLLATIONS, COLUMNS, FLOW_STATISTICS, FLOWS, REGION_PEERS, SCHEMATA,
24 STATISTICS, TABLES, VIEWS, columns, flow_statistics, flows, process_list, region_peers,
25 schemata, statistics, tables,
26};
27use common_catalog::consts::{
28 INFORMATION_SCHEMA_NAME, SEMANTIC_TYPE_FIELD, SEMANTIC_TYPE_PRIMARY_KEY,
29 SEMANTIC_TYPE_TIME_INDEX,
30};
31use common_catalog::format_full_table_name;
32use common_datasource::file_format::{FileFormat, Format, infer_schemas};
33use common_datasource::lister::{Lister, Source};
34use common_datasource::object_store::{LocalFileAccess, build_backend_with_path};
35use common_error::ext::BoxedError;
36use common_meta::SchemaOptions;
37use common_meta::ddl::create_flow::{FlowType, effective_eval_schedule_from_flow_info};
38use common_meta::key::flow::flow_info::FlowInfoValue;
39use common_query::Output;
40use common_query::prelude::greptime_timestamp;
41use common_recordbatch::RecordBatches;
42use common_recordbatch::adapter::RecordBatchStreamAdapter;
43use common_time::Timestamp;
44use common_time::timezone::get_timezone;
45use datafusion::dataframe::DataFrame;
46use datafusion::prelude::SessionContext;
47use datafusion_expr::{Expr, SortExpr, col, lit};
48use datatypes::prelude::*;
49use datatypes::schema::{ColumnDefaultConstraint, ColumnSchema, Schema};
50use datatypes::vectors::StringVector;
51use itertools::Itertools;
52use object_store::ObjectStore;
53use once_cell::sync::Lazy;
54use regex::Regex;
55use session::context::{Channel, QueryContextRef};
56pub use show_create_table::create_table_stmt;
57use snafu::{OptionExt, ResultExt, ensure};
58use sql::ast::{Ident, visit_expressions_mut};
59use sql::parser::ParserContext;
60use sql::statements::OptionMap;
61use sql::statements::create::{CreateDatabase, CreateFlow, CreateView, Partitions, SqlOrTql};
62use sql::statements::show::{
63 ShowColumns, ShowDatabases, ShowFlowStatus, ShowFlows, ShowIndex, ShowKind, ShowProcessList,
64 ShowRegion, ShowTableStatus, ShowTables, ShowVariables, ShowViews,
65};
66use sql::statements::statement::Statement;
67use sqlparser::ast::ObjectName;
68use store_api::metric_engine_consts::{is_metric_engine, is_metric_engine_internal_column};
69use table::TableRef;
70use table::metadata::TableInfoRef;
71use table::requests::{FILE_TABLE_LOCATION_KEY, FILE_TABLE_PATTERN_KEY};
72
73use crate::QueryEngineRef;
74use crate::error::{self, Result, UnsupportedVariableSnafu};
75use crate::planner::DfLogicalPlanner;
76
77const SCHEMAS_COLUMN: &str = "Database";
78const OPTIONS_COLUMN: &str = "Options";
79const VIEWS_COLUMN: &str = "Views";
80const FLOWS_COLUMN: &str = "Flows";
81const FIELD_COLUMN: &str = "Field";
82const TABLE_TYPE_COLUMN: &str = "Table_type";
83
84const COLUMN_NAME_COLUMN: &str = "Column";
85const COLUMN_GREPTIME_TYPE_COLUMN: &str = "Greptime_type";
86const COLUMN_TYPE_COLUMN: &str = "Type";
87const COLUMN_KEY_COLUMN: &str = "Key";
88const COLUMN_EXTRA_COLUMN: &str = "Extra";
89const COLUMN_PRIVILEGES_COLUMN: &str = "Privileges";
90const COLUMN_COLLATION_COLUMN: &str = "Collation";
91const COLUMN_NULLABLE_COLUMN: &str = "Null";
92const COLUMN_DEFAULT_COLUMN: &str = "Default";
93const COLUMN_COMMENT_COLUMN: &str = "Comment";
94const COLUMN_SEMANTIC_TYPE_COLUMN: &str = "Semantic Type";
95
96const YES_STR: &str = "YES";
97const NO_STR: &str = "NO";
98const PRI_KEY: &str = "PRI";
99
100const INDEX_TABLE_COLUMN: &str = "Table";
102const INDEX_NONT_UNIQUE_COLUMN: &str = "Non_unique";
103const INDEX_CARDINALITY_COLUMN: &str = "Cardinality";
104const INDEX_SUB_PART_COLUMN: &str = "Sub_part";
105const INDEX_PACKED_COLUMN: &str = "Packed";
106const INDEX_INDEX_TYPE_COLUMN: &str = "Index_type";
107const INDEX_COMMENT_COLUMN: &str = "Index_comment";
108const INDEX_VISIBLE_COLUMN: &str = "Visible";
109const INDEX_EXPRESSION_COLUMN: &str = "Expression";
110const INDEX_KEY_NAME_COLUMN: &str = "Key_name";
111const INDEX_SEQ_IN_INDEX_COLUMN: &str = "Seq_in_index";
112const INDEX_COLUMN_NAME_COLUMN: &str = "Column_name";
113
114pub static DESCRIBE_TABLE_OUTPUT_SCHEMA: Lazy<Arc<Schema>> = Lazy::new(|| {
115 Arc::new(Schema::new(vec![
116 ColumnSchema::new(
117 COLUMN_NAME_COLUMN,
118 ConcreteDataType::string_datatype(),
119 false,
120 ),
121 ColumnSchema::new(
122 COLUMN_TYPE_COLUMN,
123 ConcreteDataType::string_datatype(),
124 false,
125 ),
126 ColumnSchema::new(COLUMN_KEY_COLUMN, ConcreteDataType::string_datatype(), true),
127 ColumnSchema::new(
128 COLUMN_NULLABLE_COLUMN,
129 ConcreteDataType::string_datatype(),
130 false,
131 ),
132 ColumnSchema::new(
133 COLUMN_DEFAULT_COLUMN,
134 ConcreteDataType::string_datatype(),
135 false,
136 ),
137 ColumnSchema::new(
138 COLUMN_SEMANTIC_TYPE_COLUMN,
139 ConcreteDataType::string_datatype(),
140 false,
141 ),
142 ]))
143});
144
145static SHOW_CREATE_DATABASE_OUTPUT_SCHEMA: Lazy<Arc<Schema>> = Lazy::new(|| {
146 Arc::new(Schema::new(vec![
147 ColumnSchema::new("Database", ConcreteDataType::string_datatype(), false),
148 ColumnSchema::new(
149 "Create Database",
150 ConcreteDataType::string_datatype(),
151 false,
152 ),
153 ]))
154});
155
156static SHOW_CREATE_TABLE_OUTPUT_SCHEMA: Lazy<Arc<Schema>> = Lazy::new(|| {
157 Arc::new(Schema::new(vec![
158 ColumnSchema::new("Table", ConcreteDataType::string_datatype(), false),
159 ColumnSchema::new("Create Table", ConcreteDataType::string_datatype(), false),
160 ]))
161});
162
163static SHOW_CREATE_FLOW_OUTPUT_SCHEMA: Lazy<Arc<Schema>> = Lazy::new(|| {
164 Arc::new(Schema::new(vec![
165 ColumnSchema::new("Flow", ConcreteDataType::string_datatype(), false),
166 ColumnSchema::new("Create Flow", ConcreteDataType::string_datatype(), false),
167 ]))
168});
169
170static SHOW_CREATE_VIEW_OUTPUT_SCHEMA: Lazy<Arc<Schema>> = Lazy::new(|| {
171 Arc::new(Schema::new(vec![
172 ColumnSchema::new("View", ConcreteDataType::string_datatype(), false),
173 ColumnSchema::new("Create View", ConcreteDataType::string_datatype(), false),
174 ]))
175});
176
177pub async fn show_databases(
178 stmt: ShowDatabases,
179 query_engine: &QueryEngineRef,
180 catalog_manager: &CatalogManagerRef,
181 query_ctx: QueryContextRef,
182) -> Result<Output> {
183 let dataframe =
184 show_databases_dataframe(&stmt, query_engine, catalog_manager, query_ctx).await?;
185 dataframe_to_output(dataframe).await
186}
187
188pub async fn show_databases_dataframe(
190 stmt: &ShowDatabases,
191 query_engine: &QueryEngineRef,
192 catalog_manager: &CatalogManagerRef,
193 query_ctx: QueryContextRef,
194) -> Result<DataFrame> {
195 let projects = if stmt.full {
196 vec![
197 (schemata::SCHEMA_NAME, SCHEMAS_COLUMN),
198 (schemata::SCHEMA_OPTS, OPTIONS_COLUMN),
199 ]
200 } else {
201 vec![(schemata::SCHEMA_NAME, SCHEMAS_COLUMN)]
202 };
203
204 let filters = vec![col(schemata::CATALOG_NAME).eq(lit(query_ctx.current_catalog()))];
205 let like_field = Some(schemata::SCHEMA_NAME);
206 let sort = vec![col(schemata::SCHEMA_NAME).sort(true, true)];
207
208 query_from_information_schema_dataframe(
209 query_engine,
210 catalog_manager,
211 query_ctx,
212 SCHEMATA,
213 vec![],
214 projects,
215 filters,
216 like_field,
217 sort,
218 &stmt.kind,
219 )
220 .await
221}
222
223fn replace_column_in_expr(expr: &mut sqlparser::ast::Expr, from_column: &str, to_column: &str) {
226 let _ = visit_expressions_mut(expr, |e| {
227 match e {
228 sqlparser::ast::Expr::Identifier(ident)
229 if ident.value.eq_ignore_ascii_case(from_column) =>
230 {
231 ident.value = to_column.to_string();
232 }
233 sqlparser::ast::Expr::CompoundIdentifier(idents) => {
234 if let Some(last) = idents.last_mut()
235 && last.value.eq_ignore_ascii_case(from_column)
236 {
237 last.value = to_column.to_string();
238 }
239 }
240 _ => {}
241 }
242 ControlFlow::<()>::Continue(())
243 });
244}
245
246#[allow(clippy::too_many_arguments)]
254async fn query_from_information_schema_table(
255 query_engine: &QueryEngineRef,
256 catalog_manager: &CatalogManagerRef,
257 query_ctx: QueryContextRef,
258 table_name: &str,
259 select: Vec<Expr>,
260 projects: Vec<(&str, &str)>,
261 filters: Vec<Expr>,
262 like_field: Option<&str>,
263 sort: Vec<SortExpr>,
264 kind: ShowKind,
265) -> Result<Output> {
266 let dataframe = query_from_information_schema_dataframe(
267 query_engine,
268 catalog_manager,
269 query_ctx,
270 table_name,
271 select,
272 projects,
273 filters,
274 like_field,
275 sort,
276 &kind,
277 )
278 .await?;
279 dataframe_to_output(dataframe).await
280}
281
282async fn dataframe_to_output(dataframe: DataFrame) -> Result<Output> {
283 let stream = dataframe.execute_stream().await?;
284 Ok(Output::new_with_stream(Box::pin(
285 RecordBatchStreamAdapter::try_new(stream).context(error::CreateRecordBatchSnafu)?,
286 )))
287}
288
289#[allow(clippy::too_many_arguments)]
293async fn query_from_information_schema_dataframe(
294 query_engine: &QueryEngineRef,
295 catalog_manager: &CatalogManagerRef,
296 query_ctx: QueryContextRef,
297 table_name: &str,
298 select: Vec<Expr>,
299 projects: Vec<(&str, &str)>,
300 filters: Vec<Expr>,
301 like_field: Option<&str>,
302 sort: Vec<SortExpr>,
303 kind: &ShowKind,
304) -> Result<DataFrame> {
305 let table = catalog_manager
306 .table(
307 query_ctx.current_catalog(),
308 INFORMATION_SCHEMA_NAME,
309 table_name,
310 Some(&query_ctx),
311 )
312 .await
313 .context(error::CatalogSnafu)?
314 .with_context(|| error::TableNotFoundSnafu {
315 table: format_full_table_name(
316 query_ctx.current_catalog(),
317 INFORMATION_SCHEMA_NAME,
318 table_name,
319 ),
320 })?;
321
322 let dataframe = query_engine.read_table(table)?;
323
324 let dataframe = filters.into_iter().try_fold(dataframe, |df, expr| {
326 df.filter(expr).context(error::PlanSqlSnafu)
327 })?;
328
329 let dataframe = if let (ShowKind::Like(ident), Some(field)) = (&kind, like_field) {
331 dataframe
332 .filter(col(field).like(lit(ident.value.clone())))
333 .context(error::PlanSqlSnafu)?
334 } else {
335 dataframe
336 };
337
338 let dataframe = if sort.is_empty() {
340 dataframe
341 } else {
342 dataframe.sort(sort).context(error::PlanSqlSnafu)?
343 };
344
345 let dataframe = if select.is_empty() {
347 if projects.is_empty() {
348 dataframe
349 } else {
350 let projection = projects
351 .iter()
352 .map(|x| col(x.0).alias(x.1))
353 .collect::<Vec<_>>();
354 dataframe.select(projection).context(error::PlanSqlSnafu)?
355 }
356 } else {
357 dataframe.select(select).context(error::PlanSqlSnafu)?
358 };
359
360 let dataframe = projects
362 .into_iter()
363 .try_fold(dataframe, |df, (column, renamed_column)| {
364 df.with_column_renamed(column, renamed_column)
365 .context(error::PlanSqlSnafu)
366 })?;
367
368 let dataframe = match kind {
369 ShowKind::All | ShowKind::Like(_) => {
370 dataframe
372 }
373 ShowKind::Where(filter) => {
374 let view = dataframe.into_view();
377 let dataframe = SessionContext::new_with_state(
378 query_engine
379 .engine_context(query_ctx.clone())
380 .state()
381 .clone(),
382 )
383 .read_table(view)?;
384
385 let planner = query_engine.planner();
386 let planner = planner
387 .as_any()
388 .downcast_ref::<DfLogicalPlanner>()
389 .expect("Must be the datafusion planner");
390
391 let filter = planner
392 .sql_to_expr(filter.clone(), dataframe.schema(), false, query_ctx)
393 .await?;
394
395 dataframe.filter(filter).context(error::PlanSqlSnafu)?
397 }
398 };
399
400 Ok(dataframe)
401}
402
403pub async fn show_columns(
405 stmt: ShowColumns,
406 query_engine: &QueryEngineRef,
407 catalog_manager: &CatalogManagerRef,
408 query_ctx: QueryContextRef,
409) -> Result<Output> {
410 let dataframe = show_columns_dataframe(&stmt, query_engine, catalog_manager, query_ctx).await?;
411 dataframe_to_output(dataframe).await
412}
413
414pub async fn show_columns_dataframe(
416 stmt: &ShowColumns,
417 query_engine: &QueryEngineRef,
418 catalog_manager: &CatalogManagerRef,
419 query_ctx: QueryContextRef,
420) -> Result<DataFrame> {
421 let schema_name = if let Some(database) = &stmt.database {
422 database.clone()
423 } else {
424 query_ctx.current_schema()
425 };
426
427 let projects = if stmt.full {
428 vec![
429 (columns::COLUMN_NAME, FIELD_COLUMN),
430 (columns::DATA_TYPE, COLUMN_TYPE_COLUMN),
431 (columns::COLLATION_NAME, COLUMN_COLLATION_COLUMN),
432 (columns::IS_NULLABLE, COLUMN_NULLABLE_COLUMN),
433 (columns::COLUMN_KEY, COLUMN_KEY_COLUMN),
434 (columns::COLUMN_DEFAULT, COLUMN_DEFAULT_COLUMN),
435 (columns::COLUMN_COMMENT, COLUMN_COMMENT_COLUMN),
436 (columns::PRIVILEGES, COLUMN_PRIVILEGES_COLUMN),
437 (columns::EXTRA, COLUMN_EXTRA_COLUMN),
438 (columns::GREPTIME_DATA_TYPE, COLUMN_GREPTIME_TYPE_COLUMN),
439 ]
440 } else {
441 vec![
442 (columns::COLUMN_NAME, FIELD_COLUMN),
443 (columns::DATA_TYPE, COLUMN_TYPE_COLUMN),
444 (columns::IS_NULLABLE, COLUMN_NULLABLE_COLUMN),
445 (columns::COLUMN_KEY, COLUMN_KEY_COLUMN),
446 (columns::COLUMN_DEFAULT, COLUMN_DEFAULT_COLUMN),
447 (columns::EXTRA, COLUMN_EXTRA_COLUMN),
448 (columns::GREPTIME_DATA_TYPE, COLUMN_GREPTIME_TYPE_COLUMN),
449 ]
450 };
451
452 let filters = vec![
453 col(columns::TABLE_NAME).eq(lit(&stmt.table)),
454 col(columns::TABLE_SCHEMA).eq(lit(schema_name.clone())),
455 col(columns::TABLE_CATALOG).eq(lit(query_ctx.current_catalog())),
456 ];
457 let like_field = Some(columns::COLUMN_NAME);
458 let sort = vec![col(columns::COLUMN_NAME).sort(true, true)];
459
460 query_from_information_schema_dataframe(
461 query_engine,
462 catalog_manager,
463 query_ctx,
464 COLUMNS,
465 vec![],
466 projects,
467 filters,
468 like_field,
469 sort,
470 &stmt.kind,
471 )
472 .await
473}
474
475pub async fn show_index(
477 stmt: ShowIndex,
478 query_engine: &QueryEngineRef,
479 catalog_manager: &CatalogManagerRef,
480 query_ctx: QueryContextRef,
481) -> Result<Output> {
482 let dataframe = show_index_dataframe(&stmt, query_engine, catalog_manager, query_ctx).await?;
483 dataframe_to_output(dataframe).await
484}
485
486pub async fn show_index_dataframe(
488 stmt: &ShowIndex,
489 query_engine: &QueryEngineRef,
490 catalog_manager: &CatalogManagerRef,
491 query_ctx: QueryContextRef,
492) -> Result<DataFrame> {
493 let schema_name = if let Some(database) = &stmt.database {
494 database.clone()
495 } else {
496 query_ctx.current_schema()
497 };
498
499 let select = vec![
500 col(statistics::TABLE_NAME).alias(INDEX_TABLE_COLUMN),
501 col(statistics::NON_UNIQUE).alias(INDEX_NONT_UNIQUE_COLUMN),
502 col(statistics::INDEX_NAME).alias(INDEX_KEY_NAME_COLUMN),
503 col(statistics::SEQ_IN_INDEX).alias(INDEX_SEQ_IN_INDEX_COLUMN),
504 col(statistics::COLUMN_NAME).alias(INDEX_COLUMN_NAME_COLUMN),
505 col(statistics::COLLATION).alias(COLUMN_COLLATION_COLUMN),
506 col(statistics::CARDINALITY).alias(INDEX_CARDINALITY_COLUMN),
507 col(statistics::SUB_PART).alias(INDEX_SUB_PART_COLUMN),
508 col(statistics::PACKED).alias(INDEX_PACKED_COLUMN),
509 col(statistics::NULLABLE).alias(COLUMN_NULLABLE_COLUMN),
510 col(statistics::INDEX_TYPE).alias(INDEX_INDEX_TYPE_COLUMN),
511 col(statistics::COMMENT).alias(COLUMN_COMMENT_COLUMN),
512 col(statistics::INDEX_COMMENT).alias(INDEX_COMMENT_COLUMN),
513 col(statistics::IS_VISIBLE).alias(INDEX_VISIBLE_COLUMN),
514 col(statistics::EXPRESSION).alias(INDEX_EXPRESSION_COLUMN),
515 ];
516
517 let projects = vec![
518 (statistics::TABLE_NAME, INDEX_TABLE_COLUMN),
519 (INDEX_NONT_UNIQUE_COLUMN, INDEX_NONT_UNIQUE_COLUMN),
520 (statistics::INDEX_NAME, INDEX_KEY_NAME_COLUMN),
521 (statistics::SEQ_IN_INDEX, INDEX_SEQ_IN_INDEX_COLUMN),
522 (statistics::COLUMN_NAME, INDEX_COLUMN_NAME_COLUMN),
523 (COLUMN_COLLATION_COLUMN, COLUMN_COLLATION_COLUMN),
524 (INDEX_CARDINALITY_COLUMN, INDEX_CARDINALITY_COLUMN),
525 (INDEX_SUB_PART_COLUMN, INDEX_SUB_PART_COLUMN),
526 (INDEX_PACKED_COLUMN, INDEX_PACKED_COLUMN),
527 (COLUMN_NULLABLE_COLUMN, COLUMN_NULLABLE_COLUMN),
528 (statistics::INDEX_TYPE, INDEX_INDEX_TYPE_COLUMN),
529 (COLUMN_COMMENT_COLUMN, COLUMN_COMMENT_COLUMN),
530 (INDEX_COMMENT_COLUMN, INDEX_COMMENT_COLUMN),
531 (INDEX_VISIBLE_COLUMN, INDEX_VISIBLE_COLUMN),
532 (INDEX_EXPRESSION_COLUMN, INDEX_EXPRESSION_COLUMN),
533 ];
534
535 let filters = vec![
536 col(statistics::TABLE_NAME).eq(lit(&stmt.table)),
537 col(statistics::TABLE_SCHEMA).eq(lit(schema_name.clone())),
538 ];
539 let like_field = None;
540 let sort = vec![
541 col(statistics::INDEX_NAME).sort(true, true),
542 col(statistics::SEQ_IN_INDEX).sort(true, true),
543 ];
544
545 query_from_information_schema_dataframe(
546 query_engine,
547 catalog_manager,
548 query_ctx,
549 STATISTICS,
550 select,
551 projects,
552 filters,
553 like_field,
554 sort,
555 &stmt.kind,
556 )
557 .await
558}
559
560pub async fn show_region(
562 stmt: ShowRegion,
563 query_engine: &QueryEngineRef,
564 catalog_manager: &CatalogManagerRef,
565 query_ctx: QueryContextRef,
566) -> Result<Output> {
567 let dataframe = show_region_dataframe(&stmt, query_engine, catalog_manager, query_ctx).await?;
568 dataframe_to_output(dataframe).await
569}
570
571pub async fn show_region_dataframe(
573 stmt: &ShowRegion,
574 query_engine: &QueryEngineRef,
575 catalog_manager: &CatalogManagerRef,
576 query_ctx: QueryContextRef,
577) -> Result<DataFrame> {
578 let schema_name = if let Some(database) = &stmt.database {
579 database.clone()
580 } else {
581 query_ctx.current_schema()
582 };
583
584 let filters = vec![
585 col(region_peers::TABLE_NAME).eq(lit(&stmt.table)),
586 col(region_peers::TABLE_SCHEMA).eq(lit(schema_name.clone())),
587 col(region_peers::TABLE_CATALOG).eq(lit(query_ctx.current_catalog())),
588 ];
589 let projects = vec![
590 (region_peers::TABLE_NAME, "Table"),
591 (region_peers::REGION_ID, "Region"),
592 (region_peers::PEER_ID, "Peer"),
593 (region_peers::IS_LEADER, "Leader"),
594 ];
595
596 let like_field = None;
597 let sort = vec![
598 col(columns::REGION_ID).sort(true, true),
599 col(columns::PEER_ID).sort(true, true),
600 ];
601
602 query_from_information_schema_dataframe(
603 query_engine,
604 catalog_manager,
605 query_ctx,
606 REGION_PEERS,
607 vec![],
608 projects,
609 filters,
610 like_field,
611 sort,
612 &stmt.kind,
613 )
614 .await
615}
616
617pub async fn show_tables(
619 stmt: ShowTables,
620 query_engine: &QueryEngineRef,
621 catalog_manager: &CatalogManagerRef,
622 query_ctx: QueryContextRef,
623) -> Result<Output> {
624 let dataframe = show_tables_dataframe(&stmt, query_engine, catalog_manager, query_ctx).await?;
625 dataframe_to_output(dataframe).await
626}
627
628pub async fn show_tables_dataframe(
630 stmt: &ShowTables,
631 query_engine: &QueryEngineRef,
632 catalog_manager: &CatalogManagerRef,
633 query_ctx: QueryContextRef,
634) -> Result<DataFrame> {
635 let schema_name = if let Some(database) = &stmt.database {
636 database.clone()
637 } else {
638 query_ctx.current_schema()
639 };
640
641 let tables_column = format!("Tables_in_{}", schema_name);
643 let projects = if stmt.full {
644 vec![
645 (tables::TABLE_NAME, tables_column.as_str()),
646 (tables::TABLE_TYPE, TABLE_TYPE_COLUMN),
647 ]
648 } else {
649 vec![(tables::TABLE_NAME, tables_column.as_str())]
650 };
651 let filters = vec![
652 col(tables::TABLE_SCHEMA).eq(lit(schema_name.clone())),
653 col(tables::TABLE_CATALOG).eq(lit(query_ctx.current_catalog())),
654 ];
655 let like_field = Some(tables::TABLE_NAME);
656 let sort = vec![col(tables::TABLE_NAME).sort(true, true)];
657
658 let rewritten_kind = match &stmt.kind {
661 ShowKind::Where(filter) => {
662 let mut filter = filter.clone();
663 replace_column_in_expr(&mut filter, "Tables", &tables_column);
664 ShowKind::Where(filter)
665 }
666 kind => kind.clone(),
667 };
668
669 query_from_information_schema_dataframe(
670 query_engine,
671 catalog_manager,
672 query_ctx,
673 TABLES,
674 vec![],
675 projects,
676 filters,
677 like_field,
678 sort,
679 &rewritten_kind,
680 )
681 .await
682}
683
684pub async fn show_table_status(
686 stmt: ShowTableStatus,
687 query_engine: &QueryEngineRef,
688 catalog_manager: &CatalogManagerRef,
689 query_ctx: QueryContextRef,
690) -> Result<Output> {
691 let dataframe =
692 show_table_status_dataframe(&stmt, query_engine, catalog_manager, query_ctx).await?;
693 dataframe_to_output(dataframe).await
694}
695
696pub async fn show_table_status_dataframe(
698 stmt: &ShowTableStatus,
699 query_engine: &QueryEngineRef,
700 catalog_manager: &CatalogManagerRef,
701 query_ctx: QueryContextRef,
702) -> Result<DataFrame> {
703 let schema_name = if let Some(database) = &stmt.database {
704 database.clone()
705 } else {
706 query_ctx.current_schema()
707 };
708
709 let projects = vec![
711 (tables::TABLE_NAME, "Name"),
712 (tables::ENGINE, "Engine"),
713 (tables::VERSION, "Version"),
714 (tables::ROW_FORMAT, "Row_format"),
715 (tables::TABLE_ROWS, "Rows"),
716 (tables::AVG_ROW_LENGTH, "Avg_row_length"),
717 (tables::DATA_LENGTH, "Data_length"),
718 (tables::MAX_DATA_LENGTH, "Max_data_length"),
719 (tables::INDEX_LENGTH, "Index_length"),
720 (tables::DATA_FREE, "Data_free"),
721 (tables::AUTO_INCREMENT, "Auto_increment"),
722 (tables::CREATE_TIME, "Create_time"),
723 (tables::UPDATE_TIME, "Update_time"),
724 (tables::CHECK_TIME, "Check_time"),
725 (tables::TABLE_COLLATION, "Collation"),
726 (tables::CHECKSUM, "Checksum"),
727 (tables::CREATE_OPTIONS, "Create_options"),
728 (tables::TABLE_COMMENT, "Comment"),
729 ];
730
731 let filters = vec![
732 col(tables::TABLE_SCHEMA).eq(lit(schema_name.clone())),
733 col(tables::TABLE_CATALOG).eq(lit(query_ctx.current_catalog())),
734 ];
735 let like_field = Some(tables::TABLE_NAME);
736 let sort = vec![col(tables::TABLE_NAME).sort(true, true)];
737
738 query_from_information_schema_dataframe(
739 query_engine,
740 catalog_manager,
741 query_ctx,
742 TABLES,
743 vec![],
744 projects,
745 filters,
746 like_field,
747 sort,
748 &stmt.kind,
749 )
750 .await
751}
752
753pub async fn show_collations(
755 kind: ShowKind,
756 query_engine: &QueryEngineRef,
757 catalog_manager: &CatalogManagerRef,
758 query_ctx: QueryContextRef,
759) -> Result<Output> {
760 let dataframe =
761 show_collations_dataframe(&kind, query_engine, catalog_manager, query_ctx).await?;
762 dataframe_to_output(dataframe).await
763}
764
765pub async fn show_collations_dataframe(
767 kind: &ShowKind,
768 query_engine: &QueryEngineRef,
769 catalog_manager: &CatalogManagerRef,
770 query_ctx: QueryContextRef,
771) -> Result<DataFrame> {
772 let projects = vec![
774 ("collation_name", "Collation"),
775 ("character_set_name", "Charset"),
776 ("id", "Id"),
777 ("is_default", "Default"),
778 ("is_compiled", "Compiled"),
779 ("sortlen", "Sortlen"),
780 ];
781
782 let filters = vec![];
783 let like_field = Some("collation_name");
784 let sort = vec![];
785
786 query_from_information_schema_dataframe(
787 query_engine,
788 catalog_manager,
789 query_ctx,
790 COLLATIONS,
791 vec![],
792 projects,
793 filters,
794 like_field,
795 sort,
796 kind,
797 )
798 .await
799}
800
801pub async fn show_charsets(
803 kind: ShowKind,
804 query_engine: &QueryEngineRef,
805 catalog_manager: &CatalogManagerRef,
806 query_ctx: QueryContextRef,
807) -> Result<Output> {
808 let dataframe =
809 show_charsets_dataframe(&kind, query_engine, catalog_manager, query_ctx).await?;
810 dataframe_to_output(dataframe).await
811}
812
813pub async fn show_charsets_dataframe(
815 kind: &ShowKind,
816 query_engine: &QueryEngineRef,
817 catalog_manager: &CatalogManagerRef,
818 query_ctx: QueryContextRef,
819) -> Result<DataFrame> {
820 let projects = vec![
822 ("character_set_name", "Charset"),
823 ("description", "Description"),
824 ("default_collate_name", "Default collation"),
825 ("maxlen", "Maxlen"),
826 ];
827
828 let filters = vec![];
829 let like_field = Some("character_set_name");
830 let sort = vec![];
831
832 query_from_information_schema_dataframe(
833 query_engine,
834 catalog_manager,
835 query_ctx,
836 CHARACTER_SETS,
837 vec![],
838 projects,
839 filters,
840 like_field,
841 sort,
842 kind,
843 )
844 .await
845}
846
847pub fn show_variable(stmt: ShowVariables, query_ctx: QueryContextRef) -> Result<Output> {
848 let variable = stmt.variable.to_string().to_uppercase();
849 let value = match variable.as_str() {
850 "SYSTEM_TIME_ZONE" | "SYSTEM_TIMEZONE" => get_timezone(None).to_string(),
851 "TIME_ZONE" | "TIMEZONE" => query_ctx.timezone().to_string(),
852 "READ_PREFERENCE" => query_ctx.read_preference().to_string(),
853 "DATESTYLE" => {
854 let (style, order) = *query_ctx.configuration_parameter().pg_datetime_style();
855 format!("{}, {}", style, order)
856 }
857 "INTERVALSTYLE" => {
858 let style = *query_ctx
859 .configuration_parameter()
860 .pg_intervalstyle_format();
861 style.to_string()
862 }
863 "MAX_EXECUTION_TIME"
864 if query_ctx.channel() == Channel::Mysql => {
865 query_ctx.query_timeout_as_millis().to_string()
866 }
867 "STATEMENT_TIMEOUT"
868 if query_ctx.channel() == Channel::Postgres => {
870 let mut timeout = query_ctx.query_timeout_as_millis().to_string();
871 timeout.push_str("ms");
872 timeout
873 }
874 _ => return UnsupportedVariableSnafu { name: variable }.fail(),
875 };
876 let schema = Arc::new(Schema::new(vec![ColumnSchema::new(
877 variable,
878 ConcreteDataType::string_datatype(),
879 false,
880 )]));
881 let records = RecordBatches::try_from_columns(
882 schema,
883 vec![Arc::new(StringVector::from(vec![value])) as _],
884 )
885 .context(error::CreateRecordBatchSnafu)?;
886 Ok(Output::new_with_record_batches(records))
887}
888
889pub async fn show_status(_query_ctx: QueryContextRef) -> Result<Output> {
890 let schema = Arc::new(Schema::new(vec![
891 ColumnSchema::new("Variable_name", ConcreteDataType::string_datatype(), false),
892 ColumnSchema::new("Value", ConcreteDataType::string_datatype(), true),
893 ]));
894 let records = RecordBatches::try_from_columns(
895 schema,
896 vec![
897 Arc::new(StringVector::from(Vec::<&str>::new())) as _,
898 Arc::new(StringVector::from(Vec::<&str>::new())) as _,
899 ],
900 )
901 .context(error::CreateRecordBatchSnafu)?;
902 Ok(Output::new_with_record_batches(records))
903}
904
905pub async fn show_search_path(_query_ctx: QueryContextRef) -> Result<Output> {
906 let schema = Arc::new(Schema::new(vec![ColumnSchema::new(
907 "search_path",
908 ConcreteDataType::string_datatype(),
909 false,
910 )]));
911 let records = RecordBatches::try_from_columns(
912 schema,
913 vec![Arc::new(StringVector::from(vec![_query_ctx.current_schema()])) as _],
914 )
915 .context(error::CreateRecordBatchSnafu)?;
916 Ok(Output::new_with_record_batches(records))
917}
918
919pub fn show_create_database(database_name: &str, options: OptionMap) -> Result<Output> {
920 let stmt = CreateDatabase {
921 name: ObjectName::from(vec![Ident::new(database_name)]),
922 if_not_exists: true,
923 options,
924 };
925 let sql = format!("{stmt}");
926 let columns = vec![
927 Arc::new(StringVector::from(vec![database_name.to_string()])) as _,
928 Arc::new(StringVector::from(vec![sql])) as _,
929 ];
930 let records =
931 RecordBatches::try_from_columns(SHOW_CREATE_DATABASE_OUTPUT_SCHEMA.clone(), columns)
932 .context(error::CreateRecordBatchSnafu)?;
933 Ok(Output::new_with_record_batches(records))
934}
935
936pub fn show_create_table(
937 table_info: TableInfoRef,
938 schema_options: Option<SchemaOptions>,
939 partitions: Option<Partitions>,
940 query_ctx: QueryContextRef,
941) -> Result<Output> {
942 let table_name = table_info.name.clone();
943
944 let quote_style = query_ctx.quote_style();
945
946 let mut stmt = create_table_stmt(&table_info, schema_options, quote_style)?;
947 stmt.partitions = partitions.map(|mut p| {
948 p.set_quote(quote_style);
949 p
950 });
951 let sql = format!("{}", stmt);
952 let columns = vec![
953 Arc::new(StringVector::from(vec![table_name])) as _,
954 Arc::new(StringVector::from(vec![sql])) as _,
955 ];
956 let records = RecordBatches::try_from_columns(SHOW_CREATE_TABLE_OUTPUT_SCHEMA.clone(), columns)
957 .context(error::CreateRecordBatchSnafu)?;
958
959 Ok(Output::new_with_record_batches(records))
960}
961
962pub fn show_create_foreign_table_for_pg(
963 table: TableRef,
964 _query_ctx: QueryContextRef,
965) -> Result<Output> {
966 let table_info = table.table_info();
967
968 let table_meta = &table_info.meta;
969 let table_name = &table_info.name;
970 let schema = &table_info.meta.schema;
971 let is_metric_engine = is_metric_engine(&table_meta.engine);
972
973 let columns = schema
974 .column_schemas()
975 .iter()
976 .filter_map(|c| {
977 if is_metric_engine && is_metric_engine_internal_column(&c.name) {
978 None
979 } else {
980 Some(format!(
981 "\"{}\" {}",
982 c.name,
983 c.data_type.postgres_datatype_name()
984 ))
985 }
986 })
987 .join(",\n ");
988
989 let sql = format!(
990 r#"CREATE FOREIGN TABLE ft_{} (
991 {}
992)
993SERVER greptimedb
994OPTIONS (table_name '{}')"#,
995 table_name, columns, table_name
996 );
997
998 let columns = vec![
999 Arc::new(StringVector::from(vec![table_name.clone()])) as _,
1000 Arc::new(StringVector::from(vec![sql])) as _,
1001 ];
1002 let records = RecordBatches::try_from_columns(SHOW_CREATE_TABLE_OUTPUT_SCHEMA.clone(), columns)
1003 .context(error::CreateRecordBatchSnafu)?;
1004
1005 Ok(Output::new_with_record_batches(records))
1006}
1007
1008pub fn show_create_view(
1009 view_name: ObjectName,
1010 definition: &str,
1011 query_ctx: QueryContextRef,
1012) -> Result<Output> {
1013 let mut parser_ctx =
1014 ParserContext::new(query_ctx.sql_dialect(), definition).context(error::SqlSnafu)?;
1015
1016 let Statement::CreateView(create_view) =
1017 parser_ctx.parse_statement().context(error::SqlSnafu)?
1018 else {
1019 unreachable!();
1021 };
1022
1023 let stmt = CreateView {
1024 name: view_name.clone(),
1025 columns: create_view.columns,
1026 query: create_view.query,
1027 or_replace: create_view.or_replace,
1028 if_not_exists: create_view.if_not_exists,
1029 };
1030
1031 let sql = format!("{}", stmt);
1032 let columns = vec![
1033 Arc::new(StringVector::from(vec![view_name.to_string()])) as _,
1034 Arc::new(StringVector::from(vec![sql])) as _,
1035 ];
1036 let records = RecordBatches::try_from_columns(SHOW_CREATE_VIEW_OUTPUT_SCHEMA.clone(), columns)
1037 .context(error::CreateRecordBatchSnafu)?;
1038
1039 Ok(Output::new_with_record_batches(records))
1040}
1041
1042pub async fn show_views(
1044 stmt: ShowViews,
1045 query_engine: &QueryEngineRef,
1046 catalog_manager: &CatalogManagerRef,
1047 query_ctx: QueryContextRef,
1048) -> Result<Output> {
1049 let dataframe = show_views_dataframe(&stmt, query_engine, catalog_manager, query_ctx).await?;
1050 dataframe_to_output(dataframe).await
1051}
1052
1053pub async fn show_views_dataframe(
1055 stmt: &ShowViews,
1056 query_engine: &QueryEngineRef,
1057 catalog_manager: &CatalogManagerRef,
1058 query_ctx: QueryContextRef,
1059) -> Result<DataFrame> {
1060 let schema_name = if let Some(database) = &stmt.database {
1061 database.clone()
1062 } else {
1063 query_ctx.current_schema()
1064 };
1065
1066 let projects = vec![(tables::TABLE_NAME, VIEWS_COLUMN)];
1067 let filters = vec![
1068 col(tables::TABLE_SCHEMA).eq(lit(schema_name.clone())),
1069 col(tables::TABLE_CATALOG).eq(lit(query_ctx.current_catalog())),
1070 ];
1071 let like_field = Some(tables::TABLE_NAME);
1072 let sort = vec![col(tables::TABLE_NAME).sort(true, true)];
1073
1074 query_from_information_schema_dataframe(
1075 query_engine,
1076 catalog_manager,
1077 query_ctx,
1078 VIEWS,
1079 vec![],
1080 projects,
1081 filters,
1082 like_field,
1083 sort,
1084 &stmt.kind,
1085 )
1086 .await
1087}
1088
1089pub async fn show_flows(
1091 stmt: ShowFlows,
1092 query_engine: &QueryEngineRef,
1093 catalog_manager: &CatalogManagerRef,
1094 query_ctx: QueryContextRef,
1095) -> Result<Output> {
1096 let dataframe = show_flows_dataframe(&stmt, query_engine, catalog_manager, query_ctx).await?;
1097 dataframe_to_output(dataframe).await
1098}
1099
1100pub async fn show_flows_dataframe(
1102 stmt: &ShowFlows,
1103 query_engine: &QueryEngineRef,
1104 catalog_manager: &CatalogManagerRef,
1105 query_ctx: QueryContextRef,
1106) -> Result<DataFrame> {
1107 let projects = vec![(flows::FLOW_NAME, FLOWS_COLUMN)];
1108 let filters = vec![col(flows::TABLE_CATALOG).eq(lit(query_ctx.current_catalog()))];
1109 let like_field = Some(flows::FLOW_NAME);
1110 let sort = vec![col(flows::FLOW_NAME).sort(true, true)];
1111
1112 query_from_information_schema_dataframe(
1113 query_engine,
1114 catalog_manager,
1115 query_ctx,
1116 FLOWS,
1117 vec![],
1118 projects,
1119 filters,
1120 like_field,
1121 sort,
1122 &stmt.kind,
1123 )
1124 .await
1125}
1126
1127pub async fn show_flow_status(
1129 stmt: ShowFlowStatus,
1130 query_engine: &QueryEngineRef,
1131 catalog_manager: &CatalogManagerRef,
1132 query_ctx: QueryContextRef,
1133) -> Result<Output> {
1134 let projects = vec![
1135 (flow_statistics::FLOW_ID, flow_statistics::FLOW_ID),
1136 (flow_statistics::FLOW_NAME, flow_statistics::FLOW_NAME),
1137 (flow_statistics::START_TIME, flow_statistics::START_TIME),
1138 (
1139 flow_statistics::LAST_EXECUTION_TIME,
1140 flow_statistics::LAST_EXECUTION_TIME,
1141 ),
1142 (
1143 flow_statistics::UPTIME_SECONDS,
1144 flow_statistics::UPTIME_SECONDS,
1145 ),
1146 (flow_statistics::STATE_SIZE, flow_statistics::STATE_SIZE),
1147 ];
1148 let like_field = Some(flow_statistics::FLOW_NAME);
1149 let sort = vec![col(flow_statistics::FLOW_NAME).sort(true, true)];
1150
1151 query_from_information_schema_table(
1152 query_engine,
1153 catalog_manager,
1154 query_ctx,
1155 FLOW_STATISTICS,
1156 vec![],
1157 projects,
1158 vec![],
1159 like_field,
1160 sort,
1161 stmt.kind,
1162 )
1163 .await
1164}
1165
1166#[cfg(feature = "enterprise")]
1167pub async fn show_triggers(
1168 stmt: sql::statements::show::trigger::ShowTriggers,
1169 query_engine: &QueryEngineRef,
1170 catalog_manager: &CatalogManagerRef,
1171 query_ctx: QueryContextRef,
1172) -> Result<Output> {
1173 const TRIGGER_NAME: &str = "trigger_name";
1174 const TRIGGERS_COLUMN: &str = "Triggers";
1175
1176 let projects = vec![(TRIGGER_NAME, TRIGGERS_COLUMN)];
1177 let like_field = Some(TRIGGER_NAME);
1178 let sort = vec![col(TRIGGER_NAME).sort(true, true)];
1179
1180 query_from_information_schema_table(
1181 query_engine,
1182 catalog_manager,
1183 query_ctx,
1184 catalog::information_schema::TRIGGERS,
1185 vec![],
1186 projects,
1187 vec![],
1188 like_field,
1189 sort,
1190 stmt.kind,
1191 )
1192 .await
1193}
1194
1195pub fn show_create_flow(
1196 flow_name: ObjectName,
1197 flow_val: FlowInfoValue,
1198 query_ctx: QueryContextRef,
1199) -> Result<Output> {
1200 let mut parser_ctx =
1201 ParserContext::new(query_ctx.sql_dialect(), flow_val.raw_sql()).context(error::SqlSnafu)?;
1202
1203 let query = parser_ctx.parse_statement().context(error::SqlSnafu)?;
1204
1205 let raw_query = match &query {
1207 Statement::Tql(_) => flow_val.raw_sql().clone(),
1208 _ => query.to_string(),
1209 };
1210
1211 let query = Box::new(SqlOrTql::try_from_statement(query, &raw_query).context(error::SqlSnafu)?);
1212
1213 let comment = if flow_val.comment().is_empty() {
1214 None
1215 } else {
1216 Some(flow_val.comment().clone())
1217 };
1218
1219 let stmt = CreateFlow {
1220 flow_name,
1221 sink_table_name: ObjectName::from(vec![
1222 Ident::new(&flow_val.sink_table_name().schema_name),
1223 Ident::new(&flow_val.sink_table_name().table_name),
1224 ]),
1225 or_replace: false,
1228 if_not_exists: true,
1229 expire_after: flow_val.expire_after(),
1230 eval_interval: flow_val.eval_interval(),
1231 eval_offset: effective_eval_schedule_from_flow_info(&flow_val)
1232 .map_err(BoxedError::new)
1233 .context(error::QueryExecutionSnafu)?
1234 .map(|schedule| schedule.anchor_secs)
1235 .filter(|anchor_secs| *anchor_secs != 0),
1236 comment,
1237 flow_options: OptionMap::from_filtered_string_map(
1238 flow_val.options(),
1239 &[FlowType::FLOW_TYPE_KEY],
1240 ),
1241 query,
1242 };
1243
1244 let sql = format!("{}", stmt);
1245 let columns = vec![
1246 Arc::new(StringVector::from(vec![flow_val.flow_name().clone()])) as _,
1247 Arc::new(StringVector::from(vec![sql])) as _,
1248 ];
1249 let records = RecordBatches::try_from_columns(SHOW_CREATE_FLOW_OUTPUT_SCHEMA.clone(), columns)
1250 .context(error::CreateRecordBatchSnafu)?;
1251
1252 Ok(Output::new_with_record_batches(records))
1253}
1254
1255pub fn describe_table(table: TableRef) -> Result<Output> {
1256 let table_info = table.table_info();
1257 let columns_schemas = table_info.meta.schema.column_schemas();
1258 let columns = vec![
1259 describe_column_names(columns_schemas),
1260 describe_column_types(columns_schemas),
1261 describe_column_keys(columns_schemas, &table_info.meta.primary_key_indices),
1262 describe_column_nullables(columns_schemas),
1263 describe_column_defaults(columns_schemas),
1264 describe_column_semantic_types(columns_schemas, &table_info.meta.primary_key_indices),
1265 ];
1266 let records = RecordBatches::try_from_columns(DESCRIBE_TABLE_OUTPUT_SCHEMA.clone(), columns)
1267 .context(error::CreateRecordBatchSnafu)?;
1268 Ok(Output::new_with_record_batches(records))
1269}
1270
1271fn describe_column_names(columns_schemas: &[ColumnSchema]) -> VectorRef {
1272 Arc::new(StringVector::from_iterator(
1273 columns_schemas.iter().map(|cs| cs.name.as_str()),
1274 ))
1275}
1276
1277fn describe_column_types(columns_schemas: &[ColumnSchema]) -> VectorRef {
1278 Arc::new(StringVector::from(
1279 columns_schemas
1280 .iter()
1281 .map(|cs| describe_column_type_name(&cs.data_type))
1282 .collect::<Vec<_>>(),
1283 ))
1284}
1285
1286fn describe_column_type_name(data_type: &ConcreteDataType) -> String {
1287 match data_type {
1288 ConcreteDataType::Json(json_type) if json_type.is_json2() => "Json2".to_string(),
1289 data_type => data_type.name(),
1290 }
1291}
1292
1293fn describe_column_keys(
1294 columns_schemas: &[ColumnSchema],
1295 primary_key_indices: &[usize],
1296) -> VectorRef {
1297 Arc::new(StringVector::from_iterator(
1298 columns_schemas.iter().enumerate().map(|(i, cs)| {
1299 if cs.is_time_index() || primary_key_indices.contains(&i) {
1300 PRI_KEY
1301 } else {
1302 ""
1303 }
1304 }),
1305 ))
1306}
1307
1308fn describe_column_nullables(columns_schemas: &[ColumnSchema]) -> VectorRef {
1309 Arc::new(StringVector::from_iterator(columns_schemas.iter().map(
1310 |cs| {
1311 if cs.is_nullable() { YES_STR } else { NO_STR }
1312 },
1313 )))
1314}
1315
1316fn describe_column_defaults(columns_schemas: &[ColumnSchema]) -> VectorRef {
1317 Arc::new(StringVector::from(
1318 columns_schemas
1319 .iter()
1320 .map(|cs| {
1321 cs.default_constraint()
1322 .map_or(String::from(""), |dc| dc.to_string())
1323 })
1324 .collect::<Vec<String>>(),
1325 ))
1326}
1327
1328fn describe_column_semantic_types(
1329 columns_schemas: &[ColumnSchema],
1330 primary_key_indices: &[usize],
1331) -> VectorRef {
1332 Arc::new(StringVector::from_iterator(
1333 columns_schemas.iter().enumerate().map(|(i, cs)| {
1334 if primary_key_indices.contains(&i) {
1335 SEMANTIC_TYPE_PRIMARY_KEY
1336 } else if cs.is_time_index() {
1337 SEMANTIC_TYPE_TIME_INDEX
1338 } else {
1339 SEMANTIC_TYPE_FIELD
1340 }
1341 }),
1342 ))
1343}
1344
1345pub async fn prepare_file_table_files(
1347 options: &HashMap<String, String>,
1348 local_file_access: &LocalFileAccess,
1349) -> Result<(ObjectStore, Vec<String>)> {
1350 let url = options
1351 .get(FILE_TABLE_LOCATION_KEY)
1352 .context(error::MissingRequiredFieldSnafu {
1353 name: FILE_TABLE_LOCATION_KEY,
1354 })?;
1355
1356 let regex = options
1357 .get(FILE_TABLE_PATTERN_KEY)
1358 .map(|x| Regex::new(x))
1359 .transpose()
1360 .context(error::BuildRegexSnafu)?;
1361 let backend = build_backend_with_path(url, options, local_file_access)
1362 .await
1363 .context(error::BuildBackendSnafu)?;
1364 let source = if let Some(filename) = backend.object_path {
1365 Source::Filename(filename)
1366 } else {
1367 Source::Dir
1368 };
1369 let lister = Lister::new(backend.object_store.clone(), source, url.clone(), regex);
1370 let files = lister
1375 .list()
1376 .await
1377 .context(error::ListObjectsSnafu)?
1378 .into_iter()
1379 .filter_map(|entry| {
1380 if entry.path().ends_with('/') {
1381 None
1382 } else {
1383 Some(entry.path().to_string())
1384 }
1385 })
1386 .collect::<Vec<_>>();
1387 Ok((backend.object_store, files))
1388}
1389
1390pub async fn infer_file_table_schema(
1391 object_store: &ObjectStore,
1392 files: &[String],
1393 options: &HashMap<String, String>,
1394) -> Result<Schema> {
1395 let format = parse_file_table_format(options)?;
1396 let merged = infer_schemas(object_store, files, format.as_ref())
1397 .await
1398 .context(error::InferSchemaSnafu)?;
1399 Schema::try_from(merged).context(error::ConvertSchemaSnafu)
1400}
1401
1402pub fn file_column_schemas_to_table(
1412 file_column_schemas: &[ColumnSchema],
1413) -> (Vec<ColumnSchema>, String) {
1414 let mut column_schemas = file_column_schemas.to_owned();
1415 if let Some(time_index_column) = column_schemas.iter().find(|c| c.is_time_index()) {
1416 let time_index = time_index_column.name.clone();
1417 return (column_schemas, time_index);
1418 }
1419
1420 let timestamp_type = ConcreteDataType::timestamp_millisecond_datatype();
1421 let default_zero = Value::Timestamp(Timestamp::new_millisecond(0));
1422 let timestamp_column_schema = ColumnSchema::new(greptime_timestamp(), timestamp_type, false)
1423 .with_time_index(true)
1424 .with_default_constraint(Some(ColumnDefaultConstraint::Value(default_zero)))
1425 .unwrap();
1426
1427 if let Some(column_schema) = column_schemas
1428 .iter_mut()
1429 .find(|column_schema| column_schema.name == greptime_timestamp())
1430 {
1431 *column_schema = timestamp_column_schema;
1433 } else {
1434 column_schemas.push(timestamp_column_schema);
1435 }
1436
1437 (column_schemas, greptime_timestamp().to_string())
1438}
1439
1440pub fn check_file_to_table_schema_compatibility(
1449 file_column_schemas: &[ColumnSchema],
1450 table_column_schemas: &[ColumnSchema],
1451) -> Result<()> {
1452 let file_schemas_map = file_column_schemas
1453 .iter()
1454 .map(|s| (s.name.clone(), s))
1455 .collect::<HashMap<_, _>>();
1456
1457 for table_column in table_column_schemas {
1458 if let Some(file_column) = file_schemas_map.get(&table_column.name) {
1459 ensure!(
1461 file_column
1462 .data_type
1463 .can_arrow_type_cast_to(&table_column.data_type),
1464 error::ColumnSchemaIncompatibleSnafu {
1465 column: table_column.name.clone(),
1466 file_type: file_column.data_type.clone(),
1467 table_type: table_column.data_type.clone(),
1468 }
1469 );
1470 } else {
1471 ensure!(
1472 table_column.is_nullable() || table_column.default_constraint().is_some(),
1473 error::ColumnSchemaNoDefaultSnafu {
1474 column: table_column.name.clone(),
1475 }
1476 );
1477 }
1478 }
1479
1480 Ok(())
1481}
1482
1483fn parse_file_table_format(options: &HashMap<String, String>) -> Result<Box<dyn FileFormat>> {
1484 Ok(
1485 match Format::try_from(options).context(error::ParseFileFormatSnafu)? {
1486 Format::Csv(format) => Box::new(format),
1487 Format::Json(format) => Box::new(format),
1488 Format::Parquet(format) => Box::new(format),
1489 Format::Orc(format) => Box::new(format),
1490 },
1491 )
1492}
1493
1494pub async fn show_processlist(
1495 stmt: ShowProcessList,
1496 query_engine: &QueryEngineRef,
1497 catalog_manager: &CatalogManagerRef,
1498 query_ctx: QueryContextRef,
1499) -> Result<Output> {
1500 let dataframe =
1501 show_processlist_dataframe(&stmt, query_engine, catalog_manager, query_ctx).await?;
1502 dataframe_to_output(dataframe).await
1503}
1504
1505pub async fn show_processlist_dataframe(
1507 stmt: &ShowProcessList,
1508 query_engine: &QueryEngineRef,
1509 catalog_manager: &CatalogManagerRef,
1510 query_ctx: QueryContextRef,
1511) -> Result<DataFrame> {
1512 let projects = if stmt.full {
1513 vec![
1514 (process_list::ID, "Id"),
1515 (process_list::CATALOG, "Catalog"),
1516 (process_list::SCHEMAS, "Schema"),
1517 (process_list::CLIENT, "Client"),
1518 (process_list::FRONTEND, "Frontend"),
1519 (process_list::START_TIMESTAMP, "Start Time"),
1520 (process_list::ELAPSED_TIME, "Elapsed Time"),
1521 (process_list::QUERY, "Query"),
1522 ]
1523 } else {
1524 vec![
1525 (process_list::ID, "Id"),
1526 (process_list::CATALOG, "Catalog"),
1527 (process_list::QUERY, "Query"),
1528 (process_list::ELAPSED_TIME, "Elapsed Time"),
1529 ]
1530 };
1531
1532 let filters = if query_ctx.current_user().is_admin() {
1533 vec![]
1534 } else {
1535 vec![col(process_list::CATALOG).eq(lit(query_ctx.current_catalog()))]
1536 };
1537 let like_field = None;
1538 let sort = vec![col("id").sort(true, true)];
1539 query_from_information_schema_dataframe(
1540 query_engine,
1541 catalog_manager,
1542 query_ctx.clone(),
1543 "process_list",
1544 vec![],
1545 projects,
1546 filters,
1547 like_field,
1548 sort,
1549 &ShowKind::All,
1550 )
1551 .await
1552}
1553
1554#[cfg(test)]
1555mod test {
1556 use std::sync::Arc;
1557
1558 use common_query::{Output, OutputData};
1559 use common_recordbatch::{RecordBatch, RecordBatches};
1560 use common_time::Timezone;
1561 use common_time::timestamp::TimeUnit;
1562 use datatypes::prelude::ConcreteDataType;
1563 use datatypes::schema::{ColumnDefaultConstraint, ColumnSchema, Schema, SchemaRef};
1564 use datatypes::types::json_type::JsonNativeType;
1565 use datatypes::vectors::{StringVector, TimestampMillisecondVector, UInt32Vector, VectorRef};
1566 use session::context::QueryContextBuilder;
1567 use snafu::ResultExt;
1568 use sql::ast::{Ident, ObjectName};
1569 use sql::statements::show::ShowVariables;
1570 use table::TableRef;
1571 use table::test_util::MemTable;
1572
1573 use super::{describe_column_type_name, show_variable};
1574 use crate::error;
1575 use crate::error::Result;
1576 use crate::sql::{
1577 DESCRIBE_TABLE_OUTPUT_SCHEMA, NO_STR, SEMANTIC_TYPE_FIELD, SEMANTIC_TYPE_TIME_INDEX,
1578 YES_STR, describe_table,
1579 };
1580
1581 #[test]
1582 fn test_describe_table_multiple_columns() -> Result<()> {
1583 let table_name = "test_table";
1584 let schema = vec![
1585 ColumnSchema::new("t1", ConcreteDataType::uint32_datatype(), true),
1586 ColumnSchema::new(
1587 "t2",
1588 ConcreteDataType::timestamp_datatype(TimeUnit::Millisecond),
1589 false,
1590 )
1591 .with_default_constraint(Some(ColumnDefaultConstraint::Function(String::from(
1592 "current_timestamp()",
1593 ))))
1594 .unwrap()
1595 .with_time_index(true),
1596 ];
1597 let data = vec![
1598 Arc::new(UInt32Vector::from_slice([0])) as _,
1599 Arc::new(TimestampMillisecondVector::from_slice([0])) as _,
1600 ];
1601 let expected_columns = vec![
1602 Arc::new(StringVector::from(vec!["t1", "t2"])) as _,
1603 Arc::new(StringVector::from(vec!["UInt32", "TimestampMillisecond"])) as _,
1604 Arc::new(StringVector::from(vec!["", "PRI"])) as _,
1605 Arc::new(StringVector::from(vec![YES_STR, NO_STR])) as _,
1606 Arc::new(StringVector::from(vec!["", "current_timestamp()"])) as _,
1607 Arc::new(StringVector::from(vec![
1608 SEMANTIC_TYPE_FIELD,
1609 SEMANTIC_TYPE_TIME_INDEX,
1610 ])) as _,
1611 ];
1612
1613 describe_table_test_by_schema(table_name, schema, data, expected_columns)
1614 }
1615
1616 #[test]
1617 fn test_describe_column_type_name_json2() {
1618 assert_eq!(
1619 describe_column_type_name(&ConcreteDataType::json2(JsonNativeType::Null)),
1620 "Json2"
1621 );
1622 assert_eq!(
1623 describe_column_type_name(&ConcreteDataType::uint32_datatype()),
1624 "UInt32"
1625 );
1626 }
1627
1628 fn describe_table_test_by_schema(
1629 table_name: &str,
1630 schema: Vec<ColumnSchema>,
1631 data: Vec<VectorRef>,
1632 expected_columns: Vec<VectorRef>,
1633 ) -> Result<()> {
1634 let table_schema = SchemaRef::new(Schema::new(schema));
1635 let table = prepare_describe_table(table_name, table_schema, data);
1636
1637 let expected =
1638 RecordBatches::try_from_columns(DESCRIBE_TABLE_OUTPUT_SCHEMA.clone(), expected_columns)
1639 .context(error::CreateRecordBatchSnafu)?;
1640
1641 if let OutputData::RecordBatches(res) = describe_table(table)?.data {
1642 assert_eq!(res.take(), expected.take());
1643 } else {
1644 panic!("describe table must return record batch");
1645 }
1646
1647 Ok(())
1648 }
1649
1650 fn prepare_describe_table(
1651 table_name: &str,
1652 table_schema: SchemaRef,
1653 data: Vec<VectorRef>,
1654 ) -> TableRef {
1655 let record_batch = RecordBatch::new(table_schema, data).unwrap();
1656 MemTable::table(table_name, record_batch)
1657 }
1658
1659 #[test]
1660 fn test_show_variable() {
1661 assert_eq!(
1662 exec_show_variable("SYSTEM_TIME_ZONE", "Asia/Shanghai").unwrap(),
1663 "UTC"
1664 );
1665 assert_eq!(
1666 exec_show_variable("SYSTEM_TIMEZONE", "Asia/Shanghai").unwrap(),
1667 "UTC"
1668 );
1669 assert_eq!(
1670 exec_show_variable("TIME_ZONE", "Asia/Shanghai").unwrap(),
1671 "Asia/Shanghai"
1672 );
1673 assert_eq!(
1674 exec_show_variable("TIMEZONE", "Asia/Shanghai").unwrap(),
1675 "Asia/Shanghai"
1676 );
1677 assert!(exec_show_variable("TIME ZONE", "Asia/Shanghai").is_err());
1678 assert!(exec_show_variable("SYSTEM TIME ZONE", "Asia/Shanghai").is_err());
1679 }
1680
1681 fn exec_show_variable(variable: &str, tz: &str) -> Result<String> {
1682 let stmt = ShowVariables {
1683 variable: ObjectName::from(vec![Ident::new(variable)]),
1684 };
1685 let ctx = Arc::new(
1686 QueryContextBuilder::default()
1687 .timezone(Timezone::from_tz_string(tz).unwrap())
1688 .build(),
1689 );
1690 match show_variable(stmt, ctx) {
1691 Ok(Output {
1692 data: OutputData::RecordBatches(record),
1693 ..
1694 }) => {
1695 let record = record.take().first().cloned().unwrap();
1696 Ok(record.iter_column_as_string(0).next().unwrap().unwrap())
1697 }
1698 Ok(_) => unreachable!(),
1699 Err(e) => Err(e),
1700 }
1701 }
1702}