Skip to content

Commit 959c0a6

Browse files
committed
fix review comments
1 parent 8063080 commit 959c0a6

8 files changed

Lines changed: 330 additions & 82 deletions

File tree

Cargo.lock

Lines changed: 0 additions & 2 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

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

Lines changed: 11 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -171,33 +171,26 @@ pub async fn get_history_table_metas(
171171
Ok(tb_metas)
172172
}
173173

174-
pub async fn construct_drop_table_txn_operations(
174+
pub(crate) struct VersionedTable {
175+
pub id: TableId,
176+
pub meta: SeqV<TableMeta>,
177+
}
178+
179+
pub(crate) async fn construct_drop_table_txn_operations(
175180
kv_api: &(impl kvapi::KVApi<Error = MetaError> + ?Sized),
176181
table_name: String,
177182
tenant: &Tenant,
178183
catalog_name: Option<String>,
179-
table_id: u64,
180-
preloaded_table_meta: Option<SeqV<TableMeta>>,
184+
existing: VersionedTable,
181185
db_id: u64,
182186
if_exists: bool,
183187
if_delete: bool,
184188
txn: &mut TxnRequest,
185189
) -> Result<(u64, u64), KVAppError> {
186-
let tbid = TableId { table_id };
187-
188-
// Reuse metadata already loaded by CREATE OR REPLACE; ordinary DROP callers
189-
// do not have it and load it here. A preloaded value must belong to `table_id`.
190-
let (tb_meta_seq, tb_meta) = if let Some(preloaded_table_meta) = preloaded_table_meta {
191-
(preloaded_table_meta.seq, Some(preloaded_table_meta.data))
192-
} else {
193-
kv_api.get_pb_seq_and_value(&tbid).await?
194-
};
195-
// Check if table exists.
196-
if tb_meta_seq == 0 {
197-
return Err(KVAppError::AppError(AppError::UnknownTableId(
198-
UnknownTableId::new(table_id, "drop_table_by_id failed to find valid tb_meta"),
199-
)));
200-
}
190+
let VersionedTable { id: tbid, meta } = existing;
191+
let table_id = tbid.table_id;
192+
let tb_meta_seq = meta.seq;
193+
let mut tb_meta = meta.data;
201194

202195
// Get db name, tenant name and related info for tx.
203196
let table_id_to_name = TableIdToName { table_id };
@@ -239,7 +232,6 @@ pub async fn construct_drop_table_txn_operations(
239232
"drop table by id"
240233
);
241234

242-
let mut tb_meta = tb_meta.unwrap();
243235
// drop a table with drop_on time
244236
if tb_meta.drop_on.is_some() {
245237
return if if_exists {

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

Lines changed: 56 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -129,6 +129,7 @@ use super::database_api::DatabaseApi;
129129
use super::database_util::get_db_or_err;
130130
use super::garbage_collection_api::ORPHAN_POSTFIX;
131131
use super::garbage_collection_api::get_history_tables_for_gc;
132+
use super::schema_api::VersionedTable;
132133
use super::schema_api::build_upsert_table_deduplicated_label;
133134
use super::schema_api::construct_drop_table_txn_operations;
134135
use super::schema_api::get_db_by_id_or_err;
@@ -373,19 +374,12 @@ where
373374
"create or replace failed to find existing table meta",
374375
)))
375376
})?;
376-
if !existing_table_meta
377-
.engine
378-
.eq_ignore_ascii_case(&req.table_meta.engine)
379-
{
380-
return Err(KVAppError::AppError(
381-
TableEngineMismatch::new(
382-
req.table_name(),
383-
&existing_table_meta.engine,
384-
&req.table_meta.engine,
385-
)
386-
.into(),
387-
));
388-
}
377+
TableEngineMismatch::ensure(
378+
req.table_name(),
379+
&existing_table_meta.engine,
380+
&req.table_meta.engine,
381+
)
382+
.map_err(|e| KVAppError::AppError(e.into()))?;
389383

390384
if req.as_dropped {
391385
// If the table is being created as a dropped table, we do not
@@ -397,13 +391,16 @@ where
397391
} else {
398392
// Reuse the metadata loaded for the engine check when marking the
399393
// existing table as dropped.
394+
let existing = VersionedTable {
395+
id: TableId::new(*id.data),
396+
meta: existing_table_meta,
397+
};
400398
let (seq, id) = construct_drop_table_txn_operations(
401399
self,
402400
req.name_ident.table_name.clone(),
403401
&req.name_ident.tenant,
404402
req.catalog_name.clone(),
405-
*id.data,
406-
Some(existing_table_meta),
403+
existing,
407404
*seq_db_id.data,
408405
true,
409406
false,
@@ -634,13 +631,23 @@ where
634631

635632
let mut txn = TxnRequest::default();
636633

634+
let tbid = TableId::new(table_id);
635+
let table_meta = self.get_pb(&tbid).await?.ok_or_else(|| {
636+
KVAppError::AppError(AppError::UnknownTableId(UnknownTableId::new(
637+
table_id,
638+
"drop_table_by_id failed to find valid tb_meta",
639+
)))
640+
})?;
641+
let existing = VersionedTable {
642+
id: tbid,
643+
meta: table_meta,
644+
};
637645
let opt = construct_drop_table_txn_operations(
638646
self,
639647
req.table_name.clone(),
640648
&req.tenant,
641649
None,
642-
table_id,
643-
None,
650+
existing,
644651
req.db_id,
645652
req.if_exists,
646653
true,
@@ -1183,6 +1190,12 @@ where
11831190
// get tb_meta of the last table id
11841191
let tbid = TableId { table_id };
11851192
let (tb_meta_seq, tb_meta) = self.get_pb_seq_and_value(&tbid).await?;
1193+
let mut replacement_meta = tb_meta.ok_or_else(|| {
1194+
KVAppError::AppError(AppError::UnknownTableId(UnknownTableId::new(
1195+
table_id,
1196+
"commit_table_meta",
1197+
)))
1198+
})?;
11861199

11871200
debug!(
11881201
ident :% =(&tbid),
@@ -1191,17 +1204,36 @@ where
11911204
);
11921205

11931206
{
1194-
// reset drop on time
1195-
let mut tb_meta = tb_meta.unwrap();
1207+
let mut txn_req = TxnRequest::default();
1208+
1209+
if let Some(prev_table_id) = req.prev_table_id {
1210+
let prev_tbid = TableId::new(prev_table_id);
1211+
let previous = self.get_pb(&prev_tbid).await?.ok_or_else(|| {
1212+
KVAppError::AppError(AppError::UnknownTableId(UnknownTableId::new(
1213+
prev_table_id,
1214+
"commit_table_meta previous table",
1215+
)))
1216+
})?;
1217+
1218+
TableEngineMismatch::ensure(
1219+
req.name_ident.table_name.as_str(),
1220+
&previous.engine,
1221+
&replacement_meta.engine,
1222+
)
1223+
.map_err(|e| KVAppError::AppError(e.into()))?;
1224+
1225+
txn_req
1226+
.condition
1227+
.push(txn_cond_seq(&prev_tbid, Eq, previous.seq));
1228+
}
1229+
11961230
// undrop a table with no drop_on time
1197-
if tb_meta.drop_on.is_none() {
1231+
if replacement_meta.drop_on.is_none() {
11981232
return Err(KVAppError::AppError(AppError::UndropTableWithNoDropTime(
11991233
UndropTableWithNoDropTime::new(&tenant_dbname_tbname.table_name),
12001234
)));
12011235
}
1202-
tb_meta.drop_on = None;
1203-
1204-
let mut txn_req = TxnRequest::default();
1236+
replacement_meta.drop_on = None;
12051237

12061238
txn_req.condition.extend([
12071239
// db has not to change, i.e., no new table is created.
@@ -1220,7 +1252,7 @@ where
12201252
txn_put_pb(&DatabaseId { db_id }, &db_meta), // (db_id) -> db_meta
12211253
txn_put_pb(&dbid_tbname, &TableId::new(table_id)), /* (tenant, db_id, tb_name) -> tb_id */
12221254
// txn_put_pb(&dbid_tbname_idlist, &tb_id_list), // _fd_table_id_list/db_id/table_name -> tb_id_list
1223-
txn_put_pb(&tbid, &tb_meta), // (tenant, db_id, tb_id) -> tb_meta
1255+
txn_put_pb(&tbid, &replacement_meta), // (tenant, db_id, tb_id) -> tb_meta
12241256
txn_del(&orphan_dbid_tbname_idlist), // del orphan table idlist
12251257
txn_put_pb(&dbid_tbname_idlist, &tb_id_list.data), /* _fd_table_id_list/db_id/table_name -> tb_id_list */
12261258
]);

src/meta/app/src/app_error.rs

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -217,6 +217,14 @@ impl TableEngineMismatch {
217217
new_engine: new_engine.into(),
218218
}
219219
}
220+
221+
pub fn ensure(table_name: &str, existing_engine: &str, new_engine: &str) -> Result<(), Self> {
222+
if existing_engine.eq_ignore_ascii_case(new_engine) {
223+
return Ok(());
224+
}
225+
226+
Err(Self::new(table_name, existing_engine, new_engine))
227+
}
220228
}
221229

222230
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]

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

Lines changed: 90 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,7 @@ use databend_common_meta_api::serialize_struct;
5151
use databend_common_meta_api::util::IdempotentKVTxnSender;
5252
use databend_common_meta_app::KeyWithTenant;
5353
use databend_common_meta_app::app_error::AppError;
54+
use databend_common_meta_app::app_error::TableEngineMismatch;
5455
use databend_common_meta_app::data_mask::CreateDatamaskReq;
5556
use databend_common_meta_app::data_mask::DataMaskNameIdent;
5657
use databend_common_meta_app::data_mask::DatamaskMeta;
@@ -318,6 +319,8 @@ impl SchemaApiTestSuite {
318319
+ 'static,
319320
{
320321
self.table_commit_table_meta(&b.build().await).await?;
322+
self.table_commit_table_meta_engine_mismatch(&b.build().await)
323+
.await?;
321324
self.concurrent_commit_table_meta(b.clone()).await?;
322325
self.database_drop_out_of_retention_time_history(&b.build().await)
323326
.await?;
@@ -2634,11 +2637,11 @@ impl SchemaApiTestSuite {
26342637
table_properties: None,
26352638
table_partition: None,
26362639
};
2640+
let expected = KVAppError::AppError(AppError::from(TableEngineMismatch::new(
2641+
table, "JSON", "STREAM",
2642+
)));
26372643
let err = mt.create_table(mismatched_req.clone()).await.unwrap_err();
2638-
assert!(matches!(
2639-
err,
2640-
KVAppError::AppError(AppError::TableEngineMismatch(_))
2641-
));
2644+
assert_eq!(err, expected);
26422645
assert_eq!(table_id, get_kv_u64_data(mt, &key_dbid_tbname).await?);
26432646
assert_eq!("JSON", util.get_table_by_name(table).await?.meta.engine);
26442647

@@ -2648,10 +2651,7 @@ impl SchemaApiTestSuite {
26482651
as_dropped_req.as_dropped = true;
26492652
as_dropped_req.table_meta.drop_on = Some(Utc::now());
26502653
let err = mt.create_table(as_dropped_req).await.unwrap_err();
2651-
assert!(matches!(
2652-
err,
2653-
KVAppError::AppError(AppError::TableEngineMismatch(_))
2654-
));
2654+
assert_eq!(err, expected);
26552655
assert_eq!(table_id, get_kv_u64_data(mt, &key_dbid_tbname).await?);
26562656
}
26572657

@@ -6523,6 +6523,88 @@ impl SchemaApiTestSuite {
65236523
Ok(())
65246524
}
65256525

6526+
async fn table_commit_table_meta_engine_mismatch<MT>(&self, mt: &MT) -> anyhow::Result<()>
6527+
where MT: kvapi::KVApi<Error = MetaError> + DatabaseApi + TableApi {
6528+
let tenant_name = "table_commit_table_meta_engine_mismatch";
6529+
let db_name = "db";
6530+
let table_name = "table";
6531+
let mut util = DbTableHarness::new(mt, tenant_name, db_name, table_name, "JSON");
6532+
util.create_db().await?;
6533+
let (previous_table_id, previous_meta) = util.create_table().await?;
6534+
6535+
let mut replacement_meta = previous_meta.clone();
6536+
replacement_meta.drop_on = Some(Utc::now());
6537+
let create_req = CreateTableReq {
6538+
create_option: CreateOption::CreateOrReplace,
6539+
catalog_name: Some("default".to_string()),
6540+
name_ident: TableNameIdent {
6541+
tenant: util.tenant(),
6542+
db_name: db_name.to_string(),
6543+
table_name: table_name.to_string(),
6544+
},
6545+
table_meta: replacement_meta,
6546+
as_dropped: true,
6547+
materialized_view: None,
6548+
table_properties: None,
6549+
table_partition: None,
6550+
};
6551+
let create_reply = mt.create_table(create_req.clone()).await?;
6552+
assert_eq!(create_reply.prev_table_id, Some(previous_table_id));
6553+
6554+
let previous_tbid = TableId::new(previous_table_id);
6555+
let mut changed_previous_meta = previous_meta;
6556+
changed_previous_meta.engine = "STREAM".to_string();
6557+
upsert_test_data(mt, &previous_tbid, serialize_struct(&changed_previous_meta)).await?;
6558+
6559+
let dbid_tbname = DBIdTableName {
6560+
db_id: create_reply.db_id,
6561+
table_name: table_name.to_string(),
6562+
};
6563+
let history_ident = TableIdHistoryIdent {
6564+
database_id: create_reply.db_id,
6565+
table_name: table_name.to_string(),
6566+
};
6567+
let orphan_history_ident = TableIdHistoryIdent {
6568+
database_id: create_reply.db_id,
6569+
table_name: create_reply.orphan_table_name.clone().unwrap(),
6570+
};
6571+
let replacement_tbid = TableId::new(create_reply.table_id);
6572+
6573+
let name_before = mt.get_pb(&dbid_tbname).await?;
6574+
let history_before = mt.get_pb(&history_ident).await?;
6575+
let orphan_history_before = mt.get_pb(&orphan_history_ident).await?;
6576+
let previous_before = mt.get_pb(&previous_tbid).await?;
6577+
let replacement_before = mt.get_pb(&replacement_tbid).await?;
6578+
6579+
let err = mt
6580+
.commit_table_meta(CommitTableMetaReq {
6581+
name_ident: create_req.name_ident,
6582+
db_id: create_reply.db_id,
6583+
table_id: create_reply.table_id,
6584+
prev_table_id: create_reply.prev_table_id,
6585+
orphan_table_name: create_reply.orphan_table_name,
6586+
})
6587+
.await
6588+
.unwrap_err();
6589+
assert_eq!(
6590+
err,
6591+
KVAppError::AppError(AppError::from(TableEngineMismatch::new(
6592+
table_name, "STREAM", "JSON"
6593+
)))
6594+
);
6595+
6596+
assert_eq!(mt.get_pb(&dbid_tbname).await?, name_before);
6597+
assert_eq!(mt.get_pb(&history_ident).await?, history_before);
6598+
assert_eq!(
6599+
mt.get_pb(&orphan_history_ident).await?,
6600+
orphan_history_before
6601+
);
6602+
assert_eq!(mt.get_pb(&previous_tbid).await?, previous_before);
6603+
assert_eq!(mt.get_pb(&replacement_tbid).await?, replacement_before);
6604+
6605+
Ok(())
6606+
}
6607+
65266608
async fn concurrent_commit_table_meta<
65276609
B: kvapi::ApiBuilder<MT>,
65286610
MT: kvapi::KVApi<Error = MetaError> + DatabaseApi + TableApi + 'static,

0 commit comments

Comments
 (0)