From e1d18f02f07e3270d94908c0390fea3b5e9e98fd Mon Sep 17 00:00:00 2001 From: zhyass <34016424+zhyass@users.noreply.github.com> Date: Thu, 30 Jul 2026 19:54:12 +0800 Subject: [PATCH 1/3] fix(meta): disallow replacing tables with different engines --- src/meta/api/src/api_impl/schema_api.rs | 9 +++ src/meta/api/src/api_impl/table_api.rs | 36 ++++++++++ src/meta/app/src/app_error.rs | 30 ++++++++ .../src/schema_api_test_suite.rs | 70 +++++++++++++++++++ .../storages/common/session/src/temp_table.rs | 29 ++++++-- .../06_0016_stream_replace_engine.test | 65 +++++++++++++++++ .../temp_table/create_temp_tables.test | 6 ++ 7 files changed, 241 insertions(+), 4 deletions(-) create mode 100644 tests/sqllogictests/suites/ee/06_ee_stream/06_0016_stream_replace_engine.test diff --git a/src/meta/api/src/api_impl/schema_api.rs b/src/meta/api/src/api_impl/schema_api.rs index 5f5ca5cb2555e..f8f2af33aff5a 100644 --- a/src/meta/api/src/api_impl/schema_api.rs +++ b/src/meta/api/src/api_impl/schema_api.rs @@ -21,6 +21,7 @@ use chrono::DateTime; use chrono::Utc; use databend_common_meta_app::app_error::AppError; use databend_common_meta_app::app_error::DropTableWithDropTime; +use databend_common_meta_app::app_error::TableEngineMismatch; use databend_common_meta_app::app_error::UndropTableAlreadyExists; use databend_common_meta_app::app_error::UndropTableHasNoHistory; use databend_common_meta_app::app_error::UndropTableRetentionGuard; @@ -177,6 +178,7 @@ pub async fn construct_drop_table_txn_operations( tenant: &Tenant, catalog_name: Option, table_id: u64, + expected_engine: Option<&str>, db_id: u64, if_exists: bool, if_delete: bool, @@ -233,6 +235,13 @@ pub async fn construct_drop_table_txn_operations( ); let mut tb_meta = tb_meta.unwrap(); + if let Some(expected_engine) = expected_engine { + if !tb_meta.engine.eq_ignore_ascii_case(expected_engine) { + return Err(KVAppError::AppError( + TableEngineMismatch::new(table_name, &tb_meta.engine, expected_engine).into(), + )); + } + } // drop a table with drop_on time if tb_meta.drop_on.is_some() { return if if_exists { diff --git a/src/meta/api/src/api_impl/table_api.rs b/src/meta/api/src/api_impl/table_api.rs index 9b524c965f987..f3851c86e56ed 100644 --- a/src/meta/api/src/api_impl/table_api.rs +++ b/src/meta/api/src/api_impl/table_api.rs @@ -34,6 +34,7 @@ use databend_common_meta_app::app_error::MultiStmtTxnCommitFailed; use databend_common_meta_app::app_error::StreamAlreadyExists; use databend_common_meta_app::app_error::StreamVersionMismatched; use databend_common_meta_app::app_error::TableAlreadyExists; +use databend_common_meta_app::app_error::TableEngineMismatch; use databend_common_meta_app::app_error::TableSnapshotExpired; use databend_common_meta_app::app_error::TableVersionMismatched; use databend_common_meta_app::app_error::UndropTableHasNoHistory; @@ -362,6 +363,37 @@ where } CreateOption::CreateOrReplace => { if req.as_dropped { + // CTAS does not call construct_drop_table_txn_operations(), + // so validate its existing table here. + let existing_meta = + self.get_pb(&TableId::new(*id.data)).await?.ok_or_else(|| { + KVAppError::AppError(AppError::UnknownTableId( + UnknownTableId::new( + *id.data, + "create or replace failed to find existing table meta", + ), + )) + })?; + if !existing_meta + .engine + .eq_ignore_ascii_case(&req.table_meta.engine) + { + return Err(KVAppError::AppError( + TableEngineMismatch::new( + req.table_name(), + &existing_meta.engine, + &req.table_meta.engine, + ) + .into(), + )); + } + + // Guard the engine check against a concurrent metadata update. + txn.condition.push(txn_cond_seq( + &TableId::new(*id.data), + Eq, + existing_meta.seq, + )); // If the table is being created as a dropped table, we do not // need to combine with drop_table_txn operations, just return // the sequence number associated with the value part of @@ -369,12 +401,15 @@ where SeqV::new(id.seq, *id.data) } else { + // The drop helper validates the engine against the metadata + // sequence used by the replacement transaction. let (seq, id) = construct_drop_table_txn_operations( self, req.name_ident.table_name.clone(), &req.name_ident.tenant, req.catalog_name.clone(), *id.data, + Some(&req.table_meta.engine), *seq_db_id.data, true, false, @@ -611,6 +646,7 @@ where &req.tenant, None, table_id, + None, req.db_id, req.if_exists, true, diff --git a/src/meta/app/src/app_error.rs b/src/meta/app/src/app_error.rs index 82303633c1ddd..1272bd28f1d4a 100644 --- a/src/meta/app/src/app_error.rs +++ b/src/meta/app/src/app_error.rs @@ -195,6 +195,30 @@ impl TableAlreadyExists { } } +#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)] +#[error( + "Cannot replace '{table_name}': existing table uses engine {existing_engine}, but the new table uses engine {new_engine}" +)] +pub struct TableEngineMismatch { + table_name: String, + existing_engine: String, + new_engine: String, +} + +impl TableEngineMismatch { + pub fn new( + table_name: impl Into, + existing_engine: impl Into, + new_engine: impl Into, + ) -> Self { + Self { + table_name: table_name.into(), + existing_engine: existing_engine.into(), + new_engine: new_engine.into(), + } + } +} + #[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)] #[error("ViewAlreadyExists: {view_name} while {context}")] pub struct ViewAlreadyExists { @@ -895,6 +919,9 @@ pub enum AppError { #[error(transparent)] TableAlreadyExists(#[from] TableAlreadyExists), + #[error(transparent)] + TableEngineMismatch(#[from] TableEngineMismatch), + #[error(transparent)] ViewAlreadyExists(#[from] ViewAlreadyExists), @@ -1199,6 +1226,8 @@ impl AppErrorMessage for TableAlreadyExists { } } +impl AppErrorMessage for TableEngineMismatch {} + impl AppErrorMessage for ViewAlreadyExists { fn message(&self) -> String { format!("'{}' as view Already Exists", self.view_name) @@ -1412,6 +1441,7 @@ impl From for ErrorCode { ErrorCode::MaterializedViewAlreadyExists(err.message()) } AppError::TableAlreadyExists(err) => ErrorCode::TableAlreadyExists(err.message()), + AppError::TableEngineMismatch(err) => ErrorCode::TableEngineNotSupported(err.message()), AppError::ViewAlreadyExists(err) => ErrorCode::ViewAlreadyExists(err.message()), AppError::CreateTableWithDropTime(err) => { ErrorCode::CreateTableWithDropTime(err.message()) diff --git a/src/meta/schema-api-test-suite/src/schema_api_test_suite.rs b/src/meta/schema-api-test-suite/src/schema_api_test_suite.rs index 81b49c2ca8eb1..e14f2eade4198 100644 --- a/src/meta/schema-api-test-suite/src/schema_api_test_suite.rs +++ b/src/meta/schema-api-test-suite/src/schema_api_test_suite.rs @@ -2583,6 +2583,76 @@ impl SchemaApiTestSuite { let ret_table_name_ident: DBIdTableName = get_kv_data(mt, &key_table_id_to_name).await?; assert_eq!(ret_table_name_ident, key_dbid_tbname); + + // Engine names are case-insensitive. + let (_, replaced_meta) = util + .create_table_with( + |mut meta| { + meta.engine = "json".to_string(); + meta + }, + |mut req| { + req.create_option = CreateOption::CreateOrReplace; + req.name_ident.table_name = table.to_string(); + req + }, + ) + .await?; + assert_eq!("json", replaced_meta.engine); + + // Restore the upper-case engine for the mismatch cases below. + let (table_id, _) = util + .create_table_with( + |mut meta| { + meta.engine = "JSON".to_string(); + meta + }, + |mut req| { + req.create_option = CreateOption::CreateOrReplace; + req.name_ident.table_name = table.to_string(); + req + }, + ) + .await?; + + // Replacing a table with a different engine must fail without changing + // the existing table or its name mapping. + let mismatched_req = CreateTableReq { + create_option: CreateOption::CreateOrReplace, + catalog_name: Some("default".to_string()), + name_ident: TableNameIdent { + tenant: tenant.clone(), + db_name: db_name.to_string(), + table_name: table.to_string(), + }, + table_meta: TableMeta { + engine: "STREAM".to_string(), + ..Default::default() + }, + as_dropped: false, + materialized_view: None, + table_properties: None, + table_partition: None, + }; + let err = mt.create_table(mismatched_req.clone()).await.unwrap_err(); + assert!(matches!( + err, + KVAppError::AppError(AppError::TableEngineMismatch(_)) + )); + assert_eq!(table_id, get_kv_u64_data(mt, &key_dbid_tbname).await?); + assert_eq!("JSON", util.get_table_by_name(table).await?.meta.engine); + + // The same validation applies when CTAS creates its replacement as a + // hidden dropped table. + let mut as_dropped_req = mismatched_req; + as_dropped_req.as_dropped = true; + as_dropped_req.table_meta.drop_on = Some(Utc::now()); + let err = mt.create_table(as_dropped_req).await.unwrap_err(); + assert!(matches!( + err, + KVAppError::AppError(AppError::TableEngineMismatch(_)) + )); + assert_eq!(table_id, get_kv_u64_data(mt, &key_dbid_tbname).await?); } { diff --git a/src/query/storages/common/session/src/temp_table.rs b/src/query/storages/common/session/src/temp_table.rs index 5dbfe89bb6f35..23485f5afe432 100644 --- a/src/query/storages/common/session/src/temp_table.rs +++ b/src/query/storages/common/session/src/temp_table.rs @@ -19,6 +19,8 @@ use std::sync::Arc; use databend_common_exception::ErrorCode; use databend_common_exception::Result; +use databend_common_meta_app::app_error::AppError; +use databend_common_meta_app::app_error::TableEngineMismatch; use databend_common_meta_app::schema::CommitTableMetaReply; use databend_common_meta_app::schema::CommitTableMetaReq; use databend_common_meta_app::schema::CreateOption; @@ -117,15 +119,34 @@ impl TempTblMgr { let desc = Self::temp_table_desc(&name_ident.db_name, &name_ident.table_name); let engine = table_meta.engine.to_string(); let table_id = self.next_id; - let new_table = match (self.name_to_id.contains_key(&desc), create_option) { - (true, CreateOption::Create) => { + let new_table = match (self.name_to_id.get(&desc).copied(), create_option) { + (Some(_), CreateOption::Create) => { return Err(ErrorCode::TableAlreadyExists(format!( "Temporary table {} already exists", desc ))); } - (true, CreateOption::CreateIfNotExists) => false, - _ => { + (Some(_), CreateOption::CreateIfNotExists) => false, + (existing_id, _) => { + if let Some(existing_id) = existing_id { + let existing_table = self.id_to_table.get(&existing_id).ok_or_else(|| { + ErrorCode::Internal(format!( + "Got temporary table id {existing_id}, but its metadata was not found" + )) + })?; + if !existing_table + .meta + .engine + .eq_ignore_ascii_case(&table_meta.engine) + { + return Err(ErrorCode::from(AppError::from(TableEngineMismatch::new( + &name_ident.table_name, + &existing_table.meta.engine, + &table_meta.engine, + )))); + } + } + let desc = orphan_table_name .as_ref() .map(|o| Self::temp_table_desc(&name_ident.db_name, o)) diff --git a/tests/sqllogictests/suites/ee/06_ee_stream/06_0016_stream_replace_engine.test b/tests/sqllogictests/suites/ee/06_ee_stream/06_0016_stream_replace_engine.test new file mode 100644 index 0000000000000..443b8e2989b4b --- /dev/null +++ b/tests/sqllogictests/suites/ee/06_ee_stream/06_0016_stream_replace_engine.test @@ -0,0 +1,65 @@ +## Copyright 2023 Databend Cloud +## +## Licensed under the Elastic License, Version 2.0 (the "License"); +## you may not use this file except in compliance with the License. +## You may obtain a copy of the License at +## +## https://www.elastic.co/licensing/elastic-license +## +## Unless required by applicable law or agreed to in writing, software +## distributed under the License is distributed on an "AS IS" BASIS, +## WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +## See the License for the specific language governing permissions and +## limitations under the License. + +statement ok +DROP DATABASE IF EXISTS test_stream_replace_engine + +statement ok +CREATE DATABASE test_stream_replace_engine + +statement ok +USE test_stream_replace_engine + +statement ok +CREATE TABLE source(a INT) + +statement ok +CREATE STREAM target ON TABLE source + +statement error 1302 +CREATE OR REPLACE TABLE target(a INT) + +query TT +SHOW CREATE TABLE target +---- +target CREATE STREAM `target` ON TABLE `test_stream_replace_engine`.`source` + +statement ok +DROP STREAM target + +statement ok +CREATE TABLE target(a INT) + +statement ok +INSERT INTO target VALUES (1) + +statement error 1302 +CREATE OR REPLACE STREAM target ON TABLE source + +query I +SELECT * FROM target +---- +1 + +statement ok +CREATE OR REPLACE TABLE target(a STRING) + +statement ok +CREATE STREAM stream_target ON TABLE source + +statement ok +CREATE OR REPLACE STREAM stream_target ON TABLE source + +statement ok +DROP DATABASE test_stream_replace_engine diff --git a/tests/sqllogictests/suites/http_handler/temp_table/create_temp_tables.test b/tests/sqllogictests/suites/http_handler/temp_table/create_temp_tables.test index 66610b41b188d..fc2998d416dca 100644 --- a/tests/sqllogictests/suites/http_handler/temp_table/create_temp_tables.test +++ b/tests/sqllogictests/suites/http_handler/temp_table/create_temp_tables.test @@ -406,6 +406,12 @@ create or replace temp table IF NOT EXISTS replace_test(b int); statement ok create or replace temp table replace_test(b int); +statement error 1302 +create or replace temp table replace_test(c int) engine=memory; + +statement ok +select b from replace_test; + statement error 1065 select a from replace_test; From e6efbff5a4349ef9b3fea0c313b719c5652318f1 Mon Sep 17 00:00:00 2001 From: zhyass Date: Sun, 2 Aug 2026 00:22:15 +0800 Subject: [PATCH 2/3] fix review comments --- src/meta/api/src/api_impl/schema_api.rs | 18 ++++--- src/meta/api/src/api_impl/table_api.rs | 62 +++++++++++-------------- 2 files changed, 36 insertions(+), 44 deletions(-) diff --git a/src/meta/api/src/api_impl/schema_api.rs b/src/meta/api/src/api_impl/schema_api.rs index f8f2af33aff5a..f0014834d0d6f 100644 --- a/src/meta/api/src/api_impl/schema_api.rs +++ b/src/meta/api/src/api_impl/schema_api.rs @@ -21,7 +21,6 @@ use chrono::DateTime; use chrono::Utc; use databend_common_meta_app::app_error::AppError; use databend_common_meta_app::app_error::DropTableWithDropTime; -use databend_common_meta_app::app_error::TableEngineMismatch; use databend_common_meta_app::app_error::UndropTableAlreadyExists; use databend_common_meta_app::app_error::UndropTableHasNoHistory; use databend_common_meta_app::app_error::UndropTableRetentionGuard; @@ -178,7 +177,7 @@ pub async fn construct_drop_table_txn_operations( tenant: &Tenant, catalog_name: Option, table_id: u64, - expected_engine: Option<&str>, + preloaded_table_meta: Option>, db_id: u64, if_exists: bool, if_delete: bool, @@ -186,8 +185,14 @@ pub async fn construct_drop_table_txn_operations( ) -> Result<(u64, u64), KVAppError> { let tbid = TableId { table_id }; + // Reuse metadata already loaded by CREATE OR REPLACE; ordinary DROP callers + // do not have it and load it here. A preloaded value must belong to `table_id`. + let (tb_meta_seq, tb_meta) = if let Some(preloaded_table_meta) = preloaded_table_meta { + (preloaded_table_meta.seq, Some(preloaded_table_meta.data)) + } else { + kv_api.get_pb_seq_and_value(&tbid).await? + }; // Check if table exists. - let (tb_meta_seq, tb_meta) = kv_api.get_pb_seq_and_value(&tbid).await?; if tb_meta_seq == 0 { return Err(KVAppError::AppError(AppError::UnknownTableId( UnknownTableId::new(table_id, "drop_table_by_id failed to find valid tb_meta"), @@ -235,13 +240,6 @@ pub async fn construct_drop_table_txn_operations( ); let mut tb_meta = tb_meta.unwrap(); - if let Some(expected_engine) = expected_engine { - if !tb_meta.engine.eq_ignore_ascii_case(expected_engine) { - return Err(KVAppError::AppError( - TableEngineMismatch::new(table_name, &tb_meta.engine, expected_engine).into(), - )); - } - } // drop a table with drop_on time if tb_meta.drop_on.is_some() { return if if_exists { diff --git a/src/meta/api/src/api_impl/table_api.rs b/src/meta/api/src/api_impl/table_api.rs index f3851c86e56ed..156e443c41816 100644 --- a/src/meta/api/src/api_impl/table_api.rs +++ b/src/meta/api/src/api_impl/table_api.rs @@ -362,38 +362,32 @@ where }); } CreateOption::CreateOrReplace => { - if req.as_dropped { - // CTAS does not call construct_drop_table_txn_operations(), - // so validate its existing table here. - let existing_meta = - self.get_pb(&TableId::new(*id.data)).await?.ok_or_else(|| { - KVAppError::AppError(AppError::UnknownTableId( - UnknownTableId::new( - *id.data, - "create or replace failed to find existing table meta", - ), - )) - })?; - if !existing_meta - .engine - .eq_ignore_ascii_case(&req.table_meta.engine) - { - return Err(KVAppError::AppError( - TableEngineMismatch::new( - req.table_name(), - &existing_meta.engine, - &req.table_meta.engine, - ) - .into(), - )); - } - - // Guard the engine check against a concurrent metadata update. - txn.condition.push(txn_cond_seq( - &TableId::new(*id.data), - Eq, - existing_meta.seq, + // CREATE OR REPLACE must not replace an existing table with a + // different engine. + let existing_table_meta = self + .get_pb(&TableId::new(*id.data)) + .await? + .ok_or_else(|| { + KVAppError::AppError(AppError::UnknownTableId(UnknownTableId::new( + *id.data, + "create or replace failed to find existing table meta", + ))) + })?; + if !existing_table_meta + .engine + .eq_ignore_ascii_case(&req.table_meta.engine) + { + return Err(KVAppError::AppError( + TableEngineMismatch::new( + req.table_name(), + &existing_table_meta.engine, + &req.table_meta.engine, + ) + .into(), )); + } + + if req.as_dropped { // If the table is being created as a dropped table, we do not // need to combine with drop_table_txn operations, just return // the sequence number associated with the value part of @@ -401,15 +395,15 @@ where SeqV::new(id.seq, *id.data) } else { - // The drop helper validates the engine against the metadata - // sequence used by the replacement transaction. + // Reuse the metadata loaded for the engine check when marking the + // existing table as dropped. let (seq, id) = construct_drop_table_txn_operations( self, req.name_ident.table_name.clone(), &req.name_ident.tenant, req.catalog_name.clone(), *id.data, - Some(&req.table_meta.engine), + Some(existing_table_meta), *seq_db_id.data, true, false, From 38ba68e2c6718d9da38f56663d692d5a6e36ccff Mon Sep 17 00:00:00 2001 From: zhyass Date: Sun, 2 Aug 2026 14:33:10 +0800 Subject: [PATCH 3/3] fix review comments fix fix --- Cargo.lock | 2 - src/meta/api/src/api_impl/schema_api.rs | 30 +-- src/meta/api/src/api_impl/table_api.rs | 75 +++++-- src/meta/app/src/app_error.rs | 8 + .../src/schema_api_test_suite.rs | 98 ++++++++- .../storages/common/session/src/temp_table.rs | 189 +++++++++++++++--- src/query/storages/paimon/Cargo.toml | 2 - .../06_0016_stream_replace_engine.test | 3 + 8 files changed, 327 insertions(+), 80 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 4bd77dd53fbf2..41229298b0dc7 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -5239,13 +5239,11 @@ dependencies = [ "databend-common-meta-app", "databend-common-pipeline", "databend-common-pipeline-transforms", - "databend-common-storage", "databend-common-users", "databend-meta-client 260205.13.1", "databend-storages-common-table-meta", "databend_educe", "futures", - "log", "paimon", "pretty_assertions", "serde", diff --git a/src/meta/api/src/api_impl/schema_api.rs b/src/meta/api/src/api_impl/schema_api.rs index f0014834d0d6f..687400672d86a 100644 --- a/src/meta/api/src/api_impl/schema_api.rs +++ b/src/meta/api/src/api_impl/schema_api.rs @@ -171,33 +171,26 @@ pub async fn get_history_table_metas( Ok(tb_metas) } -pub async fn construct_drop_table_txn_operations( +pub(crate) struct VersionedTable { + pub id: TableId, + pub meta: SeqV, +} + +pub(crate) async fn construct_drop_table_txn_operations( kv_api: &(impl kvapi::KVApi + ?Sized), table_name: String, tenant: &Tenant, catalog_name: Option, - table_id: u64, - preloaded_table_meta: Option>, + existing: VersionedTable, db_id: u64, if_exists: bool, if_delete: bool, txn: &mut TxnRequest, ) -> Result<(u64, u64), KVAppError> { - let tbid = TableId { table_id }; - - // Reuse metadata already loaded by CREATE OR REPLACE; ordinary DROP callers - // do not have it and load it here. A preloaded value must belong to `table_id`. - let (tb_meta_seq, tb_meta) = if let Some(preloaded_table_meta) = preloaded_table_meta { - (preloaded_table_meta.seq, Some(preloaded_table_meta.data)) - } else { - kv_api.get_pb_seq_and_value(&tbid).await? - }; - // Check if table exists. - if tb_meta_seq == 0 { - return Err(KVAppError::AppError(AppError::UnknownTableId( - UnknownTableId::new(table_id, "drop_table_by_id failed to find valid tb_meta"), - ))); - } + let VersionedTable { id: tbid, meta } = existing; + let table_id = tbid.table_id; + let tb_meta_seq = meta.seq; + let mut tb_meta = meta.data; // Get db name, tenant name and related info for tx. let table_id_to_name = TableIdToName { table_id }; @@ -239,7 +232,6 @@ pub async fn construct_drop_table_txn_operations( "drop table by id" ); - let mut tb_meta = tb_meta.unwrap(); // drop a table with drop_on time if tb_meta.drop_on.is_some() { return if if_exists { diff --git a/src/meta/api/src/api_impl/table_api.rs b/src/meta/api/src/api_impl/table_api.rs index 156e443c41816..0ed1962ad18c3 100644 --- a/src/meta/api/src/api_impl/table_api.rs +++ b/src/meta/api/src/api_impl/table_api.rs @@ -129,6 +129,7 @@ use super::database_api::DatabaseApi; use super::database_util::get_db_or_err; use super::garbage_collection_api::ORPHAN_POSTFIX; use super::garbage_collection_api::get_history_tables_for_gc; +use super::schema_api::VersionedTable; use super::schema_api::build_upsert_table_deduplicated_label; use super::schema_api::construct_drop_table_txn_operations; use super::schema_api::get_db_by_id_or_err; @@ -373,19 +374,12 @@ where "create or replace failed to find existing table meta", ))) })?; - if !existing_table_meta - .engine - .eq_ignore_ascii_case(&req.table_meta.engine) - { - return Err(KVAppError::AppError( - TableEngineMismatch::new( - req.table_name(), - &existing_table_meta.engine, - &req.table_meta.engine, - ) - .into(), - )); - } + TableEngineMismatch::ensure( + req.table_name(), + &existing_table_meta.engine, + &req.table_meta.engine, + ) + .map_err(|e| KVAppError::AppError(e.into()))?; if req.as_dropped { // If the table is being created as a dropped table, we do not @@ -395,15 +389,16 @@ where SeqV::new(id.seq, *id.data) } else { - // Reuse the metadata loaded for the engine check when marking the - // existing table as dropped. + let existing = VersionedTable { + id: TableId::new(*id.data), + meta: existing_table_meta, + }; let (seq, id) = construct_drop_table_txn_operations( self, req.name_ident.table_name.clone(), &req.name_ident.tenant, req.catalog_name.clone(), - *id.data, - Some(existing_table_meta), + existing, *seq_db_id.data, true, false, @@ -634,13 +629,23 @@ where let mut txn = TxnRequest::default(); + let tbid = TableId::new(table_id); + let table_meta = self.get_pb(&tbid).await?.ok_or_else(|| { + KVAppError::AppError(AppError::UnknownTableId(UnknownTableId::new( + table_id, + "drop_table_by_id failed to find valid tb_meta", + ))) + })?; + let existing = VersionedTable { + id: tbid, + meta: table_meta, + }; let opt = construct_drop_table_txn_operations( self, req.table_name.clone(), &req.tenant, None, - table_id, - None, + existing, req.db_id, req.if_exists, true, @@ -1192,7 +1197,35 @@ where { // reset drop on time - let mut tb_meta = tb_meta.unwrap(); + let mut txn_req = TxnRequest::default(); + let mut tb_meta = tb_meta.ok_or_else(|| { + KVAppError::AppError(AppError::UnknownTableId(UnknownTableId::new( + table_id, + "commit_table_meta", + ))) + })?; + + if let Some(prev_table_id) = req.prev_table_id { + let prev_tbid = TableId::new(prev_table_id); + let previous = self.get_pb(&prev_tbid).await?.ok_or_else(|| { + KVAppError::AppError(AppError::UnknownTableId(UnknownTableId::new( + prev_table_id, + "commit_table_meta previous table", + ))) + })?; + + TableEngineMismatch::ensure( + req.name_ident.table_name.as_str(), + &previous.engine, + &tb_meta.engine, + ) + .map_err(|e| KVAppError::AppError(e.into()))?; + + txn_req + .condition + .push(txn_cond_seq(&prev_tbid, Eq, previous.seq)); + } + // undrop a table with no drop_on time if tb_meta.drop_on.is_none() { return Err(KVAppError::AppError(AppError::UndropTableWithNoDropTime( @@ -1201,8 +1234,6 @@ where } tb_meta.drop_on = None; - let mut txn_req = TxnRequest::default(); - txn_req.condition.extend([ // db has not to change, i.e., no new table is created. // Renaming db is OK and does not affect the seq of db_meta. diff --git a/src/meta/app/src/app_error.rs b/src/meta/app/src/app_error.rs index 1272bd28f1d4a..c25ff873e34d1 100644 --- a/src/meta/app/src/app_error.rs +++ b/src/meta/app/src/app_error.rs @@ -217,6 +217,14 @@ impl TableEngineMismatch { new_engine: new_engine.into(), } } + + pub fn ensure(table_name: &str, existing_engine: &str, new_engine: &str) -> Result<(), Self> { + if existing_engine.eq_ignore_ascii_case(new_engine) { + return Ok(()); + } + + Err(Self::new(table_name, existing_engine, new_engine)) + } } #[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)] diff --git a/src/meta/schema-api-test-suite/src/schema_api_test_suite.rs b/src/meta/schema-api-test-suite/src/schema_api_test_suite.rs index e14f2eade4198..10cda8c229d56 100644 --- a/src/meta/schema-api-test-suite/src/schema_api_test_suite.rs +++ b/src/meta/schema-api-test-suite/src/schema_api_test_suite.rs @@ -51,6 +51,7 @@ use databend_common_meta_api::serialize_struct; use databend_common_meta_api::util::IdempotentKVTxnSender; use databend_common_meta_app::KeyWithTenant; use databend_common_meta_app::app_error::AppError; +use databend_common_meta_app::app_error::TableEngineMismatch; use databend_common_meta_app::data_mask::CreateDatamaskReq; use databend_common_meta_app::data_mask::DataMaskNameIdent; use databend_common_meta_app::data_mask::DatamaskMeta; @@ -318,6 +319,8 @@ impl SchemaApiTestSuite { + 'static, { self.table_commit_table_meta(&b.build().await).await?; + self.table_commit_table_meta_engine_mismatch(&b.build().await) + .await?; self.concurrent_commit_table_meta(b.clone()).await?; self.database_drop_out_of_retention_time_history(&b.build().await) .await?; @@ -2634,11 +2637,11 @@ impl SchemaApiTestSuite { table_properties: None, table_partition: None, }; + let expected = KVAppError::AppError(AppError::from(TableEngineMismatch::new( + table, "JSON", "STREAM", + ))); let err = mt.create_table(mismatched_req.clone()).await.unwrap_err(); - assert!(matches!( - err, - KVAppError::AppError(AppError::TableEngineMismatch(_)) - )); + assert_eq!(err, expected); assert_eq!(table_id, get_kv_u64_data(mt, &key_dbid_tbname).await?); assert_eq!("JSON", util.get_table_by_name(table).await?.meta.engine); @@ -2648,10 +2651,7 @@ impl SchemaApiTestSuite { as_dropped_req.as_dropped = true; as_dropped_req.table_meta.drop_on = Some(Utc::now()); let err = mt.create_table(as_dropped_req).await.unwrap_err(); - assert!(matches!( - err, - KVAppError::AppError(AppError::TableEngineMismatch(_)) - )); + assert_eq!(err, expected); assert_eq!(table_id, get_kv_u64_data(mt, &key_dbid_tbname).await?); } @@ -6523,6 +6523,88 @@ impl SchemaApiTestSuite { Ok(()) } + async fn table_commit_table_meta_engine_mismatch(&self, mt: &MT) -> anyhow::Result<()> + where MT: kvapi::KVApi + DatabaseApi + TableApi { + let tenant_name = "table_commit_table_meta_engine_mismatch"; + let db_name = "db"; + let table_name = "table"; + let mut util = DbTableHarness::new(mt, tenant_name, db_name, table_name, "JSON"); + util.create_db().await?; + let (previous_table_id, previous_meta) = util.create_table().await?; + + let mut replacement_meta = previous_meta.clone(); + replacement_meta.drop_on = Some(Utc::now()); + let create_req = CreateTableReq { + create_option: CreateOption::CreateOrReplace, + catalog_name: Some("default".to_string()), + name_ident: TableNameIdent { + tenant: util.tenant(), + db_name: db_name.to_string(), + table_name: table_name.to_string(), + }, + table_meta: replacement_meta, + as_dropped: true, + materialized_view: None, + table_properties: None, + table_partition: None, + }; + let create_reply = mt.create_table(create_req.clone()).await?; + assert_eq!(create_reply.prev_table_id, Some(previous_table_id)); + + let previous_tbid = TableId::new(previous_table_id); + let mut changed_previous_meta = previous_meta; + changed_previous_meta.engine = "STREAM".to_string(); + upsert_test_data(mt, &previous_tbid, serialize_struct(&changed_previous_meta)).await?; + + let dbid_tbname = DBIdTableName { + db_id: create_reply.db_id, + table_name: table_name.to_string(), + }; + let history_ident = TableIdHistoryIdent { + database_id: create_reply.db_id, + table_name: table_name.to_string(), + }; + let orphan_history_ident = TableIdHistoryIdent { + database_id: create_reply.db_id, + table_name: create_reply.orphan_table_name.clone().unwrap(), + }; + let replacement_tbid = TableId::new(create_reply.table_id); + + let name_before = mt.get_pb(&dbid_tbname).await?; + let history_before = mt.get_pb(&history_ident).await?; + let orphan_history_before = mt.get_pb(&orphan_history_ident).await?; + let previous_before = mt.get_pb(&previous_tbid).await?; + let replacement_before = mt.get_pb(&replacement_tbid).await?; + + let err = mt + .commit_table_meta(CommitTableMetaReq { + name_ident: create_req.name_ident, + db_id: create_reply.db_id, + table_id: create_reply.table_id, + prev_table_id: create_reply.prev_table_id, + orphan_table_name: create_reply.orphan_table_name, + }) + .await + .unwrap_err(); + assert_eq!( + err, + KVAppError::AppError(AppError::from(TableEngineMismatch::new( + table_name, "STREAM", "JSON" + ))) + ); + + assert_eq!(mt.get_pb(&dbid_tbname).await?, name_before); + assert_eq!(mt.get_pb(&history_ident).await?, history_before); + assert_eq!( + mt.get_pb(&orphan_history_ident).await?, + orphan_history_before + ); + assert_eq!(mt.get_pb(&previous_tbid).await?, previous_before); + assert_eq!(mt.get_pb(&replacement_tbid).await?, replacement_before); + + Ok(()) + } + async fn concurrent_commit_table_meta< B: kvapi::ApiBuilder, MT: kvapi::KVApi + DatabaseApi + TableApi + 'static, diff --git a/src/query/storages/common/session/src/temp_table.rs b/src/query/storages/common/session/src/temp_table.rs index 23485f5afe432..e3bd246aa4c1a 100644 --- a/src/query/storages/common/session/src/temp_table.rs +++ b/src/query/storages/common/session/src/temp_table.rs @@ -117,9 +117,10 @@ impl TempTblMgr { let db_id = db_id.parse::()?; let desc = Self::temp_table_desc(&name_ident.db_name, &name_ident.table_name); + let existing_id = self.name_to_id.get(&desc).copied(); let engine = table_meta.engine.to_string(); let table_id = self.next_id; - let new_table = match (self.name_to_id.get(&desc).copied(), create_option) { + let new_table = match (existing_id, create_option) { (Some(_), CreateOption::Create) => { return Err(ErrorCode::TableAlreadyExists(format!( "Temporary table {} already exists", @@ -134,17 +135,12 @@ impl TempTblMgr { "Got temporary table id {existing_id}, but its metadata was not found" )) })?; - if !existing_table - .meta - .engine - .eq_ignore_ascii_case(&table_meta.engine) - { - return Err(ErrorCode::from(AppError::from(TableEngineMismatch::new( - &name_ident.table_name, - &existing_table.meta.engine, - &table_meta.engine, - )))); - } + TableEngineMismatch::ensure( + &name_ident.table_name, + &existing_table.meta.engine, + &table_meta.engine, + ) + .map_err(|e| ErrorCode::from(AppError::from(e)))?; } let desc = orphan_table_name @@ -175,7 +171,7 @@ impl TempTblMgr { table_id_seq: Some(0), db_id, new_table, - prev_table_id: None, + prev_table_id: as_dropped.then_some(existing_id).flatten(), orphan_table_name, }) } @@ -186,22 +182,44 @@ impl TempTblMgr { req.orphan_table_name.as_ref().unwrap(), ); let desc = Self::temp_table_desc(&req.name_ident.db_name, &req.name_ident.table_name); - match self.name_to_id.remove(&orphan_desc) { - Some(id) => { - if let Some(old_id) = self.name_to_id.insert(desc, id) { - self.id_to_table.remove(&old_id); - } - let table = self.id_to_table.get_mut(&id).unwrap(); - table.db_name = req.name_ident.db_name.clone(); - table.table_name = req.name_ident.table_name.clone(); + let orphan_id = self + .name_to_id + .get(&orphan_desc) + .copied() + .ok_or_else(|| ErrorCode::UnknownTable(orphan_desc.clone()))?; + + let current_id = self.name_to_id.get(&desc).copied(); + if current_id != req.prev_table_id { + return Err(ErrorCode::TableVersionMismatched(format!( + "Temporary table {desc} changed while CTAS was running" + ))); + } - Ok(CommitTableMetaReply {}) - } - None => Err(ErrorCode::UnknownTable(format!( - "Temporary table {}.{} not found", - req.name_ident.db_name, req.name_ident.table_name - ))), + if let Some(current_id) = current_id { + let current = self.id_to_table.get(¤t_id).ok_or_else(|| { + ErrorCode::Internal(format!("Missing temporary table metadata for {current_id}")) + })?; + let replacement = self.id_to_table.get(&orphan_id).ok_or_else(|| { + ErrorCode::Internal(format!("Missing temporary table metadata for {orphan_id}")) + })?; + + TableEngineMismatch::ensure( + &req.name_ident.table_name, + ¤t.meta.engine, + &replacement.meta.engine, + ) + .map_err(|e| ErrorCode::from(AppError::from(e)))?; + } + + self.name_to_id.remove(&orphan_desc); + if let Some(old_id) = self.name_to_id.insert(desc, orphan_id) { + self.id_to_table.remove(&old_id); } + let table = self.id_to_table.get_mut(&orphan_id).unwrap(); + table.db_name = req.name_ident.db_name.clone(); + table.table_name = req.name_ident.table_name.clone(); + + Ok(CommitTableMetaReply {}) } pub fn rename_table(&mut self, req: &RenameTableReq) -> Result> { @@ -534,3 +552,120 @@ pub async fn drop_all_temp_tables( } pub type TempTblMgrRef = Arc>; + +#[cfg(test)] +mod tests { + use databend_common_meta_app::schema::TableNameIdent; + use databend_common_meta_app::tenant::Tenant; + + use super::*; + + fn create_table_req( + table_name: &str, + engine: &str, + create_option: CreateOption, + as_dropped: bool, + ) -> CreateTableReq { + let mut table_meta = TableMeta { + engine: engine.to_string(), + ..Default::default() + }; + table_meta + .options + .insert(OPT_KEY_DATABASE_ID.to_string(), "1".to_string()); + + CreateTableReq { + create_option, + catalog_name: None, + name_ident: TableNameIdent { + tenant: Tenant::new_literal("tenant"), + db_name: "db".to_string(), + table_name: table_name.to_string(), + }, + table_meta, + as_dropped, + materialized_view: None, + table_properties: None, + table_partition: None, + } + } + + #[test] + fn test_commit_cross_engine_ctas_preserves_temporary_tables() { + let mut mgr = TempTblMgr { + name_to_id: HashMap::new(), + id_to_table: HashMap::new(), + next_id: TEMP_TBL_ID_BEGIN, + }; + let table_name = "t"; + let desc = TempTblMgr::temp_table_desc("db", table_name); + + let visible = mgr + .create_table( + create_table_req(table_name, "FUSE", CreateOption::Create, false), + "session".to_string(), + ) + .unwrap(); + let replacement = mgr + .create_table( + create_table_req(table_name, "FUSE", CreateOption::CreateOrReplace, true), + "session".to_string(), + ) + .unwrap(); + assert_eq!(replacement.prev_table_id, Some(visible.table_id)); + + mgr.name_to_id.remove(&desc); + mgr.id_to_table.remove(&visible.table_id); + let current = mgr + .create_table( + create_table_req(table_name, "MEMORY", CreateOption::Create, false), + "session".to_string(), + ) + .unwrap(); + + let orphan_desc = + TempTblMgr::temp_table_desc("db", replacement.orphan_table_name.as_ref().unwrap()); + let current_meta = mgr.id_to_table.get(¤t.table_id).unwrap().meta.clone(); + let replacement_meta = mgr + .id_to_table + .get(&replacement.table_id) + .unwrap() + .meta + .clone(); + + let err = mgr + .commit_table_meta(&CommitTableMetaReq { + name_ident: TableNameIdent { + tenant: Tenant::new_literal("tenant"), + db_name: "db".to_string(), + table_name: table_name.to_string(), + }, + db_id: replacement.db_id, + table_id: replacement.table_id, + prev_table_id: replacement.prev_table_id, + orphan_table_name: replacement.orphan_table_name, + }) + .unwrap_err(); + let expected = ErrorCode::TableVersionMismatched(format!( + "Temporary table {desc} changed while CTAS was running" + )); + assert_eq!( + (err.code(), err.message()), + (expected.code(), expected.message()) + ); + + assert_eq!(mgr.name_to_id.get(&desc), Some(¤t.table_id)); + assert_eq!( + mgr.name_to_id.get(&orphan_desc), + Some(&replacement.table_id) + ); + assert_eq!( + mgr.id_to_table.get(¤t.table_id).unwrap().meta, + current_meta + ); + assert_eq!( + mgr.id_to_table.get(&replacement.table_id).unwrap().meta, + replacement_meta + ); + } +} diff --git a/src/query/storages/paimon/Cargo.toml b/src/query/storages/paimon/Cargo.toml index 0471c340e05f0..4556383a7aeaa 100644 --- a/src/query/storages/paimon/Cargo.toml +++ b/src/query/storages/paimon/Cargo.toml @@ -21,13 +21,11 @@ databend-common-expression = { workspace = true } databend-common-meta-app = { workspace = true } databend-common-pipeline = { workspace = true } databend-common-pipeline-transforms = { workspace = true } -databend-common-storage = { workspace = true } databend-common-users = { workspace = true } databend-meta-client = { workspace = true } databend-storages-common-table-meta = { workspace = true } educe = { workspace = true } futures = { workspace = true } -log = { workspace = true } paimon = { workspace = true } serde = { workspace = true } serde_json = { workspace = true } diff --git a/tests/sqllogictests/suites/ee/06_ee_stream/06_0016_stream_replace_engine.test b/tests/sqllogictests/suites/ee/06_ee_stream/06_0016_stream_replace_engine.test index 443b8e2989b4b..f0380f952b1ab 100644 --- a/tests/sqllogictests/suites/ee/06_ee_stream/06_0016_stream_replace_engine.test +++ b/tests/sqllogictests/suites/ee/06_ee_stream/06_0016_stream_replace_engine.test @@ -30,6 +30,9 @@ CREATE STREAM target ON TABLE source statement error 1302 CREATE OR REPLACE TABLE target(a INT) +statement error 1302 +CREATE OR REPLACE TABLE target AS SELECT 1 AS a + query TT SHOW CREATE TABLE target ----