diff --git a/src/meta/api/src/api_impl/schema_api.rs b/src/meta/api/src/api_impl/schema_api.rs index 5f5ca5cb2555e..687400672d86a 100644 --- a/src/meta/api/src/api_impl/schema_api.rs +++ b/src/meta/api/src/api_impl/schema_api.rs @@ -171,26 +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, + existing: VersionedTable, db_id: u64, if_exists: bool, if_delete: bool, txn: &mut TxnRequest, ) -> Result<(u64, u64), KVAppError> { - let tbid = TableId { table_id }; - - // 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"), - ))); - } + 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 }; @@ -232,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 9b524c965f987..7513546403fa6 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; @@ -128,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; @@ -361,6 +363,24 @@ where }); } CreateOption::CreateOrReplace => { + // 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", + ))) + })?; + 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 // need to combine with drop_table_txn operations, just return @@ -369,12 +389,16 @@ where SeqV::new(id.seq, *id.data) } else { + 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, + existing, *seq_db_id.data, true, false, @@ -605,12 +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, + existing, req.db_id, req.if_exists, true, @@ -1078,7 +1113,8 @@ where table_name: tenant_dbname_tbname.table_name.clone(), }; - let (dbid_tbname_seq, _table_id) = self.get_pb_seq_and_value(&dbid_tbname).await?; + let (dbid_tbname_seq, visible_table_id) = + self.get_pb_seq_and_value(&dbid_tbname).await?; // get table id list from _fd_table_id_list/db_id/table_name @@ -1162,7 +1198,34 @@ 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(visible_table_id) = visible_table_id { + let previous = self.get_pb(&visible_table_id).await?.ok_or_else(|| { + KVAppError::AppError(AppError::UnknownTableId(UnknownTableId::new( + visible_table_id.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(&visible_table_id, Eq, previous.seq)); + } + // undrop a table with no drop_on time if tb_meta.drop_on.is_none() { return Err(KVAppError::AppError(AppError::UndropTableWithNoDropTime( @@ -1171,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 82303633c1ddd..c25ff873e34d1 100644 --- a/src/meta/app/src/app_error.rs +++ b/src/meta/app/src/app_error.rs @@ -195,6 +195,38 @@ 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(), + } + } + + 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)] #[error("ViewAlreadyExists: {view_name} while {context}")] pub struct ViewAlreadyExists { @@ -895,6 +927,9 @@ pub enum AppError { #[error(transparent)] TableAlreadyExists(#[from] TableAlreadyExists), + #[error(transparent)] + TableEngineMismatch(#[from] TableEngineMismatch), + #[error(transparent)] ViewAlreadyExists(#[from] ViewAlreadyExists), @@ -1199,6 +1234,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 +1449,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..652cfb2731d8b 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,10 @@ impl SchemaApiTestSuite { + 'static, { self.table_commit_table_meta(&b.build().await).await?; + self.table_commit_after_drop_different_engine(&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?; @@ -2583,6 +2588,72 @@ 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. + 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", util.get_table_by_name(table).await?.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 expected = KVAppError::AppError(AppError::from(TableEngineMismatch::new( + table, "JSON", "STREAM", + ))); + let err = mt.create_table(mismatched_req.clone()).await.unwrap_err(); + 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); + + // 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_eq!(err, expected); + assert_eq!(table_id, get_kv_u64_data(mt, &key_dbid_tbname).await?); } { @@ -6453,6 +6524,136 @@ impl SchemaApiTestSuite { Ok(()) } + async fn table_commit_after_drop_different_engine(&self, mt: &MT) -> anyhow::Result<()> + where MT: kvapi::KVApi + DatabaseApi + TableApi { + let tenant_name = "table_commit_after_drop_different_engine"; + let db_name = "db"; + let table_name = "table"; + let mut util = DbTableHarness::new(mt, tenant_name, db_name, table_name, "STREAM"); + util.create_db().await?; + let (dropped_table_id, _) = util.create_table().await?; + util.drop_table_by_id().await?; + + let mut replacement_meta = util.table_meta(); + replacement_meta.engine = "JSON".to_string(); + replacement_meta.drop_on = Some(Utc::now()); + let create_req = CreateTableReq { + create_option: CreateOption::Create, + catalog_name: None, + 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(dropped_table_id)); + + 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?; + + let table = util.get_table_by_name(table_name).await?; + assert_eq!(table.ident.table_id, create_reply.table_id); + assert_eq!(table.meta.engine, "JSON"); + + 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(); + // Simulate an arbitrary metadata client updating the full visible TableMeta + // after CTAS prepare. + 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/ee/src/stream/handler.rs b/src/query/ee/src/stream/handler.rs index cf8ec20d6c706..080d571a5a7c5 100644 --- a/src/query/ee/src/stream/handler.rs +++ b/src/query/ee/src/stream/handler.rs @@ -86,7 +86,8 @@ impl StreamHandler for RealStreamHandler { let table_id = table_info.ident.table_id; if !table.change_tracking_enabled() { let table_seq = table_info.ident.seq; - // enable change tracking. + // Enabling change tracking is independent of whether the subsequent + // stream creation succeeds and is not rolled back if it fails. let req = UpsertTableOptionReq { table_id, seq: MatchSeq::Exact(table_seq), diff --git a/src/query/storages/common/session/src/temp_table.rs b/src/query/storages/common/session/src/temp_table.rs index 5dbfe89bb6f35..51195cfe93f29 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; @@ -54,6 +56,10 @@ use log::info; use opendal::Operator; use parking_lot::Mutex; +// Staged CTAS replacements share name_to_id with user-created temporary tables, so reserve a +// namespace that users cannot enter through CREATE or RENAME. +const TEMP_ORPHAN_TABLE_PREFIX: &str = "__tmp_orphan@"; + #[derive(Debug, Clone)] pub struct TempTblMgr { name_to_id: HashMap, @@ -74,6 +80,19 @@ impl TempTblMgr { format!("'{}'.'{}'", db_name, table_name) } + fn is_orphan_table_name(table_name: &str) -> bool { + table_name.starts_with(TEMP_ORPHAN_TABLE_PREFIX) + } + + fn ensure_user_table_name(table_name: &str) -> Result<()> { + if Self::is_orphan_table_name(table_name) { + return Err(ErrorCode::BadArguments(format!( + "Temporary table names starting with '{TEMP_ORPHAN_TABLE_PREFIX}' are reserved" + ))); + } + Ok(()) + } + pub fn init() -> Arc> { Arc::new(Mutex::new(TempTblMgr { name_to_id: HashMap::new(), @@ -105,8 +124,7 @@ impl TempTblMgr { as_dropped, .. } = req; - let orphan_table_name = as_dropped.then(|| format!("orphan@{}", name_ident.table_name)); - + Self::ensure_user_table_name(&name_ident.table_name)?; let Some(db_id) = table_meta.options.get(OPT_KEY_DATABASE_ID) else { return Err(ErrorCode::Internal(format!( "Database id not set in table options" @@ -115,17 +133,35 @@ 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.contains_key(&desc), create_option) { - (true, CreateOption::Create) => { + // Keep a two-phase CTAS replacement hidden under its unique temporary table ID until + // commit_table_meta publishes it under the requested name. + let orphan_table_name = as_dropped.then(|| format!("{TEMP_ORPHAN_TABLE_PREFIX}{table_id}")); + let new_table = match (existing_id, 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" + )) + })?; + 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 .as_ref() .map(|o| Self::temp_table_desc(&name_ident.db_name, o)) @@ -154,7 +190,9 @@ impl TempTblMgr { table_id_seq: Some(0), db_id, new_table, - prev_table_id: None, + // The commit guard must compare against the visible table observed during prepare. + // Direct, single-phase creates do not need this value. + prev_table_id: as_dropped.then_some(existing_id).flatten(), orphan_table_name, }) } @@ -165,22 +203,37 @@ 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(); + // Validate the staged mapping before mutating it: concurrent CTAS requests may have the + // same visible target, but each commit must publish only its own replacement. + let orphan_id = self + .name_to_id + .get(&orphan_desc) + .copied() + .ok_or_else(|| ErrorCode::UnknownTable(orphan_desc.clone()))?; + if orphan_id != req.table_id { + return Err(ErrorCode::TableVersionMismatched(format!( + "Temporary table {desc} replacement 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 - ))), + // Do not overwrite a visible table that changed while the CTAS pipeline was running. + 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" + ))); + } + + // Both publication guards passed, so removing the orphan mapping is now safe. + 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> { @@ -191,15 +244,27 @@ impl TempTblMgr { new_table_name, } = req; let desc = Self::temp_table_desc(&name_ident.db_name, &name_ident.table_name); - match self.name_to_id.remove(&desc) { + if Self::is_orphan_table_name(&name_ident.table_name) { + // Reject access to an active staged replacement. Only names not owned by this manager + // may fall through to the underlying catalog. + if self.name_to_id.contains_key(&desc) { + Self::ensure_user_table_name(&name_ident.table_name)?; + } + return Ok(None); + } + // Keep the source mapping intact until all destination checks pass. + match self.name_to_id.get(&desc).copied() { Some(id) => { + Self::ensure_user_table_name(new_table_name)?; let new_desc = Self::temp_table_desc(new_db_name, new_table_name); - if self.name_to_id.contains_key(&new_desc) { + // A no-op rename finds the source itself and is not a destination collision. + if new_desc != desc && self.name_to_id.contains_key(&new_desc) { return Err(ErrorCode::TableAlreadyExists(format!( "Temporary table {} already exists", new_desc ))); } + self.name_to_id.remove(&desc); self.name_to_id.insert(new_desc, id); let table = self.id_to_table.get_mut(&id).unwrap(); table.db_name = new_db_name.clone(); @@ -513,3 +578,318 @@ 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 new_temp_table_mgr() -> TempTblMgr { + TempTblMgr { + name_to_id: HashMap::new(), + id_to_table: HashMap::new(), + next_id: TEMP_TBL_ID_BEGIN, + } + } + + fn table_name_ident(table_name: &str) -> TableNameIdent { + TableNameIdent { + tenant: Tenant::new_literal("tenant"), + db_name: "db".to_string(), + table_name: table_name.to_string(), + } + } + + 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: table_name_ident(table_name), + table_meta, + as_dropped, + materialized_view: None, + table_properties: None, + table_partition: None, + } + } + + fn rename_table_req(table_name: &str, new_table_name: &str) -> RenameTableReq { + RenameTableReq { + if_exists: false, + name_ident: table_name_ident(table_name), + new_db_name: "db".to_string(), + new_table_name: new_table_name.to_string(), + } + } + + fn commit_table_req(table_name: &str, reply: &CreateTableReply) -> CommitTableMetaReq { + CommitTableMetaReq { + name_ident: table_name_ident(table_name), + db_id: reply.db_id, + table_id: reply.table_id, + prev_table_id: reply.prev_table_id, + orphan_table_name: reply.orphan_table_name.clone(), + } + } + + #[test] + fn test_ctas_commit_preserves_temporary_tables_when_visible_table_changed() { + let mut mgr = new_temp_table_mgr(); + 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(&commit_table_req(table_name, &replacement)) + .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 + ); + } + + #[test] + fn test_ctas_commit_publishes_its_own_temporary_table() { + let mut mgr = new_temp_table_mgr(); + 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 first = mgr + .create_table( + create_table_req(table_name, "FUSE", CreateOption::CreateOrReplace, true), + "session".to_string(), + ) + .unwrap(); + let second = mgr + .create_table( + create_table_req(table_name, "FUSE", CreateOption::CreateOrReplace, true), + "session".to_string(), + ) + .unwrap(); + + assert_eq!(first.prev_table_id, Some(visible.table_id)); + assert_eq!(second.prev_table_id, Some(visible.table_id)); + assert_ne!(first.orphan_table_name, second.orphan_table_name); + let first_orphan_desc = + TempTblMgr::temp_table_desc("db", first.orphan_table_name.as_ref().unwrap()); + let second_orphan_desc = + TempTblMgr::temp_table_desc("db", second.orphan_table_name.as_ref().unwrap()); + assert_eq!( + mgr.name_to_id.get(&first_orphan_desc), + Some(&first.table_id) + ); + assert_eq!( + mgr.name_to_id.get(&second_orphan_desc), + Some(&second.table_id) + ); + + let name_to_id_before = mgr.name_to_id.clone(); + let table_meta_before: HashMap<_, _> = mgr + .id_to_table + .iter() + .map(|(id, table)| (*id, table.meta.clone())) + .collect(); + let mut mismatched_req = commit_table_req(table_name, &first); + mismatched_req.table_id = second.table_id; + let err = mgr.commit_table_meta(&mismatched_req).unwrap_err(); + let expected = ErrorCode::TableVersionMismatched(format!( + "Temporary table {desc} replacement changed while CTAS was running" + )); + assert_eq!( + (err.code(), err.message()), + (expected.code(), expected.message()) + ); + assert_eq!(mgr.name_to_id, name_to_id_before); + assert_eq!( + mgr.id_to_table + .iter() + .map(|(id, table)| (*id, table.meta.clone())) + .collect::>(), + table_meta_before + ); + + mgr.commit_table_meta(&commit_table_req(table_name, &first)) + .unwrap(); + + assert_eq!(mgr.name_to_id.get(&desc), Some(&first.table_id)); + assert!(!mgr.name_to_id.contains_key(&first_orphan_desc)); + assert_eq!( + mgr.name_to_id.get(&second_orphan_desc), + Some(&second.table_id) + ); + assert!(!mgr.id_to_table.contains_key(&visible.table_id)); + assert!(mgr.id_to_table.contains_key(&first.table_id)); + assert!(mgr.id_to_table.contains_key(&second.table_id)); + + let err = mgr + .commit_table_meta(&commit_table_req(table_name, &second)) + .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(&first.table_id)); + assert_eq!( + mgr.name_to_id.get(&second_orphan_desc), + Some(&second.table_id) + ); + assert!(mgr.id_to_table.contains_key(&first.table_id)); + assert!(mgr.id_to_table.contains_key(&second.table_id)); + } + + #[test] + fn test_ctas_orphan_prefix_is_reserved() { + let mut mgr = new_temp_table_mgr(); + let reserved_name = format!("{TEMP_ORPHAN_TABLE_PREFIX}user"); + + let err = mgr + .create_table( + create_table_req(&reserved_name, "MEMORY", CreateOption::Create, false), + "session".to_string(), + ) + .unwrap_err(); + let expected = ErrorCode::BadArguments(format!( + "Temporary table names starting with '{TEMP_ORPHAN_TABLE_PREFIX}' are reserved" + )); + assert_eq!( + (err.code(), err.message()), + (expected.code(), expected.message()) + ); + + let visible = mgr + .create_table( + create_table_req("t", "MEMORY", CreateOption::Create, false), + "session".to_string(), + ) + .unwrap(); + assert!( + mgr.rename_table(&rename_table_req("t", "t")) + .unwrap() + .is_some() + ); + let err = mgr + .rename_table(&rename_table_req("t", &reserved_name)) + .unwrap_err(); + assert_eq!( + (err.code(), err.message()), + (expected.code(), expected.message()) + ); + assert_eq!( + mgr.name_to_id.get(&TempTblMgr::temp_table_desc("db", "t")), + Some(&visible.table_id) + ); + + // Names not owned by the temporary manager must fall through to the inner catalog. + assert!( + mgr.rename_table(&rename_table_req("persistent", &reserved_name)) + .unwrap() + .is_none() + ); + assert!( + mgr.rename_table(&rename_table_req(&reserved_name, "persistent")) + .unwrap() + .is_none() + ); + + let replacement = mgr + .create_table( + create_table_req("t", "MEMORY", CreateOption::CreateOrReplace, true), + "session".to_string(), + ) + .unwrap(); + let orphan_table_name = replacement.orphan_table_name.as_ref().unwrap(); + let orphan_desc = TempTblMgr::temp_table_desc("db", orphan_table_name); + assert_eq!( + replacement.orphan_table_name.as_deref(), + Some(format!("{TEMP_ORPHAN_TABLE_PREFIX}{}", replacement.table_id).as_str()) + ); + assert_eq!( + mgr.name_to_id.get(&orphan_desc), + Some(&replacement.table_id) + ); + let err = mgr + .rename_table(&rename_table_req(orphan_table_name, "hijacked")) + .unwrap_err(); + assert_eq!( + (err.code(), err.message()), + (expected.code(), expected.message()) + ); + assert_eq!( + mgr.name_to_id.get(&orphan_desc), + Some(&replacement.table_id) + ); + } +} 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..f0380f952b1ab --- /dev/null +++ b/tests/sqllogictests/suites/ee/06_ee_stream/06_0016_stream_replace_engine.test @@ -0,0 +1,68 @@ +## 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) + +statement error 1302 +CREATE OR REPLACE TABLE target AS SELECT 1 AS a + +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;