Skip to content

Commit 38ba68e

Browse files
committed
fix review comments
fix fix
1 parent 8063080 commit 38ba68e

8 files changed

Lines changed: 327 additions & 80 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: 53 additions & 22 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
@@ -395,15 +389,16 @@ where
395389

396390
SeqV::new(id.seq, *id.data)
397391
} else {
398-
// Reuse the metadata loaded for the engine check when marking the
399-
// existing table as dropped.
392+
let existing = VersionedTable {
393+
id: TableId::new(*id.data),
394+
meta: existing_table_meta,
395+
};
400396
let (seq, id) = construct_drop_table_txn_operations(
401397
self,
402398
req.name_ident.table_name.clone(),
403399
&req.name_ident.tenant,
404400
req.catalog_name.clone(),
405-
*id.data,
406-
Some(existing_table_meta),
401+
existing,
407402
*seq_db_id.data,
408403
true,
409404
false,
@@ -634,13 +629,23 @@ where
634629

635630
let mut txn = TxnRequest::default();
636631

632+
let tbid = TableId::new(table_id);
633+
let table_meta = self.get_pb(&tbid).await?.ok_or_else(|| {
634+
KVAppError::AppError(AppError::UnknownTableId(UnknownTableId::new(
635+
table_id,
636+
"drop_table_by_id failed to find valid tb_meta",
637+
)))
638+
})?;
639+
let existing = VersionedTable {
640+
id: tbid,
641+
meta: table_meta,
642+
};
637643
let opt = construct_drop_table_txn_operations(
638644
self,
639645
req.table_name.clone(),
640646
&req.tenant,
641647
None,
642-
table_id,
643-
None,
648+
existing,
644649
req.db_id,
645650
req.if_exists,
646651
true,
@@ -1192,7 +1197,35 @@ where
11921197

11931198
{
11941199
// reset drop on time
1195-
let mut tb_meta = tb_meta.unwrap();
1200+
let mut txn_req = TxnRequest::default();
1201+
let mut tb_meta = tb_meta.ok_or_else(|| {
1202+
KVAppError::AppError(AppError::UnknownTableId(UnknownTableId::new(
1203+
table_id,
1204+
"commit_table_meta",
1205+
)))
1206+
})?;
1207+
1208+
if let Some(prev_table_id) = req.prev_table_id {
1209+
let prev_tbid = TableId::new(prev_table_id);
1210+
let previous = self.get_pb(&prev_tbid).await?.ok_or_else(|| {
1211+
KVAppError::AppError(AppError::UnknownTableId(UnknownTableId::new(
1212+
prev_table_id,
1213+
"commit_table_meta previous table",
1214+
)))
1215+
})?;
1216+
1217+
TableEngineMismatch::ensure(
1218+
req.name_ident.table_name.as_str(),
1219+
&previous.engine,
1220+
&tb_meta.engine,
1221+
)
1222+
.map_err(|e| KVAppError::AppError(e.into()))?;
1223+
1224+
txn_req
1225+
.condition
1226+
.push(txn_cond_seq(&prev_tbid, Eq, previous.seq));
1227+
}
1228+
11961229
// undrop a table with no drop_on time
11971230
if tb_meta.drop_on.is_none() {
11981231
return Err(KVAppError::AppError(AppError::UndropTableWithNoDropTime(
@@ -1201,8 +1234,6 @@ where
12011234
}
12021235
tb_meta.drop_on = None;
12031236

1204-
let mut txn_req = TxnRequest::default();
1205-
12061237
txn_req.condition.extend([
12071238
// db has not to change, i.e., no new table is created.
12081239
// Renaming db is OK and does not affect the seq of db_meta.

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)