From 2527bc09e714d5ba4144d2d227fb5e882b164356 Mon Sep 17 00:00:00 2001 From: Ning Sun Date: Sun, 30 Aug 2026 13:24:17 +0000 Subject: [PATCH] fix: postgres describe for more statements (#8974) * fix: postgres describe for more statements Signed-off-by: Ning Sun * fix: cover more show statements Signed-off-by: Ning Sun * fix: address review comments - add missing `clippy::too_many_arguments` allow on `query_from_information_schema_dataframe` (CI clippy failure) - take `&ShowKind` in the information-schema dataframe helper so `kind` is no longer cloned at every call site; only the WHERE arm (which needs an owned expression for `sql_to_expr`) clones internally - document why re-applying TQL explain formats never overwrites an existing value (per-query context state) Signed-off-by: Ning Sun * chore: trim comments to essentials Signed-off-by: Ning Sun --------- Signed-off-by: Ning Sun (cherry picked from commit d32cd77505d3cb304d66796b515c45ab42e739c6) --- src/common/recordbatch/src/cursor.rs | 1 - src/frontend/src/instance.rs | 101 ++++++++ src/frontend/src/lib.rs | 2 + src/operator/src/statement.rs | 2 +- src/operator/src/statement/admin.rs | 171 ++++++++----- src/query/src/analyze.rs | 10 +- src/query/src/sql.rs | 249 +++++++++++++++---- src/servers/src/postgres/handler.rs | 83 ++++--- src/session/src/lib.rs | 5 +- tests-integration/tests/sql.rs | 343 +++++++++++++++++++++++++++ 10 files changed, 825 insertions(+), 142 deletions(-) diff --git a/src/common/recordbatch/src/cursor.rs b/src/common/recordbatch/src/cursor.rs index 60ae9d54fa..125a68f4f1 100644 --- a/src/common/recordbatch/src/cursor.rs +++ b/src/common/recordbatch/src/cursor.rs @@ -47,7 +47,6 @@ impl RecordBatchStreamCursor { } } - /// Returns the schema of the underlying record batch stream. pub fn schema(&self) -> SchemaRef { self.schema.clone() } diff --git a/src/frontend/src/instance.rs b/src/frontend/src/instance.rs index 6a653b2b1d..2319f1e883 100644 --- a/src/frontend/src/instance.rs +++ b/src/frontend/src/instance.rs @@ -59,6 +59,7 @@ use common_recordbatch::error::StreamTimeoutSnafu; use common_telemetry::logging::SlowQueryOptions; use common_telemetry::{debug, error, tracing}; use dashmap::DashMap; +use datafusion::dataframe::DataFrame; use datafusion::physical_plan::ExecutionPlan; use datafusion_expr::LogicalPlan; use futures::{Stream, StreamExt, future}; @@ -847,6 +848,22 @@ impl Instance { query_interceptor.pre_execute(stmt.as_ref(), Some(&plan), query_ctx.clone())?; + // TQL EXPLAIN/ANALYZE formats are consumed from the query context at + // execution time (see `optimize_physical_plan`); re-apply the side + // effect of `plan_tql` that was lost when the plan was built during + // Describe. `explain_format` is per-query state, so this never + // overwrites anything. + if let Some(Statement::Tql(tql)) = &stmt { + let format = match tql { + Tql::Explain(explain) => explain.format.as_ref(), + Tql::Analyze(analyze) => analyze.format.as_ref(), + Tql::Eval(_) => None, + }; + if let Some(format) = format { + query_ctx.set_explain_format(format.to_string()); + } + } + let query = stmt .as_ref() .map(|s| s.to_string()) @@ -928,6 +945,59 @@ impl Instance { vec![result] } + /// Builds the [`DataFrame`] for an information-schema-backed `SHOW` + /// statement; `None` for other statements. The future is boxed to keep + /// `do_describe_inner`'s state machine small. + fn show_statement_dataframe<'a>( + &'a self, + stmt: &'a Statement, + query_ctx: &'a QueryContextRef, + ) -> Pin>> + Send + 'a>> { + Box::pin(async move { + let engine = &self.query_engine; + let catalog_manager = self.catalog_manager(); + let ctx = query_ctx.clone(); + let dataframe = match stmt { + Statement::ShowDatabases(show) => { + query::sql::show_databases_dataframe(show, engine, catalog_manager, ctx).await + } + Statement::ShowTables(show) => { + query::sql::show_tables_dataframe(show, engine, catalog_manager, ctx).await + } + Statement::ShowViews(show) => { + query::sql::show_views_dataframe(show, engine, catalog_manager, ctx).await + } + Statement::ShowFlows(show) => { + query::sql::show_flows_dataframe(show, engine, catalog_manager, ctx).await + } + Statement::ShowColumns(show) => { + query::sql::show_columns_dataframe(show, engine, catalog_manager, ctx).await + } + Statement::ShowTableStatus(show) => { + query::sql::show_table_status_dataframe(show, engine, catalog_manager, ctx) + .await + } + Statement::ShowCharset(kind) => { + query::sql::show_charsets_dataframe(kind, engine, catalog_manager, ctx).await + } + Statement::ShowCollation(kind) => { + query::sql::show_collations_dataframe(kind, engine, catalog_manager, ctx).await + } + Statement::ShowIndex(show) => { + query::sql::show_index_dataframe(show, engine, catalog_manager, ctx).await + } + Statement::ShowRegion(show) => { + query::sql::show_region_dataframe(show, engine, catalog_manager, ctx).await + } + Statement::ShowProcesslist(show) => { + query::sql::show_processlist_dataframe(show, engine, catalog_manager, ctx).await + } + _ => return None, + }; + Some(dataframe) + }) + } + async fn do_describe_inner( &self, stmt: Statement, @@ -947,6 +1017,37 @@ impl Instance { let plannable = is_inner_plannable(&stmt) || matches!(&stmt, Statement::Explain(explain) if is_inner_plannable(explain.statement.as_ref())); + if let Statement::Tql(tql) = stmt { + // TQL produces a logical plan; describe it from the plan so the + // extended-protocol RowDescription matches the executed DataRows. + self.check_sql_permission(&Statement::Tql(tql.clone()), &query_ctx) + .await?; + let plan = self.statement_executor.plan_tql(tql, &query_ctx).await?; + return self + .query_engine + .describe(plan, query_ctx) + .await + .map(Some) + .context(error::DescribeStatementSnafu); + } + + // Describe SHOW statements from the same projection the executor builds. + if let Some(dataframe) = self + .show_statement_dataframe(&stmt, &query_ctx) + .await + .transpose() + .context(PlanStatementSnafu)? + { + self.check_sql_permission(&stmt, &query_ctx).await?; + let plan = dataframe.into_unoptimized_plan(); + return self + .query_engine + .describe(plan, query_ctx) + .await + .map(Some) + .context(error::DescribeStatementSnafu); + } + if plannable { self.check_sql_permission(&stmt, &query_ctx).await?; diff --git a/src/frontend/src/lib.rs b/src/frontend/src/lib.rs index c170236073..1468fe7c0b 100644 --- a/src/frontend/src/lib.rs +++ b/src/frontend/src/lib.rs @@ -12,6 +12,8 @@ // See the License for the specific language governing permissions and // limitations under the License. +#![recursion_limit = "256"] + pub mod error; pub mod events; pub mod frontend; diff --git a/src/operator/src/statement.rs b/src/operator/src/statement.rs index 1143b1c68c..9be041435d 100644 --- a/src/operator/src/statement.rs +++ b/src/operator/src/statement.rs @@ -83,7 +83,7 @@ use table::table_reference::TableReference; pub use self::admin::{ AdminEventRecorderHandle, AdminFunctionLayer, AdminFunctionLayerRef, AdminFunctionRecordingLayer, AdminFunctionRequest, AdminFunctionResponse, AdminFunctionService, - AdminFunctionServiceRef, + AdminFunctionServiceRef, admin_output_schema, }; use self::set::{ set_bytea_output, set_datestyle, set_intervalstyle, set_timezone, validate_client_encoding, diff --git a/src/operator/src/statement/admin.rs b/src/operator/src/statement/admin.rs index 2b4a909e3c..ff8b26915c 100644 --- a/src/operator/src/statement/admin.rs +++ b/src/operator/src/statement/admin.rs @@ -19,6 +19,7 @@ use std::sync::Arc; use common_function::function::FunctionContext; use common_function::function_registry::{FUNCTION_REGISTRY, get_admin_function}; +use common_function::state::FunctionState; use common_query::Output; use common_recordbatch::{RecordBatch, RecordBatches}; use common_sql::convert::sql_value_to_value; @@ -77,6 +78,103 @@ struct CoreAdminFunctionService { query_engine: query::QueryEngineRef, } +/// Parts of an `ADMIN` call needed both for execution and schema derivation. +struct ResolvedAdminFunction { + admin_udf: datafusion_expr::ScalarUDF, + fn_name: String, + args: Vec, + arg_types: Vec, + ret_type: ArrowDataType, +} + +/// Resolves the function, parses its literal arguments and derives its +/// return type, without executing it. +fn resolve_admin_function( + stmt: &Admin, + query_ctx: &QueryContextRef, + state: Arc, +) -> Result { + let Admin::Func(func) = stmt; + // the function name should be in lower case. + let func_name = func.name.to_string().to_lowercase(); + let factory = get_admin_function(&func_name) + .or_else(|| FUNCTION_REGISTRY.get_function(&func_name)) + .context(error::AdminFunctionNotFoundSnafu { + name: func_name.clone(), + })?; + + let func_ctx = FunctionContext { + query_ctx: query_ctx.clone(), + state, + }; + + let admin_udf = factory.provide(func_ctx); + admin_udf + .as_async() + .context(error::AdminFunctionNotFoundSnafu { name: func_name })?; + + let fn_name = admin_udf.name().to_string(); + let signature = admin_udf.signature(); + + // Parse function arguments + let FunctionArguments::List(args) = &func.args else { + return error::BuildAdminFunctionArgsSnafu { + msg: format!("unsupported function args {} for {}", func.args, fn_name), + } + .fail(); + }; + let arg_values = args + .args + .iter() + .map(|arg| { + let FunctionArg::Unnamed(FunctionArgExpr::Expr(Expr::Value(value))) = arg else { + return error::BuildAdminFunctionArgsSnafu { + msg: format!("unsupported function arg {arg} for {}", fn_name), + } + .fail(); + }; + Ok(&value.value) + }) + .collect::>>()?; + + let args = args_to_vector(&signature.type_signature, &arg_values, query_ctx)?; + let arg_types = args + .iter() + .map(|arg| arg.data_type().as_arrow_type()) + .collect::>(); + let ret_type = + admin_udf + .return_type(&arg_types) + .map_err(|e| error::Error::BuildAdminFunctionArgs { + msg: format!( + "Failed to get return type of admin function {}: {}", + fn_name, e + ), + })?; + + Ok(ResolvedAdminFunction { + admin_udf, + fn_name, + args, + arg_types, + ret_type, + }) +} + +/// Output schema of an `ADMIN` statement, mirroring what +/// [`CoreAdminFunctionService::execute`] produces. `None` if the function +/// or arguments are unresolvable; execution will then surface the error. +pub fn admin_output_schema(stmt: &Admin, query_ctx: &QueryContextRef) -> Option { + let resolved = + resolve_admin_function(stmt, query_ctx, Arc::new(FunctionState::default())).ok()?; + Some(Schema::new(vec![ColumnSchema::new( + // Use statement as the result column name + stmt.to_string(), + ConcreteDataType::from_arrow_type(&resolved.ret_type), + false, + )])) +} + impl CoreAdminFunctionService { fn new(query_engine: query::QueryEngineRef) -> Self { Self { query_engine } @@ -88,62 +186,23 @@ impl CoreAdminFunctionService { query_ctx, } = request; - let Admin::Func(func) = &stmt; - // the function name should be in lower case. - let func_name = func.name.to_string().to_lowercase(); - let factory = get_admin_function(&func_name) - .or_else(|| FUNCTION_REGISTRY.get_function(&func_name)) - .context(error::AdminFunctionNotFoundSnafu { - name: func_name.clone(), - })?; - - let func_ctx = FunctionContext { - query_ctx: query_ctx.clone(), - state: self.query_engine.engine_state().function_state(), - }; - - let admin_udf = factory.provide(func_ctx); + let resolved = resolve_admin_function( + &stmt, + &query_ctx, + self.query_engine.engine_state().function_state(), + )?; + let ResolvedAdminFunction { + admin_udf, + fn_name, + args, + arg_types, + ret_type, + } = resolved; let admin_async_fn = admin_udf .as_async() - .context(error::AdminFunctionNotFoundSnafu { name: func_name })?; - - let fn_name = admin_udf.name(); - let signature = admin_udf.signature(); - - // Parse function arguments - let FunctionArguments::List(args) = &func.args else { - return error::BuildAdminFunctionArgsSnafu { - msg: format!("unsupported function args {} for {}", func.args, fn_name), - } - .fail(); - }; - let arg_values = args - .args - .iter() - .map(|arg| { - let FunctionArg::Unnamed(FunctionArgExpr::Expr(Expr::Value(value))) = arg else { - return error::BuildAdminFunctionArgsSnafu { - msg: format!("unsupported function arg {arg} for {}", fn_name), - } - .fail(); - }; - Ok(&value.value) - }) - .collect::>>()?; - - let args = args_to_vector(&signature.type_signature, &arg_values, &query_ctx)?; - let arg_types = args - .iter() - .map(|arg| arg.data_type().as_arrow_type()) - .collect::>(); - let ret_type = admin_udf.return_type(&arg_types).map_err(|e| { - error::Error::BuildAdminFunctionArgs { - msg: format!( - "Failed to get return type of admin function {}: {}", - fn_name, e - ), - } - })?; + .context(error::AdminFunctionNotFoundSnafu { + name: fn_name.clone(), + })?; // Convert arguments to DataFusion ColumnarValue format let columnar_args: Vec = args @@ -174,9 +233,7 @@ impl CoreAdminFunctionService { let result_columnar = admin_async_fn .invoke_async_with_args(func_args) .await - .with_context(|_| ExecuteAdminFunctionSnafu { - msg: fn_name.to_string(), - })?; + .with_context(|_| ExecuteAdminFunctionSnafu { msg: fn_name })?; // Convert result back to VectorRef let result_columnar: common_query::prelude::ColumnarValue = diff --git a/src/query/src/analyze.rs b/src/query/src/analyze.rs index b1522fc05d..2ab3ec5176 100644 --- a/src/query/src/analyze.rs +++ b/src/query/src/analyze.rs @@ -46,13 +46,9 @@ const STAGE: &str = "stage"; const NODE: &str = "node"; const PLAN: &str = "plan"; -/// Returns the fixed output schema of [`DistAnalyzeExec`]: -/// (`stage`: UInt32, `node`: UInt32, `plan`: Utf8). -/// -/// Exposed so protocol-level `Describe` handlers can advertise the schema -/// clients will actually receive (the physical plan is rewritten in -/// `optimize_physical_plan` and its schema differs from the DataFusion -/// `Analyze` logical plan schema). +/// Fixed output schema of [`DistAnalyzeExec`], for `Describe` handlers: +/// execution rewrites the plan in `optimize_physical_plan`, so this schema +/// differs from the logical `Analyze` plan's. pub fn dist_analyze_output_schema() -> SchemaRef { SchemaRef::new(Schema::new(vec![ Field::new(STAGE, DataType::UInt32, true), diff --git a/src/query/src/sql.rs b/src/query/src/sql.rs index b2d979c666..93e74ca33e 100644 --- a/src/query/src/sql.rs +++ b/src/query/src/sql.rs @@ -41,6 +41,7 @@ use common_recordbatch::RecordBatches; use common_recordbatch::adapter::RecordBatchStreamAdapter; use common_time::Timestamp; use common_time::timezone::get_timezone; +use datafusion::dataframe::DataFrame; use datafusion::prelude::SessionContext; use datafusion_expr::{Expr, SortExpr, col, lit}; use datatypes::prelude::*; @@ -109,7 +110,7 @@ const INDEX_KEY_NAME_COLUMN: &str = "Key_name"; const INDEX_SEQ_IN_INDEX_COLUMN: &str = "Seq_in_index"; const INDEX_COLUMN_NAME_COLUMN: &str = "Column_name"; -static DESCRIBE_TABLE_OUTPUT_SCHEMA: Lazy> = Lazy::new(|| { +pub static DESCRIBE_TABLE_OUTPUT_SCHEMA: Lazy> = Lazy::new(|| { Arc::new(Schema::new(vec![ ColumnSchema::new( COLUMN_NAME_COLUMN, @@ -178,6 +179,18 @@ pub async fn show_databases( catalog_manager: &CatalogManagerRef, query_ctx: QueryContextRef, ) -> Result { + let dataframe = + show_databases_dataframe(&stmt, query_engine, catalog_manager, query_ctx).await?; + dataframe_to_output(dataframe).await +} + +/// Builds the [`DataFrame`] for `SHOW DATABASES` without executing it. +pub async fn show_databases_dataframe( + stmt: &ShowDatabases, + query_engine: &QueryEngineRef, + catalog_manager: &CatalogManagerRef, + query_ctx: QueryContextRef, +) -> Result { let projects = if stmt.full { vec![ (schemata::SCHEMA_NAME, SCHEMAS_COLUMN), @@ -191,7 +204,7 @@ pub async fn show_databases( let like_field = Some(schemata::SCHEMA_NAME); let sort = vec![col(schemata::SCHEMA_NAME).sort(true, true)]; - query_from_information_schema_table( + query_from_information_schema_dataframe( query_engine, catalog_manager, query_ctx, @@ -201,7 +214,7 @@ pub async fn show_databases( filters, like_field, sort, - stmt.kind, + &stmt.kind, ) .await } @@ -249,6 +262,45 @@ async fn query_from_information_schema_table( sort: Vec, kind: ShowKind, ) -> Result { + let dataframe = query_from_information_schema_dataframe( + query_engine, + catalog_manager, + query_ctx, + table_name, + select, + projects, + filters, + like_field, + sort, + &kind, + ) + .await?; + dataframe_to_output(dataframe).await +} + +async fn dataframe_to_output(dataframe: DataFrame) -> Result { + let stream = dataframe.execute_stream().await?; + Ok(Output::new_with_stream(Box::pin( + RecordBatchStreamAdapter::try_new(stream).context(error::CreateRecordBatchSnafu)?, + ))) +} + +/// Builds the [`DataFrame`] for a `SHOW` statement without executing it, +/// so `Describe` handlers can derive the output schema from the same +/// projection the executor uses. +#[allow(clippy::too_many_arguments)] +async fn query_from_information_schema_dataframe( + query_engine: &QueryEngineRef, + catalog_manager: &CatalogManagerRef, + query_ctx: QueryContextRef, + table_name: &str, + select: Vec, + projects: Vec<(&str, &str)>, + filters: Vec, + like_field: Option<&str>, + sort: Vec, + kind: &ShowKind, +) -> Result { let table = catalog_manager .table( query_ctx.current_catalog(), @@ -336,7 +388,7 @@ async fn query_from_information_schema_table( .expect("Must be the datafusion planner"); let filter = planner - .sql_to_expr(filter, dataframe.schema(), false, query_ctx) + .sql_to_expr(filter.clone(), dataframe.schema(), false, query_ctx) .await?; // Apply the `where` clause filters @@ -344,11 +396,7 @@ async fn query_from_information_schema_table( } }; - let stream = dataframe.execute_stream().await?; - - Ok(Output::new_with_stream(Box::pin( - RecordBatchStreamAdapter::try_new(stream).context(error::CreateRecordBatchSnafu)?, - ))) + Ok(dataframe) } /// Execute `SHOW COLUMNS` statement. @@ -358,8 +406,19 @@ pub async fn show_columns( catalog_manager: &CatalogManagerRef, query_ctx: QueryContextRef, ) -> Result { - let schema_name = if let Some(database) = stmt.database { - database + let dataframe = show_columns_dataframe(&stmt, query_engine, catalog_manager, query_ctx).await?; + dataframe_to_output(dataframe).await +} + +/// Builds the [`DataFrame`] for `SHOW COLUMNS` without executing it. +pub async fn show_columns_dataframe( + stmt: &ShowColumns, + query_engine: &QueryEngineRef, + catalog_manager: &CatalogManagerRef, + query_ctx: QueryContextRef, +) -> Result { + let schema_name = if let Some(database) = &stmt.database { + database.clone() } else { query_ctx.current_schema() }; @@ -397,7 +456,7 @@ pub async fn show_columns( let like_field = Some(columns::COLUMN_NAME); let sort = vec![col(columns::COLUMN_NAME).sort(true, true)]; - query_from_information_schema_table( + query_from_information_schema_dataframe( query_engine, catalog_manager, query_ctx, @@ -407,7 +466,7 @@ pub async fn show_columns( filters, like_field, sort, - stmt.kind, + &stmt.kind, ) .await } @@ -419,8 +478,19 @@ pub async fn show_index( catalog_manager: &CatalogManagerRef, query_ctx: QueryContextRef, ) -> Result { - let schema_name = if let Some(database) = stmt.database { - database + let dataframe = show_index_dataframe(&stmt, query_engine, catalog_manager, query_ctx).await?; + dataframe_to_output(dataframe).await +} + +/// Builds the [`DataFrame`] for `SHOW INDEX` without executing it. +pub async fn show_index_dataframe( + stmt: &ShowIndex, + query_engine: &QueryEngineRef, + catalog_manager: &CatalogManagerRef, + query_ctx: QueryContextRef, +) -> Result { + let schema_name = if let Some(database) = &stmt.database { + database.clone() } else { query_ctx.current_schema() }; @@ -471,7 +541,7 @@ pub async fn show_index( col(statistics::SEQ_IN_INDEX).sort(true, true), ]; - query_from_information_schema_table( + query_from_information_schema_dataframe( query_engine, catalog_manager, query_ctx, @@ -481,7 +551,7 @@ pub async fn show_index( filters, like_field, sort, - stmt.kind, + &stmt.kind, ) .await } @@ -493,8 +563,19 @@ pub async fn show_region( catalog_manager: &CatalogManagerRef, query_ctx: QueryContextRef, ) -> Result { - let schema_name = if let Some(database) = stmt.database { - database + let dataframe = show_region_dataframe(&stmt, query_engine, catalog_manager, query_ctx).await?; + dataframe_to_output(dataframe).await +} + +/// Builds the [`DataFrame`] for `SHOW REGION` without executing it. +pub async fn show_region_dataframe( + stmt: &ShowRegion, + query_engine: &QueryEngineRef, + catalog_manager: &CatalogManagerRef, + query_ctx: QueryContextRef, +) -> Result { + let schema_name = if let Some(database) = &stmt.database { + database.clone() } else { query_ctx.current_schema() }; @@ -517,7 +598,7 @@ pub async fn show_region( col(columns::PEER_ID).sort(true, true), ]; - query_from_information_schema_table( + query_from_information_schema_dataframe( query_engine, catalog_manager, query_ctx, @@ -527,7 +608,7 @@ pub async fn show_region( filters, like_field, sort, - stmt.kind, + &stmt.kind, ) .await } @@ -539,8 +620,19 @@ pub async fn show_tables( catalog_manager: &CatalogManagerRef, query_ctx: QueryContextRef, ) -> Result { - let schema_name = if let Some(database) = stmt.database { - database + let dataframe = show_tables_dataframe(&stmt, query_engine, catalog_manager, query_ctx).await?; + dataframe_to_output(dataframe).await +} + +/// Builds the [`DataFrame`] for [`ShowTables`] without executing it. +pub async fn show_tables_dataframe( + stmt: &ShowTables, + query_engine: &QueryEngineRef, + catalog_manager: &CatalogManagerRef, + query_ctx: QueryContextRef, +) -> Result { + let schema_name = if let Some(database) = &stmt.database { + database.clone() } else { query_ctx.current_schema() }; @@ -564,15 +656,16 @@ pub async fn show_tables( // Transform the WHERE clause for backward compatibility: // Replace "Tables" with "Tables_in_{schema}" to support old queries - let kind = match stmt.kind { - ShowKind::Where(mut filter) => { + let rewritten_kind = match &stmt.kind { + ShowKind::Where(filter) => { + let mut filter = filter.clone(); replace_column_in_expr(&mut filter, "Tables", &tables_column); ShowKind::Where(filter) } - other => other, + kind => kind.clone(), }; - query_from_information_schema_table( + query_from_information_schema_dataframe( query_engine, catalog_manager, query_ctx, @@ -582,7 +675,7 @@ pub async fn show_tables( filters, like_field, sort, - kind, + &rewritten_kind, ) .await } @@ -594,8 +687,20 @@ pub async fn show_table_status( catalog_manager: &CatalogManagerRef, query_ctx: QueryContextRef, ) -> Result { - let schema_name = if let Some(database) = stmt.database { - database + let dataframe = + show_table_status_dataframe(&stmt, query_engine, catalog_manager, query_ctx).await?; + dataframe_to_output(dataframe).await +} + +/// Builds the [`DataFrame`] for [`ShowTableStatus`] without executing it. +pub async fn show_table_status_dataframe( + stmt: &ShowTableStatus, + query_engine: &QueryEngineRef, + catalog_manager: &CatalogManagerRef, + query_ctx: QueryContextRef, +) -> Result { + let schema_name = if let Some(database) = &stmt.database { + database.clone() } else { query_ctx.current_schema() }; @@ -629,7 +734,7 @@ pub async fn show_table_status( let like_field = Some(tables::TABLE_NAME); let sort = vec![col(tables::TABLE_NAME).sort(true, true)]; - query_from_information_schema_table( + query_from_information_schema_dataframe( query_engine, catalog_manager, query_ctx, @@ -639,7 +744,7 @@ pub async fn show_table_status( filters, like_field, sort, - stmt.kind, + &stmt.kind, ) .await } @@ -651,6 +756,18 @@ pub async fn show_collations( catalog_manager: &CatalogManagerRef, query_ctx: QueryContextRef, ) -> Result { + let dataframe = + show_collations_dataframe(&kind, query_engine, catalog_manager, query_ctx).await?; + dataframe_to_output(dataframe).await +} + +/// Builds the [`DataFrame`] for `SHOW COLLATION` without executing it. +pub async fn show_collations_dataframe( + kind: &ShowKind, + query_engine: &QueryEngineRef, + catalog_manager: &CatalogManagerRef, + query_ctx: QueryContextRef, +) -> Result { // Refer to https://dev.mysql.com/doc/refman/8.0/en/show-collation.html let projects = vec![ ("collation_name", "Collation"), @@ -665,7 +782,7 @@ pub async fn show_collations( let like_field = Some("collation_name"); let sort = vec![]; - query_from_information_schema_table( + query_from_information_schema_dataframe( query_engine, catalog_manager, query_ctx, @@ -687,6 +804,18 @@ pub async fn show_charsets( catalog_manager: &CatalogManagerRef, query_ctx: QueryContextRef, ) -> Result { + let dataframe = + show_charsets_dataframe(&kind, query_engine, catalog_manager, query_ctx).await?; + dataframe_to_output(dataframe).await +} + +/// Builds the [`DataFrame`] for `SHOW CHARSET` without executing it. +pub async fn show_charsets_dataframe( + kind: &ShowKind, + query_engine: &QueryEngineRef, + catalog_manager: &CatalogManagerRef, + query_ctx: QueryContextRef, +) -> Result { // Refer to https://dev.mysql.com/doc/refman/8.0/en/show-character-set.html let projects = vec![ ("character_set_name", "Charset"), @@ -699,7 +828,7 @@ pub async fn show_charsets( let like_field = Some("character_set_name"); let sort = vec![]; - query_from_information_schema_table( + query_from_information_schema_dataframe( query_engine, catalog_manager, query_ctx, @@ -916,8 +1045,19 @@ pub async fn show_views( catalog_manager: &CatalogManagerRef, query_ctx: QueryContextRef, ) -> Result { - let schema_name = if let Some(database) = stmt.database { - database + let dataframe = show_views_dataframe(&stmt, query_engine, catalog_manager, query_ctx).await?; + dataframe_to_output(dataframe).await +} + +/// Builds the [`DataFrame`] for [`ShowViews`] without executing it. +pub async fn show_views_dataframe( + stmt: &ShowViews, + query_engine: &QueryEngineRef, + catalog_manager: &CatalogManagerRef, + query_ctx: QueryContextRef, +) -> Result { + let schema_name = if let Some(database) = &stmt.database { + database.clone() } else { query_ctx.current_schema() }; @@ -930,7 +1070,7 @@ pub async fn show_views( let like_field = Some(tables::TABLE_NAME); let sort = vec![col(tables::TABLE_NAME).sort(true, true)]; - query_from_information_schema_table( + query_from_information_schema_dataframe( query_engine, catalog_manager, query_ctx, @@ -940,7 +1080,7 @@ pub async fn show_views( filters, like_field, sort, - stmt.kind, + &stmt.kind, ) .await } @@ -952,12 +1092,23 @@ pub async fn show_flows( catalog_manager: &CatalogManagerRef, query_ctx: QueryContextRef, ) -> Result { + let dataframe = show_flows_dataframe(&stmt, query_engine, catalog_manager, query_ctx).await?; + dataframe_to_output(dataframe).await +} + +/// Builds the [`DataFrame`] for [`ShowFlows`] without executing it. +pub async fn show_flows_dataframe( + stmt: &ShowFlows, + query_engine: &QueryEngineRef, + catalog_manager: &CatalogManagerRef, + query_ctx: QueryContextRef, +) -> Result { let projects = vec![(flows::FLOW_NAME, FLOWS_COLUMN)]; let filters = vec![col(flows::TABLE_CATALOG).eq(lit(query_ctx.current_catalog()))]; let like_field = Some(flows::FLOW_NAME); let sort = vec![col(flows::FLOW_NAME).sort(true, true)]; - query_from_information_schema_table( + query_from_information_schema_dataframe( query_engine, catalog_manager, query_ctx, @@ -967,7 +1118,7 @@ pub async fn show_flows( filters, like_field, sort, - stmt.kind, + &stmt.kind, ) .await } @@ -1340,6 +1491,18 @@ pub async fn show_processlist( catalog_manager: &CatalogManagerRef, query_ctx: QueryContextRef, ) -> Result { + let dataframe = + show_processlist_dataframe(&stmt, query_engine, catalog_manager, query_ctx).await?; + dataframe_to_output(dataframe).await +} + +/// Builds the [`DataFrame`] for `SHOW PROCESSLIST` without executing it. +pub async fn show_processlist_dataframe( + stmt: &ShowProcessList, + query_engine: &QueryEngineRef, + catalog_manager: &CatalogManagerRef, + query_ctx: QueryContextRef, +) -> Result { let projects = if stmt.full { vec![ (process_list::ID, "Id"), @@ -1367,17 +1530,17 @@ pub async fn show_processlist( }; let like_field = None; let sort = vec![col("id").sort(true, true)]; - query_from_information_schema_table( + query_from_information_schema_dataframe( query_engine, catalog_manager, query_ctx.clone(), "process_list", vec![], - projects.clone(), + projects, filters, like_field, sort, - ShowKind::All, + &ShowKind::All, ) .await } diff --git a/src/servers/src/postgres/handler.rs b/src/servers/src/postgres/handler.rs index 433eb11cda..aca4362bd3 100644 --- a/src/servers/src/postgres/handler.rs +++ b/src/servers/src/postgres/handler.rs @@ -28,6 +28,7 @@ use datafusion_pg_catalog::sql::PostgresCompatibilityParser; use datatypes::prelude::ConcreteDataType; use datatypes::schema::{Schema, SchemaRef}; use futures::{Sink, SinkExt, Stream, StreamExt, future, stream}; +use operator::statement::admin_output_schema; use pgwire::api::portal::{Format, Portal}; use pgwire::api::query::{ExtendedQueryHandler, SimpleQueryHandler}; use pgwire::api::results::{ @@ -43,6 +44,7 @@ use pgwire::messages::data::DataRow; use query::dist_analyze_output_schema; use query::planner::DfLogicalPlanner; use query::query_engine::DescribeResult; +use query::sql::DESCRIBE_TABLE_OUTPUT_SCHEMA; use session::Session; use session::context::QueryContextRef; use snafu::ResultExt; @@ -555,11 +557,8 @@ fn describe_fields( session: &Arc, ) -> PgWireResult> { match sql_plan { - // EXPLAIN ANALYZE: at execution time the physical plan is replaced with - // GreptimeDB's `DistAnalyzeExec` (see `optimize_physical_plan`), whose - // output schema (stage/node/plan) differs from the DataFusion `Analyze` - // logical plan schema (plan_type/plan). Describe with the schema the - // client will actually receive so the DataRow field count matches. + // Execution swaps in DistAnalyzeExec (stage/node/plan), whose schema + // differs from the logical `Analyze` plan's (plan_type/plan). SqlPlan::Plan(LogicalPlan::Analyze(_), _) => { let schema: Schema = Schema::try_from(dist_analyze_output_schema()).map_err(convert_err)?; @@ -658,17 +657,6 @@ fn describe_fields( ), ]), - // single column show statements - SqlPlan::Statement( - Statement::ShowTables(_) | Statement::ShowFlows(_) | Statement::ShowViews(_), - _, - ) => Ok(vec![FieldInfo::new( - "name".to_string(), - None, - None, - Type::TEXT, - format.format_for(0), - )]), #[cfg(feature = "enterprise")] SqlPlan::Statement(Statement::ShowTriggers(_), _) => Ok(vec![FieldInfo::new( "name".to_string(), @@ -690,22 +678,59 @@ fn describe_fields( Ok(vec![]) } } - // FETCH cursor: return the cursor's schema so the RowDescription - // matches the DataRow field count sent during Execute. + // Single column named after the variable (see `query::sql::show_variable`). + SqlPlan::Statement(Statement::ShowVariables(show), _) => Ok(vec![FieldInfo::new( + show.variable.to_string().to_uppercase(), + None, + None, + Type::TEXT, + format.format_for(0), + )]), + // Mirrors `query::sql::show_status` (currently always empty). + SqlPlan::Statement(Statement::ShowStatus(_), _) => Ok(vec![ + FieldInfo::new( + "Variable_name".to_string(), + None, + None, + Type::TEXT, + format.format_for(0), + ), + FieldInfo::new( + "Value".to_string(), + None, + None, + Type::TEXT, + format.format_for(1), + ), + ]), + SqlPlan::Statement(Statement::ShowSearchPath(_), _) => Ok(vec![FieldInfo::new( + "search_path".to_string(), + None, + None, + Type::TEXT, + format.format_for(0), + )]), + // Mirrors `query::sql::describe_table`. + SqlPlan::Statement(Statement::DescribeTable(_), _) => { + schema_to_pg(&DESCRIBE_TABLE_OUTPUT_SCHEMA, format, None).map_err(convert_err) + } + // Single column typed with the function's return type (see + // `operator::statement::admin_output_schema`). + SqlPlan::Statement(Statement::Admin(admin), _) => { + let query_ctx = session.new_query_context(); + match admin_output_schema(admin, &query_ctx) { + Some(schema) => schema_to_pg(&schema, format, None).map_err(convert_err), + // Unresolvable; execution will surface the error. + None => Ok(vec![]), + } + } + // Describe from the declared cursor's schema. SqlPlan::Statement(Statement::FetchCursor(fetch), _) => { let cursor_name = fetch.cursor_name.to_string(); match session.get_cursor(&cursor_name) { - Some(cursor) => { - // `cursor.schema()` is the GreptimeDB `SchemaRef` captured - // when the cursor was declared; `schema_to_pg` accepts it - // directly (no DataFusion -> GreptimeDB conversion needed). - schema_to_pg(&cursor.schema(), format, None).map_err(convert_err) - } - None => { - // Cursor not found (e.g. DECLARE hasn't executed yet in - // an extended-protocol batch). Return NoData as a fallback. - Ok(vec![]) - } + Some(cursor) => schema_to_pg(&cursor.schema(), format, None).map_err(convert_err), + // Cursor not declared yet; execution will error. + None => Ok(vec![]), } } _ => { diff --git a/src/session/src/lib.rs b/src/session/src/lib.rs index 44f3ee6689..c1414fc5a3 100644 --- a/src/session/src/lib.rs +++ b/src/session/src/lib.rs @@ -111,10 +111,7 @@ impl Session { .into() } - /// Returns the cursor with the given name, if it exists. - /// - /// Cursors are stored in the session's mutable inner data and shared across - /// all query contexts created from this session. + /// Cursors are shared across query contexts created from this session. pub fn get_cursor(&self, name: &str) -> Option> { let guard = self.mutable_inner.read().unwrap(); guard.cursors.get(name).cloned() diff --git a/tests-integration/tests/sql.rs b/tests-integration/tests/sql.rs index 32a2881dd7..c95adb483e 100644 --- a/tests-integration/tests/sql.rs +++ b/tests-integration/tests/sql.rs @@ -91,6 +91,8 @@ macro_rules! sql_tests { test_mysql_prepare_stmt_insert_timestamp, test_mysql_prepare_stmt_timezone, test_mysql_federated_prepare_stmt, + test_mysql_prepare_tql_and_show, + test_postgres_extended_query_row_returning_statements, test_declare_fetch_close_cursor, test_alter_update_on, ); @@ -1487,6 +1489,280 @@ fn assert_pg_numeric_range_error(error: tokio_postgres::Error) { assert_eq!("numeric_value_out_of_range", error.message()); } +pub async fn test_postgres_extended_query_row_returning_statements(store_type: StorageType) { + // Regression test for the tokio-postgres >= 0.7.14 DataRow/RowDescription + // mismatch: statements answered with NoData at Describe but emitting + // DataRows at Execute must describe their real output schema. + let (mut guard, fe_pg_server) = setup_pg_server(store_type, "test_pg_extended_row_stmts").await; + let addr = fe_pg_server.bind_addr().unwrap().to_string(); + + let (client, connection) = tokio_postgres::connect(&format!("postgres://{addr}/public"), NoTls) + .await + .unwrap(); + + let (tx, rx) = tokio::sync::oneshot::channel(); + tokio::spawn(async move { + connection.await.unwrap(); + tx.send(()).unwrap(); + }); + + client + .execute( + "CREATE TABLE demo_metrics (ts timestamp time index, val double, host string primary key skipping index)", + &[], + ) + .await + .unwrap(); + client + .execute( + "INSERT INTO demo_metrics (ts, host, val) VALUES (1000, 'host-a', 1.0), (2000, 'host-b', 2.0)", + &[], + ) + .await + .unwrap(); + + // ---- SHOW DATABASES: single `Database` column ---- + let rows = client.query("SHOW DATABASES", &[]).await.unwrap(); + assert!(!rows.is_empty()); + assert_eq!(1, rows[0].columns().len()); + assert_eq!("Database", rows[0].columns()[0].name()); + assert!(rows.iter().any(|r| r.get::<_, String>(0) == "public")); + + // ---- SHOW FULL DATABASES: `Database` + `Options` columns ---- + let rows = client.query("SHOW FULL DATABASES", &[]).await.unwrap(); + assert!(!rows.is_empty()); + assert_eq!(2, rows[0].columns().len()); + assert_eq!("Database", rows[0].columns()[0].name()); + assert_eq!("Options", rows[0].columns()[1].name()); + assert!(rows.iter().any(|r| r.get::<_, String>(0) == "public")); + + // ---- SHOW TABLES: single `Tables_in_` column ---- + let rows = client.query("SHOW TABLES", &[]).await.unwrap(); + assert!(!rows.is_empty()); + assert_eq!(1, rows[0].columns().len()); + assert_eq!("Tables_in_public", rows[0].columns()[0].name()); + assert!(rows.iter().any(|r| r.get::<_, String>(0) == "demo_metrics")); + + // ---- SHOW FULL TABLES: `Tables_in_` + `Table_type` columns ---- + let rows = client.query("SHOW FULL TABLES", &[]).await.unwrap(); + assert!(!rows.is_empty()); + assert_eq!(2, rows[0].columns().len()); + assert_eq!("Tables_in_public", rows[0].columns()[0].name()); + assert_eq!("Table_type", rows[0].columns()[1].name()); + + // ---- SHOW VIEWS / SHOW FLOWS: empty results, still described with one column ---- + let rows = client.query("SHOW VIEWS", &[]).await.unwrap(); + assert!(rows.is_empty()); + let stmt = client.prepare("SHOW VIEWS").await.unwrap(); + assert_eq!(1, stmt.columns().len()); + assert_eq!("Views", stmt.columns()[0].name()); + let rows = client.query("SHOW FLOWS", &[]).await.unwrap(); + assert!(rows.is_empty()); + let stmt = client.prepare("SHOW FLOWS").await.unwrap(); + assert_eq!(1, stmt.columns().len()); + assert_eq!("Flows", stmt.columns()[0].name()); + + // ---- SHOW TABLE STATUS: fixed eighteen-column schema ---- + let rows = client.query("SHOW TABLE STATUS", &[]).await.unwrap(); + assert!(!rows.is_empty()); + let names: Vec<&str> = rows[0].columns().iter().map(|c| c.name()).collect(); + assert_eq!( + vec![ + "Name", + "Engine", + "Version", + "Row_format", + "Rows", + "Avg_row_length", + "Data_length", + "Max_data_length", + "Index_length", + "Data_free", + "Auto_increment", + "Create_time", + "Update_time", + "Check_time", + "Collation", + "Checksum", + "Create_options", + "Comment", + ], + names + ); + + // ---- SHOW COLUMNS / SHOW FULL COLUMNS ---- + let rows = client + .query("SHOW COLUMNS FROM demo_metrics", &[]) + .await + .unwrap(); + assert_eq!(3, rows.len()); + let names: Vec<&str> = rows[0].columns().iter().map(|c| c.name()).collect(); + assert_eq!( + vec![ + "Field", + "Type", + "Null", + "Key", + "Default", + "Extra", + "Greptime_type" + ], + names + ); + let rows = client + .query("SHOW FULL COLUMNS FROM demo_metrics", &[]) + .await + .unwrap(); + assert_eq!(10, rows[0].columns().len()); + + // ---- SHOW CHARSET / SHOW COLLATION ---- + let rows = client.query("SHOW CHARSET", &[]).await.unwrap(); + assert!(!rows.is_empty()); + assert_eq!(4, rows[0].columns().len()); + let rows = client.query("SHOW COLLATION", &[]).await.unwrap(); + assert!(!rows.is_empty()); + assert_eq!(6, rows[0].columns().len()); + + // ---- SHOW INDEX: fixed fifteen-column schema ---- + let rows = client + .query("SHOW INDEX IN demo_metrics", &[]) + .await + .unwrap(); + assert!(!rows.is_empty()); + let names: Vec<&str> = rows[0].columns().iter().map(|c| c.name()).collect(); + assert_eq!( + vec![ + "Table", + "Non_unique", + "Key_name", + "Seq_in_index", + "Column_name", + "Collation", + "Cardinality", + "Sub_part", + "Packed", + "Null", + "Index_type", + "Comment", + "Index_comment", + "Visible", + "Expression", + ], + names + ); + + // ---- SHOW REGION ---- + let rows = client + .query("SHOW REGION IN demo_metrics", &[]) + .await + .unwrap(); + assert!(!rows.is_empty()); + assert_eq!(4, rows[0].columns().len()); + + // ---- SHOW SEARCH_PATH: single string column ---- + let rows = client.query("SHOW SEARCH_PATH", &[]).await.unwrap(); + assert_eq!(1, rows.len()); + assert_eq!(1, rows[0].columns().len()); + assert_eq!("search_path", rows[0].columns()[0].name()); + assert_eq!("public", rows[0].get::<_, String>(0)); + + // ---- SHOW VARIABLES: single column named after the variable ---- + let rows = client.query("SHOW VARIABLES timezone", &[]).await.unwrap(); + assert_eq!(1, rows.len()); + assert_eq!(1, rows[0].columns().len()); + assert_eq!("TIMEZONE", rows[0].columns()[0].name()); + let _ = rows[0].get::<_, String>(0); + + // ---- DESCRIBE TABLE: fixed six-column string schema ---- + let rows = client + .query("DESCRIBE TABLE demo_metrics", &[]) + .await + .unwrap(); + assert_eq!(3, rows.len()); + let names: Vec<&str> = rows[0].columns().iter().map(|c| c.name()).collect(); + assert_eq!( + vec!["Column", "Type", "Key", "Null", "Default", "Semantic Type"], + names + ); + // first column of each row is the column name; ensure values decode as TEXT + let columns: Vec = rows.iter().map(|r| r.get::<_, String>(0)).collect(); + assert!(columns.contains(&"ts".to_string())); + assert!(columns.contains(&"val".to_string())); + assert!(columns.contains(&"host".to_string())); + + // ---- ADMIN function: single column named after the statement ---- + let rows = client + .query("ADMIN flush_table('demo_metrics')", &[]) + .await + .unwrap(); + assert_eq!(1, rows.len()); + assert_eq!(1, rows[0].columns().len()); + assert!(rows[0].columns()[0].name().contains("flush_table")); + + // ---- TQL EVAL / EXPLAIN / ANALYZE: described from the planned query ---- + let rows = client + .query("TQL EVAL (0, 3000, '1s') demo_metrics", &[]) + .await + .unwrap(); + assert!(!rows.is_empty()); + let names: Vec<&str> = rows[0].columns().iter().map(|c| c.name()).collect(); + assert_eq!(vec!["ts", "val", "host"], names); + + let rows = client + .query("TQL EXPLAIN (0, 3000, '1s') demo_metrics", &[]) + .await + .unwrap(); + assert!(!rows.is_empty()); + assert_eq!(2, rows[0].columns().len()); + + let rows = client + .query("TQL ANALYZE (0, 3000, '1s') demo_metrics", &[]) + .await + .unwrap(); + assert!(!rows.is_empty()); + assert_eq!(3, rows[0].columns().len()); + + // FORMAT JSON variants keep working through the plan-based execution path + let rows = client + .query("TQL EXPLAIN FORMAT JSON (0, 3000, '1s') demo_metrics", &[]) + .await + .unwrap(); + assert!(!rows.is_empty()); + let rows = client + .query("TQL ANALYZE FORMAT JSON (0, 3000, '1s') demo_metrics", &[]) + .await + .unwrap(); + assert!(!rows.is_empty()); + + // ---- the same statements also work over the simple query protocol ---- + for sql in [ + "SHOW DATABASES", + "SHOW FULL TABLES", + "SHOW TABLE STATUS", + "SHOW COLUMNS FROM demo_metrics", + "SHOW CHARSET", + "SHOW COLLATION", + "SHOW INDEX IN demo_metrics", + "SHOW REGION IN demo_metrics", + "SHOW SEARCH_PATH", + "DESCRIBE TABLE demo_metrics", + "ADMIN flush_table('demo_metrics')", + "TQL EVAL (0, 3000, '1s') demo_metrics", + ] { + let msgs = client.simple_query(sql).await.unwrap(); + assert!( + msgs.iter().any(|m| matches!(m, SimpleQueryMessage::Row(_))), + "simple query {sql} should return rows" + ); + } + + drop(client); + rx.await.unwrap(); + + let _ = fe_pg_server.shutdown().await; + guard.remove_all().await; +} + pub async fn test_postgres_explain_bind_parameter(store_type: StorageType) { // Regression test for #8029: EXPLAIN / EXPLAIN ANALYZE must accept bind // parameters over the Postgres extended query protocol. @@ -1890,6 +2166,73 @@ pub async fn test_mysql_federated_prepare_stmt(store_type: StorageType) { guard.remove_all().await; } +pub async fn test_mysql_prepare_tql_and_show(store_type: StorageType) { + // `do_describe` now plans TQL and information-schema-backed SHOW + // statements, so MySQL prepared statements derive their column metadata + // from the planned query and execute the plan directly. This exercises + // that path end-to-end. + common_telemetry::init_default_ut_logging(); + + let (mut guard, fe_mysql_server) = + setup_mysql_server(store_type, "test_mysql_prepare_tql_and_show").await; + let addr = fe_mysql_server.bind_addr().unwrap().to_string(); + + let pool = MySqlPoolOptions::new() + .max_connections(2) + .connect(&format!("mysql://{addr}/public")) + .await + .unwrap(); + + sqlx::query( + "CREATE TABLE demo_metrics (ts timestamp time index, val double, host string primary key)", + ) + .execute(&pool) + .await + .unwrap(); + sqlx::query("INSERT INTO demo_metrics (ts, host, val) VALUES (1000, 'host-a', 1.0), (2000, 'host-b', 2.0)") + .execute(&pool) + .await + .unwrap(); + + // sqlx::query uses the binary prepared statement protocol + // (COM_STMT_PREPARE + COM_STMT_EXECUTE). + let rows = sqlx::query("TQL EVAL (0, 3000, '1s') demo_metrics") + .fetch_all(&pool) + .await + .unwrap(); + assert!(!rows.is_empty()); + assert_eq!(3, rows[0].columns().len()); + + let rows = sqlx::query("TQL ANALYZE (0, 3000, '1s') demo_metrics") + .fetch_all(&pool) + .await + .unwrap(); + assert!(!rows.is_empty()); + + // SHOW statements share the same describe path; prepared SHOW FULL + // TABLES must report both columns. + let rows = sqlx::query("SHOW TABLES").fetch_all(&pool).await.unwrap(); + assert!(!rows.is_empty()); + assert_eq!(1, rows[0].columns().len()); + + let rows = sqlx::query("SHOW FULL TABLES") + .fetch_all(&pool) + .await + .unwrap(); + assert!(!rows.is_empty()); + assert_eq!(2, rows[0].columns().len()); + + let rows = sqlx::query("SHOW DATABASES") + .fetch_all(&pool) + .await + .unwrap(); + assert!(!rows.is_empty()); + assert_eq!(1, rows[0].columns().len()); + + let _ = fe_mysql_server.shutdown().await; + guard.remove_all().await; +} + pub async fn test_postgres_array_types(store_type: StorageType) { let (mut guard, fe_pg_server) = setup_pg_server(store_type, "test_postgres_array_types").await; let addr = fe_pg_server.bind_addr().unwrap().to_string();