Skip to content

Commit e1d18f0

Browse files
committed
fix(meta): disallow replacing tables with different engines
1 parent ca7d1d9 commit e1d18f0

7 files changed

Lines changed: 241 additions & 4 deletions

File tree

src/meta/api/src/api_impl/schema_api.rs

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@ use chrono::DateTime;
2121
use chrono::Utc;
2222
use databend_common_meta_app::app_error::AppError;
2323
use databend_common_meta_app::app_error::DropTableWithDropTime;
24+
use databend_common_meta_app::app_error::TableEngineMismatch;
2425
use databend_common_meta_app::app_error::UndropTableAlreadyExists;
2526
use databend_common_meta_app::app_error::UndropTableHasNoHistory;
2627
use databend_common_meta_app::app_error::UndropTableRetentionGuard;
@@ -177,6 +178,7 @@ pub async fn construct_drop_table_txn_operations(
177178
tenant: &Tenant,
178179
catalog_name: Option<String>,
179180
table_id: u64,
181+
expected_engine: Option<&str>,
180182
db_id: u64,
181183
if_exists: bool,
182184
if_delete: bool,
@@ -233,6 +235,13 @@ pub async fn construct_drop_table_txn_operations(
233235
);
234236

235237
let mut tb_meta = tb_meta.unwrap();
238+
if let Some(expected_engine) = expected_engine {
239+
if !tb_meta.engine.eq_ignore_ascii_case(expected_engine) {
240+
return Err(KVAppError::AppError(
241+
TableEngineMismatch::new(table_name, &tb_meta.engine, expected_engine).into(),
242+
));
243+
}
244+
}
236245
// drop a table with drop_on time
237246
if tb_meta.drop_on.is_some() {
238247
return if if_exists {

src/meta/api/src/api_impl/table_api.rs

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -34,6 +34,7 @@ use databend_common_meta_app::app_error::MultiStmtTxnCommitFailed;
3434
use databend_common_meta_app::app_error::StreamAlreadyExists;
3535
use databend_common_meta_app::app_error::StreamVersionMismatched;
3636
use databend_common_meta_app::app_error::TableAlreadyExists;
37+
use databend_common_meta_app::app_error::TableEngineMismatch;
3738
use databend_common_meta_app::app_error::TableSnapshotExpired;
3839
use databend_common_meta_app::app_error::TableVersionMismatched;
3940
use databend_common_meta_app::app_error::UndropTableHasNoHistory;
@@ -362,19 +363,53 @@ where
362363
}
363364
CreateOption::CreateOrReplace => {
364365
if req.as_dropped {
366+
// CTAS does not call construct_drop_table_txn_operations(),
367+
// so validate its existing table here.
368+
let existing_meta =
369+
self.get_pb(&TableId::new(*id.data)).await?.ok_or_else(|| {
370+
KVAppError::AppError(AppError::UnknownTableId(
371+
UnknownTableId::new(
372+
*id.data,
373+
"create or replace failed to find existing table meta",
374+
),
375+
))
376+
})?;
377+
if !existing_meta
378+
.engine
379+
.eq_ignore_ascii_case(&req.table_meta.engine)
380+
{
381+
return Err(KVAppError::AppError(
382+
TableEngineMismatch::new(
383+
req.table_name(),
384+
&existing_meta.engine,
385+
&req.table_meta.engine,
386+
)
387+
.into(),
388+
));
389+
}
390+
391+
// Guard the engine check against a concurrent metadata update.
392+
txn.condition.push(txn_cond_seq(
393+
&TableId::new(*id.data),
394+
Eq,
395+
existing_meta.seq,
396+
));
365397
// If the table is being created as a dropped table, we do not
366398
// need to combine with drop_table_txn operations, just return
367399
// the sequence number associated with the value part of
368400
// the key-value pair (key_dbid_tbname, table_id).
369401

370402
SeqV::new(id.seq, *id.data)
371403
} else {
404+
// The drop helper validates the engine against the metadata
405+
// sequence used by the replacement transaction.
372406
let (seq, id) = construct_drop_table_txn_operations(
373407
self,
374408
req.name_ident.table_name.clone(),
375409
&req.name_ident.tenant,
376410
req.catalog_name.clone(),
377411
*id.data,
412+
Some(&req.table_meta.engine),
378413
*seq_db_id.data,
379414
true,
380415
false,
@@ -611,6 +646,7 @@ where
611646
&req.tenant,
612647
None,
613648
table_id,
649+
None,
614650
req.db_id,
615651
req.if_exists,
616652
true,

src/meta/app/src/app_error.rs

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -195,6 +195,30 @@ impl TableAlreadyExists {
195195
}
196196
}
197197

198+
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
199+
#[error(
200+
"Cannot replace '{table_name}': existing table uses engine {existing_engine}, but the new table uses engine {new_engine}"
201+
)]
202+
pub struct TableEngineMismatch {
203+
table_name: String,
204+
existing_engine: String,
205+
new_engine: String,
206+
}
207+
208+
impl TableEngineMismatch {
209+
pub fn new(
210+
table_name: impl Into<String>,
211+
existing_engine: impl Into<String>,
212+
new_engine: impl Into<String>,
213+
) -> Self {
214+
Self {
215+
table_name: table_name.into(),
216+
existing_engine: existing_engine.into(),
217+
new_engine: new_engine.into(),
218+
}
219+
}
220+
}
221+
198222
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
199223
#[error("ViewAlreadyExists: {view_name} while {context}")]
200224
pub struct ViewAlreadyExists {
@@ -895,6 +919,9 @@ pub enum AppError {
895919
#[error(transparent)]
896920
TableAlreadyExists(#[from] TableAlreadyExists),
897921

922+
#[error(transparent)]
923+
TableEngineMismatch(#[from] TableEngineMismatch),
924+
898925
#[error(transparent)]
899926
ViewAlreadyExists(#[from] ViewAlreadyExists),
900927

@@ -1199,6 +1226,8 @@ impl AppErrorMessage for TableAlreadyExists {
11991226
}
12001227
}
12011228

1229+
impl AppErrorMessage for TableEngineMismatch {}
1230+
12021231
impl AppErrorMessage for ViewAlreadyExists {
12031232
fn message(&self) -> String {
12041233
format!("'{}' as view Already Exists", self.view_name)
@@ -1412,6 +1441,7 @@ impl From<AppError> for ErrorCode {
14121441
ErrorCode::MaterializedViewAlreadyExists(err.message())
14131442
}
14141443
AppError::TableAlreadyExists(err) => ErrorCode::TableAlreadyExists(err.message()),
1444+
AppError::TableEngineMismatch(err) => ErrorCode::TableEngineNotSupported(err.message()),
14151445
AppError::ViewAlreadyExists(err) => ErrorCode::ViewAlreadyExists(err.message()),
14161446
AppError::CreateTableWithDropTime(err) => {
14171447
ErrorCode::CreateTableWithDropTime(err.message())

src/meta/schema-api-test-suite/src/schema_api_test_suite.rs

Lines changed: 70 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2583,6 +2583,76 @@ impl SchemaApiTestSuite {
25832583
let ret_table_name_ident: DBIdTableName =
25842584
get_kv_data(mt, &key_table_id_to_name).await?;
25852585
assert_eq!(ret_table_name_ident, key_dbid_tbname);
2586+
2587+
// Engine names are case-insensitive.
2588+
let (_, replaced_meta) = util
2589+
.create_table_with(
2590+
|mut meta| {
2591+
meta.engine = "json".to_string();
2592+
meta
2593+
},
2594+
|mut req| {
2595+
req.create_option = CreateOption::CreateOrReplace;
2596+
req.name_ident.table_name = table.to_string();
2597+
req
2598+
},
2599+
)
2600+
.await?;
2601+
assert_eq!("json", replaced_meta.engine);
2602+
2603+
// Restore the upper-case engine for the mismatch cases below.
2604+
let (table_id, _) = util
2605+
.create_table_with(
2606+
|mut meta| {
2607+
meta.engine = "JSON".to_string();
2608+
meta
2609+
},
2610+
|mut req| {
2611+
req.create_option = CreateOption::CreateOrReplace;
2612+
req.name_ident.table_name = table.to_string();
2613+
req
2614+
},
2615+
)
2616+
.await?;
2617+
2618+
// Replacing a table with a different engine must fail without changing
2619+
// the existing table or its name mapping.
2620+
let mismatched_req = CreateTableReq {
2621+
create_option: CreateOption::CreateOrReplace,
2622+
catalog_name: Some("default".to_string()),
2623+
name_ident: TableNameIdent {
2624+
tenant: tenant.clone(),
2625+
db_name: db_name.to_string(),
2626+
table_name: table.to_string(),
2627+
},
2628+
table_meta: TableMeta {
2629+
engine: "STREAM".to_string(),
2630+
..Default::default()
2631+
},
2632+
as_dropped: false,
2633+
materialized_view: None,
2634+
table_properties: None,
2635+
table_partition: None,
2636+
};
2637+
let err = mt.create_table(mismatched_req.clone()).await.unwrap_err();
2638+
assert!(matches!(
2639+
err,
2640+
KVAppError::AppError(AppError::TableEngineMismatch(_))
2641+
));
2642+
assert_eq!(table_id, get_kv_u64_data(mt, &key_dbid_tbname).await?);
2643+
assert_eq!("JSON", util.get_table_by_name(table).await?.meta.engine);
2644+
2645+
// The same validation applies when CTAS creates its replacement as a
2646+
// hidden dropped table.
2647+
let mut as_dropped_req = mismatched_req;
2648+
as_dropped_req.as_dropped = true;
2649+
as_dropped_req.table_meta.drop_on = Some(Utc::now());
2650+
let err = mt.create_table(as_dropped_req).await.unwrap_err();
2651+
assert!(matches!(
2652+
err,
2653+
KVAppError::AppError(AppError::TableEngineMismatch(_))
2654+
));
2655+
assert_eq!(table_id, get_kv_u64_data(mt, &key_dbid_tbname).await?);
25862656
}
25872657

25882658
{

src/query/storages/common/session/src/temp_table.rs

Lines changed: 25 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,8 @@ use std::sync::Arc;
1919

2020
use databend_common_exception::ErrorCode;
2121
use databend_common_exception::Result;
22+
use databend_common_meta_app::app_error::AppError;
23+
use databend_common_meta_app::app_error::TableEngineMismatch;
2224
use databend_common_meta_app::schema::CommitTableMetaReply;
2325
use databend_common_meta_app::schema::CommitTableMetaReq;
2426
use databend_common_meta_app::schema::CreateOption;
@@ -117,15 +119,34 @@ impl TempTblMgr {
117119
let desc = Self::temp_table_desc(&name_ident.db_name, &name_ident.table_name);
118120
let engine = table_meta.engine.to_string();
119121
let table_id = self.next_id;
120-
let new_table = match (self.name_to_id.contains_key(&desc), create_option) {
121-
(true, CreateOption::Create) => {
122+
let new_table = match (self.name_to_id.get(&desc).copied(), create_option) {
123+
(Some(_), CreateOption::Create) => {
122124
return Err(ErrorCode::TableAlreadyExists(format!(
123125
"Temporary table {} already exists",
124126
desc
125127
)));
126128
}
127-
(true, CreateOption::CreateIfNotExists) => false,
128-
_ => {
129+
(Some(_), CreateOption::CreateIfNotExists) => false,
130+
(existing_id, _) => {
131+
if let Some(existing_id) = existing_id {
132+
let existing_table = self.id_to_table.get(&existing_id).ok_or_else(|| {
133+
ErrorCode::Internal(format!(
134+
"Got temporary table id {existing_id}, but its metadata was not found"
135+
))
136+
})?;
137+
if !existing_table
138+
.meta
139+
.engine
140+
.eq_ignore_ascii_case(&table_meta.engine)
141+
{
142+
return Err(ErrorCode::from(AppError::from(TableEngineMismatch::new(
143+
&name_ident.table_name,
144+
&existing_table.meta.engine,
145+
&table_meta.engine,
146+
))));
147+
}
148+
}
149+
129150
let desc = orphan_table_name
130151
.as_ref()
131152
.map(|o| Self::temp_table_desc(&name_ident.db_name, o))
Lines changed: 65 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,65 @@
1+
## Copyright 2023 Databend Cloud
2+
##
3+
## Licensed under the Elastic 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+
## https://www.elastic.co/licensing/elastic-license
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+
15+
statement ok
16+
DROP DATABASE IF EXISTS test_stream_replace_engine
17+
18+
statement ok
19+
CREATE DATABASE test_stream_replace_engine
20+
21+
statement ok
22+
USE test_stream_replace_engine
23+
24+
statement ok
25+
CREATE TABLE source(a INT)
26+
27+
statement ok
28+
CREATE STREAM target ON TABLE source
29+
30+
statement error 1302
31+
CREATE OR REPLACE TABLE target(a INT)
32+
33+
query TT
34+
SHOW CREATE TABLE target
35+
----
36+
target CREATE STREAM `target` ON TABLE `test_stream_replace_engine`.`source`
37+
38+
statement ok
39+
DROP STREAM target
40+
41+
statement ok
42+
CREATE TABLE target(a INT)
43+
44+
statement ok
45+
INSERT INTO target VALUES (1)
46+
47+
statement error 1302
48+
CREATE OR REPLACE STREAM target ON TABLE source
49+
50+
query I
51+
SELECT * FROM target
52+
----
53+
1
54+
55+
statement ok
56+
CREATE OR REPLACE TABLE target(a STRING)
57+
58+
statement ok
59+
CREATE STREAM stream_target ON TABLE source
60+
61+
statement ok
62+
CREATE OR REPLACE STREAM stream_target ON TABLE source
63+
64+
statement ok
65+
DROP DATABASE test_stream_replace_engine

tests/sqllogictests/suites/http_handler/temp_table/create_temp_tables.test

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -406,6 +406,12 @@ create or replace temp table IF NOT EXISTS replace_test(b int);
406406
statement ok
407407
create or replace temp table replace_test(b int);
408408

409+
statement error 1302
410+
create or replace temp table replace_test(c int) engine=memory;
411+
412+
statement ok
413+
select b from replace_test;
414+
409415
statement error 1065
410416
select a from replace_test;
411417

0 commit comments

Comments
 (0)