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_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
100/// SHOW index columns
101const 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
188/// Builds the [`DataFrame`] for `SHOW DATABASES` without executing it.
189pub 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
223/// Replaces column identifier references in a SQL expression.
224/// Used for backward compatibility where old column names should work with new ones.
225fn 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/// Cast a `show` statement execution into a query from tables in  `information_schema`.
247/// - `table_name`: the table name in `information_schema`,
248/// - `projects`: query projection, a list of `(column, renamed_column)`,
249/// - `filters`: filter expressions for query,
250/// - `like_field`: the field to filter by the predicate `ShowKind::Like`,
251/// - `sort`: sort the results by the specified sorting expressions,
252/// - `kind`: the show kind
253#[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/// Builds the [`DataFrame`] for a `SHOW` statement without executing it,
290/// so `Describe` handlers can derive the output schema from the same
291/// projection the executor uses.
292#[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    // Apply filters
325    let dataframe = filters.into_iter().try_fold(dataframe, |df, expr| {
326        df.filter(expr).context(error::PlanSqlSnafu)
327    })?;
328
329    // Apply `like` predicate if exists
330    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    // Apply sorting
339    let dataframe = if sort.is_empty() {
340        dataframe
341    } else {
342        dataframe.sort(sort).context(error::PlanSqlSnafu)?
343    };
344
345    // Apply select
346    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    // Apply projection
361    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            // Like kind is processed above
371            dataframe
372        }
373        ShowKind::Where(filter) => {
374            // Cast the results into VIEW for `where` clause,
375            // which is evaluated against the column names displayed by the SHOW statement.
376            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            // Apply the `where` clause filters
396            dataframe.filter(filter).context(error::PlanSqlSnafu)?
397        }
398    };
399
400    Ok(dataframe)
401}
402
403/// Execute `SHOW COLUMNS` statement.
404pub 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
414/// Builds the [`DataFrame`] for `SHOW COLUMNS` without executing it.
415pub 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
475/// Execute `SHOW INDEX` statement.
476pub 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
486/// Builds the [`DataFrame`] for `SHOW INDEX` without executing it.
487pub 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
560/// Execute `SHOW REGION` statement.
561pub 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
571/// Builds the [`DataFrame`] for `SHOW REGION` without executing it.
572pub 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
617/// Execute [`ShowTables`] statement and return the [`Output`] if success.
618pub 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
628/// Builds the [`DataFrame`] for [`ShowTables`] without executing it.
629pub 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    // MySQL renames `table_name` to `Tables_in_{schema}` for protocol compatibility
642    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    // Transform the WHERE clause for backward compatibility:
659    // Replace "Tables" with "Tables_in_{schema}" to support old queries
660    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
684/// Execute [`ShowTableStatus`] statement and return the [`Output`] if success.
685pub 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
696/// Builds the [`DataFrame`] for [`ShowTableStatus`] without executing it.
697pub 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    // Refer to https://dev.mysql.com/doc/refman/8.4/en/show-table-status.html
710    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
753/// Execute `SHOW COLLATION` statement and returns the `Output` if success.
754pub 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
765/// Builds the [`DataFrame`] for `SHOW COLLATION` without executing it.
766pub async fn show_collations_dataframe(
767    kind: &ShowKind,
768    query_engine: &QueryEngineRef,
769    catalog_manager: &CatalogManagerRef,
770    query_ctx: QueryContextRef,
771) -> Result<DataFrame> {
772    // Refer to https://dev.mysql.com/doc/refman/8.0/en/show-collation.html
773    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
801/// Execute `SHOW CHARSET` statement and returns the `Output` if success.
802pub 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
813/// Builds the [`DataFrame`] for `SHOW CHARSET` without executing it.
814pub async fn show_charsets_dataframe(
815    kind: &ShowKind,
816    query_engine: &QueryEngineRef,
817    catalog_manager: &CatalogManagerRef,
818    query_ctx: QueryContextRef,
819) -> Result<DataFrame> {
820    // Refer to https://dev.mysql.com/doc/refman/8.0/en/show-character-set.html
821    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            // Add time units to postgres query timeout display.
869            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        // MUST be `CreateView` statement.
1020        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
1042/// Execute [`ShowViews`] statement and return the [`Output`] if success.
1043pub 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
1053/// Builds the [`DataFrame`] for [`ShowViews`] without executing it.
1054pub 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
1089/// Execute [`ShowFlows`] statement and return the [`Output`] if success.
1090pub 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
1100/// Builds the [`DataFrame`] for [`ShowFlows`] without executing it.
1101pub 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
1127/// Execute [`ShowFlowStatus`] statement and return the [`Output`] if success.
1128pub 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    // since prom ql will parse `now()` to a fixed time, we need to not use it for generating raw query
1206    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        // notice we don't want `OR REPLACE` and `IF NOT EXISTS` in same sql since it's unclear what to do
1226        // so we set `or_replace` to false.
1227        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
1345// lists files in the frontend to reduce unnecessary scan requests repeated in each datanode.
1346pub 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    // If we scan files in a directory every time the database restarts,
1371    // then it might lead to a potential undefined behavior:
1372    // If a user adds a file with an incompatible schema to that directory,
1373    // it will make the external table unavailable.
1374    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
1402// Converts the file column schemas to table column schemas.
1403// Returns the column schemas and the time index column name.
1404//
1405// More specifically, this function will do the following:
1406// 1. Add a default time index column if there is no time index column
1407//    in the file column schemas, or
1408// 2. If the file column schemas contain a column with name conflicts with
1409//    the default time index column, it will replace the column schema
1410//    with the default one.
1411pub 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        // Replace the column schema with the default one
1432        *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
1440/// This function checks if the column schemas from a file can be matched with
1441/// the column schemas of a table.
1442///
1443/// More specifically, for each column seen in the table schema,
1444/// - If the same column does exist in the file schema, it checks if the data
1445///   type of the file column can be casted into the form of the table column.
1446/// - If the same column does not exist in the file schema, it checks if the
1447///   table column is nullable or has a default constraint.
1448pub 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            // TODO(zhongzc): a temporary solution, we should use `can_cast_to` once it's ready.
1460            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
1505/// Builds the [`DataFrame`] for `SHOW PROCESSLIST` without executing it.
1506pub 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}