Skip to main content

query/
sql.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15mod 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_meta::SchemaOptions;
36use common_meta::ddl::create_flow::FlowType;
37use common_meta::key::flow::flow_info::FlowInfoValue;
38use common_query::Output;
39use common_query::prelude::greptime_timestamp;
40use common_recordbatch::RecordBatches;
41use common_recordbatch::adapter::RecordBatchStreamAdapter;
42use common_time::Timestamp;
43use common_time::timezone::get_timezone;
44use datafusion::prelude::SessionContext;
45use datafusion_expr::{Expr, SortExpr, col, lit};
46use datatypes::prelude::*;
47use datatypes::schema::{ColumnDefaultConstraint, ColumnSchema, Schema};
48use datatypes::vectors::StringVector;
49use itertools::Itertools;
50use object_store::ObjectStore;
51use once_cell::sync::Lazy;
52use regex::Regex;
53use session::context::{Channel, QueryContextRef};
54pub use show_create_table::create_table_stmt;
55use snafu::{OptionExt, ResultExt, ensure};
56use sql::ast::{Ident, visit_expressions_mut};
57use sql::parser::ParserContext;
58use sql::statements::OptionMap;
59use sql::statements::create::{CreateDatabase, CreateFlow, CreateView, Partitions, SqlOrTql};
60use sql::statements::show::{
61    ShowColumns, ShowDatabases, ShowFlowStatus, ShowFlows, ShowIndex, ShowKind, ShowProcessList,
62    ShowRegion, ShowTableStatus, ShowTables, ShowVariables, ShowViews,
63};
64use sql::statements::statement::Statement;
65use sqlparser::ast::ObjectName;
66use store_api::metric_engine_consts::{is_metric_engine, is_metric_engine_internal_column};
67use table::TableRef;
68use table::metadata::TableInfoRef;
69use table::requests::{FILE_TABLE_LOCATION_KEY, FILE_TABLE_PATTERN_KEY};
70
71use crate::QueryEngineRef;
72use crate::error::{self, Result, UnsupportedVariableSnafu};
73use crate::planner::DfLogicalPlanner;
74
75const SCHEMAS_COLUMN: &str = "Database";
76const OPTIONS_COLUMN: &str = "Options";
77const VIEWS_COLUMN: &str = "Views";
78const FLOWS_COLUMN: &str = "Flows";
79const FIELD_COLUMN: &str = "Field";
80const TABLE_TYPE_COLUMN: &str = "Table_type";
81
82const COLUMN_NAME_COLUMN: &str = "Column";
83const COLUMN_GREPTIME_TYPE_COLUMN: &str = "Greptime_type";
84const COLUMN_TYPE_COLUMN: &str = "Type";
85const COLUMN_KEY_COLUMN: &str = "Key";
86const COLUMN_EXTRA_COLUMN: &str = "Extra";
87const COLUMN_PRIVILEGES_COLUMN: &str = "Privileges";
88const COLUMN_COLLATION_COLUMN: &str = "Collation";
89const COLUMN_NULLABLE_COLUMN: &str = "Null";
90const COLUMN_DEFAULT_COLUMN: &str = "Default";
91const COLUMN_COMMENT_COLUMN: &str = "Comment";
92const COLUMN_SEMANTIC_TYPE_COLUMN: &str = "Semantic Type";
93
94const YES_STR: &str = "YES";
95const NO_STR: &str = "NO";
96const PRI_KEY: &str = "PRI";
97
98/// SHOW index columns
99const INDEX_TABLE_COLUMN: &str = "Table";
100const INDEX_NONT_UNIQUE_COLUMN: &str = "Non_unique";
101const INDEX_CARDINALITY_COLUMN: &str = "Cardinality";
102const INDEX_SUB_PART_COLUMN: &str = "Sub_part";
103const INDEX_PACKED_COLUMN: &str = "Packed";
104const INDEX_INDEX_TYPE_COLUMN: &str = "Index_type";
105const INDEX_COMMENT_COLUMN: &str = "Index_comment";
106const INDEX_VISIBLE_COLUMN: &str = "Visible";
107const INDEX_EXPRESSION_COLUMN: &str = "Expression";
108const INDEX_KEY_NAME_COLUMN: &str = "Key_name";
109const INDEX_SEQ_IN_INDEX_COLUMN: &str = "Seq_in_index";
110const INDEX_COLUMN_NAME_COLUMN: &str = "Column_name";
111
112static DESCRIBE_TABLE_OUTPUT_SCHEMA: Lazy<Arc<Schema>> = Lazy::new(|| {
113    Arc::new(Schema::new(vec![
114        ColumnSchema::new(
115            COLUMN_NAME_COLUMN,
116            ConcreteDataType::string_datatype(),
117            false,
118        ),
119        ColumnSchema::new(
120            COLUMN_TYPE_COLUMN,
121            ConcreteDataType::string_datatype(),
122            false,
123        ),
124        ColumnSchema::new(COLUMN_KEY_COLUMN, ConcreteDataType::string_datatype(), true),
125        ColumnSchema::new(
126            COLUMN_NULLABLE_COLUMN,
127            ConcreteDataType::string_datatype(),
128            false,
129        ),
130        ColumnSchema::new(
131            COLUMN_DEFAULT_COLUMN,
132            ConcreteDataType::string_datatype(),
133            false,
134        ),
135        ColumnSchema::new(
136            COLUMN_SEMANTIC_TYPE_COLUMN,
137            ConcreteDataType::string_datatype(),
138            false,
139        ),
140    ]))
141});
142
143static SHOW_CREATE_DATABASE_OUTPUT_SCHEMA: Lazy<Arc<Schema>> = Lazy::new(|| {
144    Arc::new(Schema::new(vec![
145        ColumnSchema::new("Database", ConcreteDataType::string_datatype(), false),
146        ColumnSchema::new(
147            "Create Database",
148            ConcreteDataType::string_datatype(),
149            false,
150        ),
151    ]))
152});
153
154static SHOW_CREATE_TABLE_OUTPUT_SCHEMA: Lazy<Arc<Schema>> = Lazy::new(|| {
155    Arc::new(Schema::new(vec![
156        ColumnSchema::new("Table", ConcreteDataType::string_datatype(), false),
157        ColumnSchema::new("Create Table", ConcreteDataType::string_datatype(), false),
158    ]))
159});
160
161static SHOW_CREATE_FLOW_OUTPUT_SCHEMA: Lazy<Arc<Schema>> = Lazy::new(|| {
162    Arc::new(Schema::new(vec![
163        ColumnSchema::new("Flow", ConcreteDataType::string_datatype(), false),
164        ColumnSchema::new("Create Flow", ConcreteDataType::string_datatype(), false),
165    ]))
166});
167
168static SHOW_CREATE_VIEW_OUTPUT_SCHEMA: Lazy<Arc<Schema>> = Lazy::new(|| {
169    Arc::new(Schema::new(vec![
170        ColumnSchema::new("View", ConcreteDataType::string_datatype(), false),
171        ColumnSchema::new("Create View", ConcreteDataType::string_datatype(), false),
172    ]))
173});
174
175pub async fn show_databases(
176    stmt: ShowDatabases,
177    query_engine: &QueryEngineRef,
178    catalog_manager: &CatalogManagerRef,
179    query_ctx: QueryContextRef,
180) -> Result<Output> {
181    let projects = if stmt.full {
182        vec![
183            (schemata::SCHEMA_NAME, SCHEMAS_COLUMN),
184            (schemata::SCHEMA_OPTS, OPTIONS_COLUMN),
185        ]
186    } else {
187        vec![(schemata::SCHEMA_NAME, SCHEMAS_COLUMN)]
188    };
189
190    let filters = vec![col(schemata::CATALOG_NAME).eq(lit(query_ctx.current_catalog()))];
191    let like_field = Some(schemata::SCHEMA_NAME);
192    let sort = vec![col(schemata::SCHEMA_NAME).sort(true, true)];
193
194    query_from_information_schema_table(
195        query_engine,
196        catalog_manager,
197        query_ctx,
198        SCHEMATA,
199        vec![],
200        projects,
201        filters,
202        like_field,
203        sort,
204        stmt.kind,
205    )
206    .await
207}
208
209/// Replaces column identifier references in a SQL expression.
210/// Used for backward compatibility where old column names should work with new ones.
211fn replace_column_in_expr(expr: &mut sqlparser::ast::Expr, from_column: &str, to_column: &str) {
212    let _ = visit_expressions_mut(expr, |e| {
213        match e {
214            sqlparser::ast::Expr::Identifier(ident)
215                if ident.value.eq_ignore_ascii_case(from_column) =>
216            {
217                ident.value = to_column.to_string();
218            }
219            sqlparser::ast::Expr::CompoundIdentifier(idents) => {
220                if let Some(last) = idents.last_mut()
221                    && last.value.eq_ignore_ascii_case(from_column)
222                {
223                    last.value = to_column.to_string();
224                }
225            }
226            _ => {}
227        }
228        ControlFlow::<()>::Continue(())
229    });
230}
231
232/// Cast a `show` statement execution into a query from tables in  `information_schema`.
233/// - `table_name`: the table name in `information_schema`,
234/// - `projects`: query projection, a list of `(column, renamed_column)`,
235/// - `filters`: filter expressions for query,
236/// - `like_field`: the field to filter by the predicate `ShowKind::Like`,
237/// - `sort`: sort the results by the specified sorting expressions,
238/// - `kind`: the show kind
239#[allow(clippy::too_many_arguments)]
240async fn query_from_information_schema_table(
241    query_engine: &QueryEngineRef,
242    catalog_manager: &CatalogManagerRef,
243    query_ctx: QueryContextRef,
244    table_name: &str,
245    select: Vec<Expr>,
246    projects: Vec<(&str, &str)>,
247    filters: Vec<Expr>,
248    like_field: Option<&str>,
249    sort: Vec<SortExpr>,
250    kind: ShowKind,
251) -> Result<Output> {
252    let table = catalog_manager
253        .table(
254            query_ctx.current_catalog(),
255            INFORMATION_SCHEMA_NAME,
256            table_name,
257            Some(&query_ctx),
258        )
259        .await
260        .context(error::CatalogSnafu)?
261        .with_context(|| error::TableNotFoundSnafu {
262            table: format_full_table_name(
263                query_ctx.current_catalog(),
264                INFORMATION_SCHEMA_NAME,
265                table_name,
266            ),
267        })?;
268
269    let dataframe = query_engine.read_table(table)?;
270
271    // Apply filters
272    let dataframe = filters.into_iter().try_fold(dataframe, |df, expr| {
273        df.filter(expr).context(error::PlanSqlSnafu)
274    })?;
275
276    // Apply `like` predicate if exists
277    let dataframe = if let (ShowKind::Like(ident), Some(field)) = (&kind, like_field) {
278        dataframe
279            .filter(col(field).like(lit(ident.value.clone())))
280            .context(error::PlanSqlSnafu)?
281    } else {
282        dataframe
283    };
284
285    // Apply sorting
286    let dataframe = if sort.is_empty() {
287        dataframe
288    } else {
289        dataframe.sort(sort).context(error::PlanSqlSnafu)?
290    };
291
292    // Apply select
293    let dataframe = if select.is_empty() {
294        if projects.is_empty() {
295            dataframe
296        } else {
297            let projection = projects
298                .iter()
299                .map(|x| col(x.0).alias(x.1))
300                .collect::<Vec<_>>();
301            dataframe.select(projection).context(error::PlanSqlSnafu)?
302        }
303    } else {
304        dataframe.select(select).context(error::PlanSqlSnafu)?
305    };
306
307    // Apply projection
308    let dataframe = projects
309        .into_iter()
310        .try_fold(dataframe, |df, (column, renamed_column)| {
311            df.with_column_renamed(column, renamed_column)
312                .context(error::PlanSqlSnafu)
313        })?;
314
315    let dataframe = match kind {
316        ShowKind::All | ShowKind::Like(_) => {
317            // Like kind is processed above
318            dataframe
319        }
320        ShowKind::Where(filter) => {
321            // Cast the results into VIEW for `where` clause,
322            // which is evaluated against the column names displayed by the SHOW statement.
323            let view = dataframe.into_view();
324            let dataframe = SessionContext::new_with_state(
325                query_engine
326                    .engine_context(query_ctx.clone())
327                    .state()
328                    .clone(),
329            )
330            .read_table(view)?;
331
332            let planner = query_engine.planner();
333            let planner = planner
334                .as_any()
335                .downcast_ref::<DfLogicalPlanner>()
336                .expect("Must be the datafusion planner");
337
338            let filter = planner
339                .sql_to_expr(filter, dataframe.schema(), false, query_ctx)
340                .await?;
341
342            // Apply the `where` clause filters
343            dataframe.filter(filter).context(error::PlanSqlSnafu)?
344        }
345    };
346
347    let stream = dataframe.execute_stream().await?;
348
349    Ok(Output::new_with_stream(Box::pin(
350        RecordBatchStreamAdapter::try_new(stream).context(error::CreateRecordBatchSnafu)?,
351    )))
352}
353
354/// Execute `SHOW COLUMNS` statement.
355pub async fn show_columns(
356    stmt: ShowColumns,
357    query_engine: &QueryEngineRef,
358    catalog_manager: &CatalogManagerRef,
359    query_ctx: QueryContextRef,
360) -> Result<Output> {
361    let schema_name = if let Some(database) = stmt.database {
362        database
363    } else {
364        query_ctx.current_schema()
365    };
366
367    let projects = if stmt.full {
368        vec![
369            (columns::COLUMN_NAME, FIELD_COLUMN),
370            (columns::DATA_TYPE, COLUMN_TYPE_COLUMN),
371            (columns::COLLATION_NAME, COLUMN_COLLATION_COLUMN),
372            (columns::IS_NULLABLE, COLUMN_NULLABLE_COLUMN),
373            (columns::COLUMN_KEY, COLUMN_KEY_COLUMN),
374            (columns::COLUMN_DEFAULT, COLUMN_DEFAULT_COLUMN),
375            (columns::COLUMN_COMMENT, COLUMN_COMMENT_COLUMN),
376            (columns::PRIVILEGES, COLUMN_PRIVILEGES_COLUMN),
377            (columns::EXTRA, COLUMN_EXTRA_COLUMN),
378            (columns::GREPTIME_DATA_TYPE, COLUMN_GREPTIME_TYPE_COLUMN),
379        ]
380    } else {
381        vec![
382            (columns::COLUMN_NAME, FIELD_COLUMN),
383            (columns::DATA_TYPE, COLUMN_TYPE_COLUMN),
384            (columns::IS_NULLABLE, COLUMN_NULLABLE_COLUMN),
385            (columns::COLUMN_KEY, COLUMN_KEY_COLUMN),
386            (columns::COLUMN_DEFAULT, COLUMN_DEFAULT_COLUMN),
387            (columns::EXTRA, COLUMN_EXTRA_COLUMN),
388            (columns::GREPTIME_DATA_TYPE, COLUMN_GREPTIME_TYPE_COLUMN),
389        ]
390    };
391
392    let filters = vec![
393        col(columns::TABLE_NAME).eq(lit(&stmt.table)),
394        col(columns::TABLE_SCHEMA).eq(lit(schema_name.clone())),
395        col(columns::TABLE_CATALOG).eq(lit(query_ctx.current_catalog())),
396    ];
397    let like_field = Some(columns::COLUMN_NAME);
398    let sort = vec![col(columns::COLUMN_NAME).sort(true, true)];
399
400    query_from_information_schema_table(
401        query_engine,
402        catalog_manager,
403        query_ctx,
404        COLUMNS,
405        vec![],
406        projects,
407        filters,
408        like_field,
409        sort,
410        stmt.kind,
411    )
412    .await
413}
414
415/// Execute `SHOW INDEX` statement.
416pub async fn show_index(
417    stmt: ShowIndex,
418    query_engine: &QueryEngineRef,
419    catalog_manager: &CatalogManagerRef,
420    query_ctx: QueryContextRef,
421) -> Result<Output> {
422    let schema_name = if let Some(database) = stmt.database {
423        database
424    } else {
425        query_ctx.current_schema()
426    };
427
428    let select = vec![
429        col(statistics::TABLE_NAME).alias(INDEX_TABLE_COLUMN),
430        col(statistics::NON_UNIQUE).alias(INDEX_NONT_UNIQUE_COLUMN),
431        col(statistics::INDEX_NAME).alias(INDEX_KEY_NAME_COLUMN),
432        col(statistics::SEQ_IN_INDEX).alias(INDEX_SEQ_IN_INDEX_COLUMN),
433        col(statistics::COLUMN_NAME).alias(INDEX_COLUMN_NAME_COLUMN),
434        col(statistics::COLLATION).alias(COLUMN_COLLATION_COLUMN),
435        col(statistics::CARDINALITY).alias(INDEX_CARDINALITY_COLUMN),
436        col(statistics::SUB_PART).alias(INDEX_SUB_PART_COLUMN),
437        col(statistics::PACKED).alias(INDEX_PACKED_COLUMN),
438        col(statistics::NULLABLE).alias(COLUMN_NULLABLE_COLUMN),
439        col(statistics::INDEX_TYPE).alias(INDEX_INDEX_TYPE_COLUMN),
440        col(statistics::COMMENT).alias(COLUMN_COMMENT_COLUMN),
441        col(statistics::INDEX_COMMENT).alias(INDEX_COMMENT_COLUMN),
442        col(statistics::IS_VISIBLE).alias(INDEX_VISIBLE_COLUMN),
443        col(statistics::EXPRESSION).alias(INDEX_EXPRESSION_COLUMN),
444    ];
445
446    let projects = vec![
447        (statistics::TABLE_NAME, INDEX_TABLE_COLUMN),
448        (INDEX_NONT_UNIQUE_COLUMN, INDEX_NONT_UNIQUE_COLUMN),
449        (statistics::INDEX_NAME, INDEX_KEY_NAME_COLUMN),
450        (statistics::SEQ_IN_INDEX, INDEX_SEQ_IN_INDEX_COLUMN),
451        (statistics::COLUMN_NAME, INDEX_COLUMN_NAME_COLUMN),
452        (COLUMN_COLLATION_COLUMN, COLUMN_COLLATION_COLUMN),
453        (INDEX_CARDINALITY_COLUMN, INDEX_CARDINALITY_COLUMN),
454        (INDEX_SUB_PART_COLUMN, INDEX_SUB_PART_COLUMN),
455        (INDEX_PACKED_COLUMN, INDEX_PACKED_COLUMN),
456        (COLUMN_NULLABLE_COLUMN, COLUMN_NULLABLE_COLUMN),
457        (statistics::INDEX_TYPE, INDEX_INDEX_TYPE_COLUMN),
458        (COLUMN_COMMENT_COLUMN, COLUMN_COMMENT_COLUMN),
459        (INDEX_COMMENT_COLUMN, INDEX_COMMENT_COLUMN),
460        (INDEX_VISIBLE_COLUMN, INDEX_VISIBLE_COLUMN),
461        (INDEX_EXPRESSION_COLUMN, INDEX_EXPRESSION_COLUMN),
462    ];
463
464    let filters = vec![
465        col(statistics::TABLE_NAME).eq(lit(&stmt.table)),
466        col(statistics::TABLE_SCHEMA).eq(lit(schema_name.clone())),
467    ];
468    let like_field = None;
469    let sort = vec![
470        col(statistics::INDEX_NAME).sort(true, true),
471        col(statistics::SEQ_IN_INDEX).sort(true, true),
472    ];
473
474    query_from_information_schema_table(
475        query_engine,
476        catalog_manager,
477        query_ctx,
478        STATISTICS,
479        select,
480        projects,
481        filters,
482        like_field,
483        sort,
484        stmt.kind,
485    )
486    .await
487}
488
489/// Execute `SHOW REGION` statement.
490pub async fn show_region(
491    stmt: ShowRegion,
492    query_engine: &QueryEngineRef,
493    catalog_manager: &CatalogManagerRef,
494    query_ctx: QueryContextRef,
495) -> Result<Output> {
496    let schema_name = if let Some(database) = stmt.database {
497        database
498    } else {
499        query_ctx.current_schema()
500    };
501
502    let filters = vec![
503        col(region_peers::TABLE_NAME).eq(lit(&stmt.table)),
504        col(region_peers::TABLE_SCHEMA).eq(lit(schema_name.clone())),
505        col(region_peers::TABLE_CATALOG).eq(lit(query_ctx.current_catalog())),
506    ];
507    let projects = vec![
508        (region_peers::TABLE_NAME, "Table"),
509        (region_peers::REGION_ID, "Region"),
510        (region_peers::PEER_ID, "Peer"),
511        (region_peers::IS_LEADER, "Leader"),
512    ];
513
514    let like_field = None;
515    let sort = vec![
516        col(columns::REGION_ID).sort(true, true),
517        col(columns::PEER_ID).sort(true, true),
518    ];
519
520    query_from_information_schema_table(
521        query_engine,
522        catalog_manager,
523        query_ctx,
524        REGION_PEERS,
525        vec![],
526        projects,
527        filters,
528        like_field,
529        sort,
530        stmt.kind,
531    )
532    .await
533}
534
535/// Execute [`ShowTables`] statement and return the [`Output`] if success.
536pub async fn show_tables(
537    stmt: ShowTables,
538    query_engine: &QueryEngineRef,
539    catalog_manager: &CatalogManagerRef,
540    query_ctx: QueryContextRef,
541) -> Result<Output> {
542    let schema_name = if let Some(database) = stmt.database {
543        database
544    } else {
545        query_ctx.current_schema()
546    };
547
548    // MySQL renames `table_name` to `Tables_in_{schema}` for protocol compatibility
549    let tables_column = format!("Tables_in_{}", schema_name);
550    let projects = if stmt.full {
551        vec![
552            (tables::TABLE_NAME, tables_column.as_str()),
553            (tables::TABLE_TYPE, TABLE_TYPE_COLUMN),
554        ]
555    } else {
556        vec![(tables::TABLE_NAME, tables_column.as_str())]
557    };
558    let filters = vec![
559        col(tables::TABLE_SCHEMA).eq(lit(schema_name.clone())),
560        col(tables::TABLE_CATALOG).eq(lit(query_ctx.current_catalog())),
561    ];
562    let like_field = Some(tables::TABLE_NAME);
563    let sort = vec![col(tables::TABLE_NAME).sort(true, true)];
564
565    // Transform the WHERE clause for backward compatibility:
566    // Replace "Tables" with "Tables_in_{schema}" to support old queries
567    let kind = match stmt.kind {
568        ShowKind::Where(mut filter) => {
569            replace_column_in_expr(&mut filter, "Tables", &tables_column);
570            ShowKind::Where(filter)
571        }
572        other => other,
573    };
574
575    query_from_information_schema_table(
576        query_engine,
577        catalog_manager,
578        query_ctx,
579        TABLES,
580        vec![],
581        projects,
582        filters,
583        like_field,
584        sort,
585        kind,
586    )
587    .await
588}
589
590/// Execute [`ShowTableStatus`] statement and return the [`Output`] if success.
591pub async fn show_table_status(
592    stmt: ShowTableStatus,
593    query_engine: &QueryEngineRef,
594    catalog_manager: &CatalogManagerRef,
595    query_ctx: QueryContextRef,
596) -> Result<Output> {
597    let schema_name = if let Some(database) = stmt.database {
598        database
599    } else {
600        query_ctx.current_schema()
601    };
602
603    // Refer to https://dev.mysql.com/doc/refman/8.4/en/show-table-status.html
604    let projects = vec![
605        (tables::TABLE_NAME, "Name"),
606        (tables::ENGINE, "Engine"),
607        (tables::VERSION, "Version"),
608        (tables::ROW_FORMAT, "Row_format"),
609        (tables::TABLE_ROWS, "Rows"),
610        (tables::AVG_ROW_LENGTH, "Avg_row_length"),
611        (tables::DATA_LENGTH, "Data_length"),
612        (tables::MAX_DATA_LENGTH, "Max_data_length"),
613        (tables::INDEX_LENGTH, "Index_length"),
614        (tables::DATA_FREE, "Data_free"),
615        (tables::AUTO_INCREMENT, "Auto_increment"),
616        (tables::CREATE_TIME, "Create_time"),
617        (tables::UPDATE_TIME, "Update_time"),
618        (tables::CHECK_TIME, "Check_time"),
619        (tables::TABLE_COLLATION, "Collation"),
620        (tables::CHECKSUM, "Checksum"),
621        (tables::CREATE_OPTIONS, "Create_options"),
622        (tables::TABLE_COMMENT, "Comment"),
623    ];
624
625    let filters = vec![
626        col(tables::TABLE_SCHEMA).eq(lit(schema_name.clone())),
627        col(tables::TABLE_CATALOG).eq(lit(query_ctx.current_catalog())),
628    ];
629    let like_field = Some(tables::TABLE_NAME);
630    let sort = vec![col(tables::TABLE_NAME).sort(true, true)];
631
632    query_from_information_schema_table(
633        query_engine,
634        catalog_manager,
635        query_ctx,
636        TABLES,
637        vec![],
638        projects,
639        filters,
640        like_field,
641        sort,
642        stmt.kind,
643    )
644    .await
645}
646
647/// Execute `SHOW COLLATION` statement and returns the `Output` if success.
648pub async fn show_collations(
649    kind: ShowKind,
650    query_engine: &QueryEngineRef,
651    catalog_manager: &CatalogManagerRef,
652    query_ctx: QueryContextRef,
653) -> Result<Output> {
654    // Refer to https://dev.mysql.com/doc/refman/8.0/en/show-collation.html
655    let projects = vec![
656        ("collation_name", "Collation"),
657        ("character_set_name", "Charset"),
658        ("id", "Id"),
659        ("is_default", "Default"),
660        ("is_compiled", "Compiled"),
661        ("sortlen", "Sortlen"),
662    ];
663
664    let filters = vec![];
665    let like_field = Some("collation_name");
666    let sort = vec![];
667
668    query_from_information_schema_table(
669        query_engine,
670        catalog_manager,
671        query_ctx,
672        COLLATIONS,
673        vec![],
674        projects,
675        filters,
676        like_field,
677        sort,
678        kind,
679    )
680    .await
681}
682
683/// Execute `SHOW CHARSET` statement and returns the `Output` if success.
684pub async fn show_charsets(
685    kind: ShowKind,
686    query_engine: &QueryEngineRef,
687    catalog_manager: &CatalogManagerRef,
688    query_ctx: QueryContextRef,
689) -> Result<Output> {
690    // Refer to https://dev.mysql.com/doc/refman/8.0/en/show-character-set.html
691    let projects = vec![
692        ("character_set_name", "Charset"),
693        ("description", "Description"),
694        ("default_collate_name", "Default collation"),
695        ("maxlen", "Maxlen"),
696    ];
697
698    let filters = vec![];
699    let like_field = Some("character_set_name");
700    let sort = vec![];
701
702    query_from_information_schema_table(
703        query_engine,
704        catalog_manager,
705        query_ctx,
706        CHARACTER_SETS,
707        vec![],
708        projects,
709        filters,
710        like_field,
711        sort,
712        kind,
713    )
714    .await
715}
716
717pub fn show_variable(stmt: ShowVariables, query_ctx: QueryContextRef) -> Result<Output> {
718    let variable = stmt.variable.to_string().to_uppercase();
719    let value = match variable.as_str() {
720        "SYSTEM_TIME_ZONE" | "SYSTEM_TIMEZONE" => get_timezone(None).to_string(),
721        "TIME_ZONE" | "TIMEZONE" => query_ctx.timezone().to_string(),
722        "READ_PREFERENCE" => query_ctx.read_preference().to_string(),
723        "DATESTYLE" => {
724            let (style, order) = *query_ctx.configuration_parameter().pg_datetime_style();
725            format!("{}, {}", style, order)
726        }
727        "INTERVALSTYLE" => {
728            let style = *query_ctx
729                .configuration_parameter()
730                .pg_intervalstyle_format();
731            style.to_string()
732        }
733        "MAX_EXECUTION_TIME"
734            if query_ctx.channel() == Channel::Mysql => {
735                query_ctx.query_timeout_as_millis().to_string()
736            }
737        "STATEMENT_TIMEOUT"
738            // Add time units to postgres query timeout display.
739            if query_ctx.channel() == Channel::Postgres => {
740                let mut timeout = query_ctx.query_timeout_as_millis().to_string();
741                timeout.push_str("ms");
742                timeout
743            }
744        _ => return UnsupportedVariableSnafu { name: variable }.fail(),
745    };
746    let schema = Arc::new(Schema::new(vec![ColumnSchema::new(
747        variable,
748        ConcreteDataType::string_datatype(),
749        false,
750    )]));
751    let records = RecordBatches::try_from_columns(
752        schema,
753        vec![Arc::new(StringVector::from(vec![value])) as _],
754    )
755    .context(error::CreateRecordBatchSnafu)?;
756    Ok(Output::new_with_record_batches(records))
757}
758
759pub async fn show_status(_query_ctx: QueryContextRef) -> Result<Output> {
760    let schema = Arc::new(Schema::new(vec![
761        ColumnSchema::new("Variable_name", ConcreteDataType::string_datatype(), false),
762        ColumnSchema::new("Value", ConcreteDataType::string_datatype(), true),
763    ]));
764    let records = RecordBatches::try_from_columns(
765        schema,
766        vec![
767            Arc::new(StringVector::from(Vec::<&str>::new())) as _,
768            Arc::new(StringVector::from(Vec::<&str>::new())) as _,
769        ],
770    )
771    .context(error::CreateRecordBatchSnafu)?;
772    Ok(Output::new_with_record_batches(records))
773}
774
775pub async fn show_search_path(_query_ctx: QueryContextRef) -> Result<Output> {
776    let schema = Arc::new(Schema::new(vec![ColumnSchema::new(
777        "search_path",
778        ConcreteDataType::string_datatype(),
779        false,
780    )]));
781    let records = RecordBatches::try_from_columns(
782        schema,
783        vec![Arc::new(StringVector::from(vec![_query_ctx.current_schema()])) as _],
784    )
785    .context(error::CreateRecordBatchSnafu)?;
786    Ok(Output::new_with_record_batches(records))
787}
788
789pub fn show_create_database(database_name: &str, options: OptionMap) -> Result<Output> {
790    let stmt = CreateDatabase {
791        name: ObjectName::from(vec![Ident::new(database_name)]),
792        if_not_exists: true,
793        options,
794    };
795    let sql = format!("{stmt}");
796    let columns = vec![
797        Arc::new(StringVector::from(vec![database_name.to_string()])) as _,
798        Arc::new(StringVector::from(vec![sql])) as _,
799    ];
800    let records =
801        RecordBatches::try_from_columns(SHOW_CREATE_DATABASE_OUTPUT_SCHEMA.clone(), columns)
802            .context(error::CreateRecordBatchSnafu)?;
803    Ok(Output::new_with_record_batches(records))
804}
805
806pub fn show_create_table(
807    table_info: TableInfoRef,
808    schema_options: Option<SchemaOptions>,
809    partitions: Option<Partitions>,
810    query_ctx: QueryContextRef,
811) -> Result<Output> {
812    let table_name = table_info.name.clone();
813
814    let quote_style = query_ctx.quote_style();
815
816    let mut stmt = create_table_stmt(&table_info, schema_options, quote_style)?;
817    stmt.partitions = partitions.map(|mut p| {
818        p.set_quote(quote_style);
819        p
820    });
821    let sql = format!("{}", stmt);
822    let columns = vec![
823        Arc::new(StringVector::from(vec![table_name])) as _,
824        Arc::new(StringVector::from(vec![sql])) as _,
825    ];
826    let records = RecordBatches::try_from_columns(SHOW_CREATE_TABLE_OUTPUT_SCHEMA.clone(), columns)
827        .context(error::CreateRecordBatchSnafu)?;
828
829    Ok(Output::new_with_record_batches(records))
830}
831
832pub fn show_create_foreign_table_for_pg(
833    table: TableRef,
834    _query_ctx: QueryContextRef,
835) -> Result<Output> {
836    let table_info = table.table_info();
837
838    let table_meta = &table_info.meta;
839    let table_name = &table_info.name;
840    let schema = &table_info.meta.schema;
841    let is_metric_engine = is_metric_engine(&table_meta.engine);
842
843    let columns = schema
844        .column_schemas()
845        .iter()
846        .filter_map(|c| {
847            if is_metric_engine && is_metric_engine_internal_column(&c.name) {
848                None
849            } else {
850                Some(format!(
851                    "\"{}\" {}",
852                    c.name,
853                    c.data_type.postgres_datatype_name()
854                ))
855            }
856        })
857        .join(",\n  ");
858
859    let sql = format!(
860        r#"CREATE FOREIGN TABLE ft_{} (
861  {}
862)
863SERVER greptimedb
864OPTIONS (table_name '{}')"#,
865        table_name, columns, table_name
866    );
867
868    let columns = vec![
869        Arc::new(StringVector::from(vec![table_name.clone()])) as _,
870        Arc::new(StringVector::from(vec![sql])) as _,
871    ];
872    let records = RecordBatches::try_from_columns(SHOW_CREATE_TABLE_OUTPUT_SCHEMA.clone(), columns)
873        .context(error::CreateRecordBatchSnafu)?;
874
875    Ok(Output::new_with_record_batches(records))
876}
877
878pub fn show_create_view(
879    view_name: ObjectName,
880    definition: &str,
881    query_ctx: QueryContextRef,
882) -> Result<Output> {
883    let mut parser_ctx =
884        ParserContext::new(query_ctx.sql_dialect(), definition).context(error::SqlSnafu)?;
885
886    let Statement::CreateView(create_view) =
887        parser_ctx.parse_statement().context(error::SqlSnafu)?
888    else {
889        // MUST be `CreateView` statement.
890        unreachable!();
891    };
892
893    let stmt = CreateView {
894        name: view_name.clone(),
895        columns: create_view.columns,
896        query: create_view.query,
897        or_replace: create_view.or_replace,
898        if_not_exists: create_view.if_not_exists,
899    };
900
901    let sql = format!("{}", stmt);
902    let columns = vec![
903        Arc::new(StringVector::from(vec![view_name.to_string()])) as _,
904        Arc::new(StringVector::from(vec![sql])) as _,
905    ];
906    let records = RecordBatches::try_from_columns(SHOW_CREATE_VIEW_OUTPUT_SCHEMA.clone(), columns)
907        .context(error::CreateRecordBatchSnafu)?;
908
909    Ok(Output::new_with_record_batches(records))
910}
911
912/// Execute [`ShowViews`] statement and return the [`Output`] if success.
913pub async fn show_views(
914    stmt: ShowViews,
915    query_engine: &QueryEngineRef,
916    catalog_manager: &CatalogManagerRef,
917    query_ctx: QueryContextRef,
918) -> Result<Output> {
919    let schema_name = if let Some(database) = stmt.database {
920        database
921    } else {
922        query_ctx.current_schema()
923    };
924
925    let projects = vec![(tables::TABLE_NAME, VIEWS_COLUMN)];
926    let filters = vec![
927        col(tables::TABLE_SCHEMA).eq(lit(schema_name.clone())),
928        col(tables::TABLE_CATALOG).eq(lit(query_ctx.current_catalog())),
929    ];
930    let like_field = Some(tables::TABLE_NAME);
931    let sort = vec![col(tables::TABLE_NAME).sort(true, true)];
932
933    query_from_information_schema_table(
934        query_engine,
935        catalog_manager,
936        query_ctx,
937        VIEWS,
938        vec![],
939        projects,
940        filters,
941        like_field,
942        sort,
943        stmt.kind,
944    )
945    .await
946}
947
948/// Execute [`ShowFlows`] statement and return the [`Output`] if success.
949pub async fn show_flows(
950    stmt: ShowFlows,
951    query_engine: &QueryEngineRef,
952    catalog_manager: &CatalogManagerRef,
953    query_ctx: QueryContextRef,
954) -> Result<Output> {
955    let projects = vec![(flows::FLOW_NAME, FLOWS_COLUMN)];
956    let filters = vec![col(flows::TABLE_CATALOG).eq(lit(query_ctx.current_catalog()))];
957    let like_field = Some(flows::FLOW_NAME);
958    let sort = vec![col(flows::FLOW_NAME).sort(true, true)];
959
960    query_from_information_schema_table(
961        query_engine,
962        catalog_manager,
963        query_ctx,
964        FLOWS,
965        vec![],
966        projects,
967        filters,
968        like_field,
969        sort,
970        stmt.kind,
971    )
972    .await
973}
974
975/// Execute [`ShowFlowStatus`] statement and return the [`Output`] if success.
976pub async fn show_flow_status(
977    stmt: ShowFlowStatus,
978    query_engine: &QueryEngineRef,
979    catalog_manager: &CatalogManagerRef,
980    query_ctx: QueryContextRef,
981) -> Result<Output> {
982    let projects = vec![
983        (flow_statistics::FLOW_ID, flow_statistics::FLOW_ID),
984        (flow_statistics::FLOW_NAME, flow_statistics::FLOW_NAME),
985        (flow_statistics::START_TIME, flow_statistics::START_TIME),
986        (
987            flow_statistics::LAST_EXECUTION_TIME,
988            flow_statistics::LAST_EXECUTION_TIME,
989        ),
990        (
991            flow_statistics::UPTIME_SECONDS,
992            flow_statistics::UPTIME_SECONDS,
993        ),
994        (flow_statistics::STATE_SIZE, flow_statistics::STATE_SIZE),
995    ];
996    let like_field = Some(flow_statistics::FLOW_NAME);
997    let sort = vec![col(flow_statistics::FLOW_NAME).sort(true, true)];
998
999    query_from_information_schema_table(
1000        query_engine,
1001        catalog_manager,
1002        query_ctx,
1003        FLOW_STATISTICS,
1004        vec![],
1005        projects,
1006        vec![],
1007        like_field,
1008        sort,
1009        stmt.kind,
1010    )
1011    .await
1012}
1013
1014#[cfg(feature = "enterprise")]
1015pub async fn show_triggers(
1016    stmt: sql::statements::show::trigger::ShowTriggers,
1017    query_engine: &QueryEngineRef,
1018    catalog_manager: &CatalogManagerRef,
1019    query_ctx: QueryContextRef,
1020) -> Result<Output> {
1021    const TRIGGER_NAME: &str = "trigger_name";
1022    const TRIGGERS_COLUMN: &str = "Triggers";
1023
1024    let projects = vec![(TRIGGER_NAME, TRIGGERS_COLUMN)];
1025    let like_field = Some(TRIGGER_NAME);
1026    let sort = vec![col(TRIGGER_NAME).sort(true, true)];
1027
1028    query_from_information_schema_table(
1029        query_engine,
1030        catalog_manager,
1031        query_ctx,
1032        catalog::information_schema::TRIGGERS,
1033        vec![],
1034        projects,
1035        vec![],
1036        like_field,
1037        sort,
1038        stmt.kind,
1039    )
1040    .await
1041}
1042
1043pub fn show_create_flow(
1044    flow_name: ObjectName,
1045    flow_val: FlowInfoValue,
1046    query_ctx: QueryContextRef,
1047) -> Result<Output> {
1048    let mut parser_ctx =
1049        ParserContext::new(query_ctx.sql_dialect(), flow_val.raw_sql()).context(error::SqlSnafu)?;
1050
1051    let query = parser_ctx.parse_statement().context(error::SqlSnafu)?;
1052
1053    // since prom ql will parse `now()` to a fixed time, we need to not use it for generating raw query
1054    let raw_query = match &query {
1055        Statement::Tql(_) => flow_val.raw_sql().clone(),
1056        _ => query.to_string(),
1057    };
1058
1059    let query = Box::new(SqlOrTql::try_from_statement(query, &raw_query).context(error::SqlSnafu)?);
1060
1061    let comment = if flow_val.comment().is_empty() {
1062        None
1063    } else {
1064        Some(flow_val.comment().clone())
1065    };
1066
1067    let stmt = CreateFlow {
1068        flow_name,
1069        sink_table_name: ObjectName::from(vec![
1070            Ident::new(&flow_val.sink_table_name().schema_name),
1071            Ident::new(&flow_val.sink_table_name().table_name),
1072        ]),
1073        // notice we don't want `OR REPLACE` and `IF NOT EXISTS` in same sql since it's unclear what to do
1074        // so we set `or_replace` to false.
1075        or_replace: false,
1076        if_not_exists: true,
1077        expire_after: flow_val.expire_after(),
1078        eval_interval: flow_val.eval_interval(),
1079        comment,
1080        flow_options: OptionMap::from_filtered_string_map(
1081            flow_val.options(),
1082            &[FlowType::FLOW_TYPE_KEY],
1083        ),
1084        query,
1085    };
1086
1087    let sql = format!("{}", stmt);
1088    let columns = vec![
1089        Arc::new(StringVector::from(vec![flow_val.flow_name().clone()])) as _,
1090        Arc::new(StringVector::from(vec![sql])) as _,
1091    ];
1092    let records = RecordBatches::try_from_columns(SHOW_CREATE_FLOW_OUTPUT_SCHEMA.clone(), columns)
1093        .context(error::CreateRecordBatchSnafu)?;
1094
1095    Ok(Output::new_with_record_batches(records))
1096}
1097
1098pub fn describe_table(table: TableRef) -> Result<Output> {
1099    let table_info = table.table_info();
1100    let columns_schemas = table_info.meta.schema.column_schemas();
1101    let columns = vec![
1102        describe_column_names(columns_schemas),
1103        describe_column_types(columns_schemas),
1104        describe_column_keys(columns_schemas, &table_info.meta.primary_key_indices),
1105        describe_column_nullables(columns_schemas),
1106        describe_column_defaults(columns_schemas),
1107        describe_column_semantic_types(columns_schemas, &table_info.meta.primary_key_indices),
1108    ];
1109    let records = RecordBatches::try_from_columns(DESCRIBE_TABLE_OUTPUT_SCHEMA.clone(), columns)
1110        .context(error::CreateRecordBatchSnafu)?;
1111    Ok(Output::new_with_record_batches(records))
1112}
1113
1114fn describe_column_names(columns_schemas: &[ColumnSchema]) -> VectorRef {
1115    Arc::new(StringVector::from_iterator(
1116        columns_schemas.iter().map(|cs| cs.name.as_str()),
1117    ))
1118}
1119
1120fn describe_column_types(columns_schemas: &[ColumnSchema]) -> VectorRef {
1121    Arc::new(StringVector::from(
1122        columns_schemas
1123            .iter()
1124            .map(|cs| describe_column_type_name(&cs.data_type))
1125            .collect::<Vec<_>>(),
1126    ))
1127}
1128
1129fn describe_column_type_name(data_type: &ConcreteDataType) -> String {
1130    match data_type {
1131        ConcreteDataType::Json(json_type) if json_type.is_json2() => "Json2".to_string(),
1132        data_type => data_type.name(),
1133    }
1134}
1135
1136fn describe_column_keys(
1137    columns_schemas: &[ColumnSchema],
1138    primary_key_indices: &[usize],
1139) -> VectorRef {
1140    Arc::new(StringVector::from_iterator(
1141        columns_schemas.iter().enumerate().map(|(i, cs)| {
1142            if cs.is_time_index() || primary_key_indices.contains(&i) {
1143                PRI_KEY
1144            } else {
1145                ""
1146            }
1147        }),
1148    ))
1149}
1150
1151fn describe_column_nullables(columns_schemas: &[ColumnSchema]) -> VectorRef {
1152    Arc::new(StringVector::from_iterator(columns_schemas.iter().map(
1153        |cs| {
1154            if cs.is_nullable() { YES_STR } else { NO_STR }
1155        },
1156    )))
1157}
1158
1159fn describe_column_defaults(columns_schemas: &[ColumnSchema]) -> VectorRef {
1160    Arc::new(StringVector::from(
1161        columns_schemas
1162            .iter()
1163            .map(|cs| {
1164                cs.default_constraint()
1165                    .map_or(String::from(""), |dc| dc.to_string())
1166            })
1167            .collect::<Vec<String>>(),
1168    ))
1169}
1170
1171fn describe_column_semantic_types(
1172    columns_schemas: &[ColumnSchema],
1173    primary_key_indices: &[usize],
1174) -> VectorRef {
1175    Arc::new(StringVector::from_iterator(
1176        columns_schemas.iter().enumerate().map(|(i, cs)| {
1177            if primary_key_indices.contains(&i) {
1178                SEMANTIC_TYPE_PRIMARY_KEY
1179            } else if cs.is_time_index() {
1180                SEMANTIC_TYPE_TIME_INDEX
1181            } else {
1182                SEMANTIC_TYPE_FIELD
1183            }
1184        }),
1185    ))
1186}
1187
1188// lists files in the frontend to reduce unnecessary scan requests repeated in each datanode.
1189pub async fn prepare_file_table_files(
1190    options: &HashMap<String, String>,
1191    local_file_access: &LocalFileAccess,
1192) -> Result<(ObjectStore, Vec<String>)> {
1193    let url = options
1194        .get(FILE_TABLE_LOCATION_KEY)
1195        .context(error::MissingRequiredFieldSnafu {
1196            name: FILE_TABLE_LOCATION_KEY,
1197        })?;
1198
1199    let regex = options
1200        .get(FILE_TABLE_PATTERN_KEY)
1201        .map(|x| Regex::new(x))
1202        .transpose()
1203        .context(error::BuildRegexSnafu)?;
1204    let backend = build_backend_with_path(url, options, local_file_access)
1205        .await
1206        .context(error::BuildBackendSnafu)?;
1207    let source = if let Some(filename) = backend.object_path {
1208        Source::Filename(filename)
1209    } else {
1210        Source::Dir
1211    };
1212    let lister = Lister::new(backend.object_store.clone(), source, url.clone(), regex);
1213    // If we scan files in a directory every time the database restarts,
1214    // then it might lead to a potential undefined behavior:
1215    // If a user adds a file with an incompatible schema to that directory,
1216    // it will make the external table unavailable.
1217    let files = lister
1218        .list()
1219        .await
1220        .context(error::ListObjectsSnafu)?
1221        .into_iter()
1222        .filter_map(|entry| {
1223            if entry.path().ends_with('/') {
1224                None
1225            } else {
1226                Some(entry.path().to_string())
1227            }
1228        })
1229        .collect::<Vec<_>>();
1230    Ok((backend.object_store, files))
1231}
1232
1233pub async fn infer_file_table_schema(
1234    object_store: &ObjectStore,
1235    files: &[String],
1236    options: &HashMap<String, String>,
1237) -> Result<Schema> {
1238    let format = parse_file_table_format(options)?;
1239    let merged = infer_schemas(object_store, files, format.as_ref())
1240        .await
1241        .context(error::InferSchemaSnafu)?;
1242    Schema::try_from(merged).context(error::ConvertSchemaSnafu)
1243}
1244
1245// Converts the file column schemas to table column schemas.
1246// Returns the column schemas and the time index column name.
1247//
1248// More specifically, this function will do the following:
1249// 1. Add a default time index column if there is no time index column
1250//    in the file column schemas, or
1251// 2. If the file column schemas contain a column with name conflicts with
1252//    the default time index column, it will replace the column schema
1253//    with the default one.
1254pub fn file_column_schemas_to_table(
1255    file_column_schemas: &[ColumnSchema],
1256) -> (Vec<ColumnSchema>, String) {
1257    let mut column_schemas = file_column_schemas.to_owned();
1258    if let Some(time_index_column) = column_schemas.iter().find(|c| c.is_time_index()) {
1259        let time_index = time_index_column.name.clone();
1260        return (column_schemas, time_index);
1261    }
1262
1263    let timestamp_type = ConcreteDataType::timestamp_millisecond_datatype();
1264    let default_zero = Value::Timestamp(Timestamp::new_millisecond(0));
1265    let timestamp_column_schema = ColumnSchema::new(greptime_timestamp(), timestamp_type, false)
1266        .with_time_index(true)
1267        .with_default_constraint(Some(ColumnDefaultConstraint::Value(default_zero)))
1268        .unwrap();
1269
1270    if let Some(column_schema) = column_schemas
1271        .iter_mut()
1272        .find(|column_schema| column_schema.name == greptime_timestamp())
1273    {
1274        // Replace the column schema with the default one
1275        *column_schema = timestamp_column_schema;
1276    } else {
1277        column_schemas.push(timestamp_column_schema);
1278    }
1279
1280    (column_schemas, greptime_timestamp().to_string())
1281}
1282
1283/// This function checks if the column schemas from a file can be matched with
1284/// the column schemas of a table.
1285///
1286/// More specifically, for each column seen in the table schema,
1287/// - If the same column does exist in the file schema, it checks if the data
1288///   type of the file column can be casted into the form of the table column.
1289/// - If the same column does not exist in the file schema, it checks if the
1290///   table column is nullable or has a default constraint.
1291pub fn check_file_to_table_schema_compatibility(
1292    file_column_schemas: &[ColumnSchema],
1293    table_column_schemas: &[ColumnSchema],
1294) -> Result<()> {
1295    let file_schemas_map = file_column_schemas
1296        .iter()
1297        .map(|s| (s.name.clone(), s))
1298        .collect::<HashMap<_, _>>();
1299
1300    for table_column in table_column_schemas {
1301        if let Some(file_column) = file_schemas_map.get(&table_column.name) {
1302            // TODO(zhongzc): a temporary solution, we should use `can_cast_to` once it's ready.
1303            ensure!(
1304                file_column
1305                    .data_type
1306                    .can_arrow_type_cast_to(&table_column.data_type),
1307                error::ColumnSchemaIncompatibleSnafu {
1308                    column: table_column.name.clone(),
1309                    file_type: file_column.data_type.clone(),
1310                    table_type: table_column.data_type.clone(),
1311                }
1312            );
1313        } else {
1314            ensure!(
1315                table_column.is_nullable() || table_column.default_constraint().is_some(),
1316                error::ColumnSchemaNoDefaultSnafu {
1317                    column: table_column.name.clone(),
1318                }
1319            );
1320        }
1321    }
1322
1323    Ok(())
1324}
1325
1326fn parse_file_table_format(options: &HashMap<String, String>) -> Result<Box<dyn FileFormat>> {
1327    Ok(
1328        match Format::try_from(options).context(error::ParseFileFormatSnafu)? {
1329            Format::Csv(format) => Box::new(format),
1330            Format::Json(format) => Box::new(format),
1331            Format::Parquet(format) => Box::new(format),
1332            Format::Orc(format) => Box::new(format),
1333        },
1334    )
1335}
1336
1337pub async fn show_processlist(
1338    stmt: ShowProcessList,
1339    query_engine: &QueryEngineRef,
1340    catalog_manager: &CatalogManagerRef,
1341    query_ctx: QueryContextRef,
1342) -> Result<Output> {
1343    let projects = if stmt.full {
1344        vec![
1345            (process_list::ID, "Id"),
1346            (process_list::CATALOG, "Catalog"),
1347            (process_list::SCHEMAS, "Schema"),
1348            (process_list::CLIENT, "Client"),
1349            (process_list::FRONTEND, "Frontend"),
1350            (process_list::START_TIMESTAMP, "Start Time"),
1351            (process_list::ELAPSED_TIME, "Elapsed Time"),
1352            (process_list::QUERY, "Query"),
1353        ]
1354    } else {
1355        vec![
1356            (process_list::ID, "Id"),
1357            (process_list::CATALOG, "Catalog"),
1358            (process_list::QUERY, "Query"),
1359            (process_list::ELAPSED_TIME, "Elapsed Time"),
1360        ]
1361    };
1362
1363    let filters = vec![];
1364    let like_field = None;
1365    let sort = vec![col("id").sort(true, true)];
1366    query_from_information_schema_table(
1367        query_engine,
1368        catalog_manager,
1369        query_ctx.clone(),
1370        "process_list",
1371        vec![],
1372        projects.clone(),
1373        filters,
1374        like_field,
1375        sort,
1376        ShowKind::All,
1377    )
1378    .await
1379}
1380
1381#[cfg(test)]
1382mod test {
1383    use std::sync::Arc;
1384
1385    use common_query::{Output, OutputData};
1386    use common_recordbatch::{RecordBatch, RecordBatches};
1387    use common_time::Timezone;
1388    use common_time::timestamp::TimeUnit;
1389    use datatypes::prelude::ConcreteDataType;
1390    use datatypes::schema::{ColumnDefaultConstraint, ColumnSchema, Schema, SchemaRef};
1391    use datatypes::types::json_type::JsonNativeType;
1392    use datatypes::vectors::{StringVector, TimestampMillisecondVector, UInt32Vector, VectorRef};
1393    use session::context::QueryContextBuilder;
1394    use snafu::ResultExt;
1395    use sql::ast::{Ident, ObjectName};
1396    use sql::statements::show::ShowVariables;
1397    use table::TableRef;
1398    use table::test_util::MemTable;
1399
1400    use super::{describe_column_type_name, show_variable};
1401    use crate::error;
1402    use crate::error::Result;
1403    use crate::sql::{
1404        DESCRIBE_TABLE_OUTPUT_SCHEMA, NO_STR, SEMANTIC_TYPE_FIELD, SEMANTIC_TYPE_TIME_INDEX,
1405        YES_STR, describe_table,
1406    };
1407
1408    #[test]
1409    fn test_describe_table_multiple_columns() -> Result<()> {
1410        let table_name = "test_table";
1411        let schema = vec![
1412            ColumnSchema::new("t1", ConcreteDataType::uint32_datatype(), true),
1413            ColumnSchema::new(
1414                "t2",
1415                ConcreteDataType::timestamp_datatype(TimeUnit::Millisecond),
1416                false,
1417            )
1418            .with_default_constraint(Some(ColumnDefaultConstraint::Function(String::from(
1419                "current_timestamp()",
1420            ))))
1421            .unwrap()
1422            .with_time_index(true),
1423        ];
1424        let data = vec![
1425            Arc::new(UInt32Vector::from_slice([0])) as _,
1426            Arc::new(TimestampMillisecondVector::from_slice([0])) as _,
1427        ];
1428        let expected_columns = vec![
1429            Arc::new(StringVector::from(vec!["t1", "t2"])) as _,
1430            Arc::new(StringVector::from(vec!["UInt32", "TimestampMillisecond"])) as _,
1431            Arc::new(StringVector::from(vec!["", "PRI"])) as _,
1432            Arc::new(StringVector::from(vec![YES_STR, NO_STR])) as _,
1433            Arc::new(StringVector::from(vec!["", "current_timestamp()"])) as _,
1434            Arc::new(StringVector::from(vec![
1435                SEMANTIC_TYPE_FIELD,
1436                SEMANTIC_TYPE_TIME_INDEX,
1437            ])) as _,
1438        ];
1439
1440        describe_table_test_by_schema(table_name, schema, data, expected_columns)
1441    }
1442
1443    #[test]
1444    fn test_describe_column_type_name_json2() {
1445        assert_eq!(
1446            describe_column_type_name(&ConcreteDataType::json2(JsonNativeType::Null)),
1447            "Json2"
1448        );
1449        assert_eq!(
1450            describe_column_type_name(&ConcreteDataType::uint32_datatype()),
1451            "UInt32"
1452        );
1453    }
1454
1455    fn describe_table_test_by_schema(
1456        table_name: &str,
1457        schema: Vec<ColumnSchema>,
1458        data: Vec<VectorRef>,
1459        expected_columns: Vec<VectorRef>,
1460    ) -> Result<()> {
1461        let table_schema = SchemaRef::new(Schema::new(schema));
1462        let table = prepare_describe_table(table_name, table_schema, data);
1463
1464        let expected =
1465            RecordBatches::try_from_columns(DESCRIBE_TABLE_OUTPUT_SCHEMA.clone(), expected_columns)
1466                .context(error::CreateRecordBatchSnafu)?;
1467
1468        if let OutputData::RecordBatches(res) = describe_table(table)?.data {
1469            assert_eq!(res.take(), expected.take());
1470        } else {
1471            panic!("describe table must return record batch");
1472        }
1473
1474        Ok(())
1475    }
1476
1477    fn prepare_describe_table(
1478        table_name: &str,
1479        table_schema: SchemaRef,
1480        data: Vec<VectorRef>,
1481    ) -> TableRef {
1482        let record_batch = RecordBatch::new(table_schema, data).unwrap();
1483        MemTable::table(table_name, record_batch)
1484    }
1485
1486    #[test]
1487    fn test_show_variable() {
1488        assert_eq!(
1489            exec_show_variable("SYSTEM_TIME_ZONE", "Asia/Shanghai").unwrap(),
1490            "UTC"
1491        );
1492        assert_eq!(
1493            exec_show_variable("SYSTEM_TIMEZONE", "Asia/Shanghai").unwrap(),
1494            "UTC"
1495        );
1496        assert_eq!(
1497            exec_show_variable("TIME_ZONE", "Asia/Shanghai").unwrap(),
1498            "Asia/Shanghai"
1499        );
1500        assert_eq!(
1501            exec_show_variable("TIMEZONE", "Asia/Shanghai").unwrap(),
1502            "Asia/Shanghai"
1503        );
1504        assert!(exec_show_variable("TIME ZONE", "Asia/Shanghai").is_err());
1505        assert!(exec_show_variable("SYSTEM TIME ZONE", "Asia/Shanghai").is_err());
1506    }
1507
1508    fn exec_show_variable(variable: &str, tz: &str) -> Result<String> {
1509        let stmt = ShowVariables {
1510            variable: ObjectName::from(vec![Ident::new(variable)]),
1511        };
1512        let ctx = Arc::new(
1513            QueryContextBuilder::default()
1514                .timezone(Timezone::from_tz_string(tz).unwrap())
1515                .build(),
1516        );
1517        match show_variable(stmt, ctx) {
1518            Ok(Output {
1519                data: OutputData::RecordBatches(record),
1520                ..
1521            }) => {
1522                let record = record.take().first().cloned().unwrap();
1523                Ok(record.iter_column_as_string(0).next().unwrap().unwrap())
1524            }
1525            Ok(_) => unreachable!(),
1526            Err(e) => Err(e),
1527        }
1528    }
1529}