1pub mod catalog_name;
101pub mod datanode_table;
102pub mod flow;
103pub mod node_address;
104pub mod runtime_switch;
105mod schema_metadata_manager;
106pub mod schema_name;
107pub mod table_info;
108pub mod table_name;
109pub mod table_repart;
110pub mod table_route;
111#[cfg(any(test, feature = "testing"))]
112pub mod test_utils;
113pub mod tombstone;
114pub mod topic_name;
115pub mod topic_region;
116pub mod txn_helper;
117pub mod view_info;
118
119use std::collections::{BTreeMap, HashMap, HashSet};
120use std::fmt::Debug;
121use std::ops::{Deref, DerefMut};
122use std::sync::Arc;
123
124use bytes::Bytes;
125use common_base::regex_pattern::NAME_PATTERN;
126use common_catalog::consts::{
127 DEFAULT_CATALOG_NAME, DEFAULT_PRIVATE_SCHEMA_NAME, DEFAULT_SCHEMA_NAME, INFORMATION_SCHEMA_NAME,
128};
129use common_telemetry::warn;
130use common_wal::options::WalOptions;
131use datanode_table::{DatanodeTableKey, DatanodeTableManager, DatanodeTableValue};
132use flow::flow_route::FlowRouteValue;
133use flow::table_flow::TableFlowValue;
134use futures_util::TryStreamExt;
135use futures_util::stream::BoxStream;
136use lazy_static::lazy_static;
137use regex::Regex;
138pub use schema_metadata_manager::{SchemaMetadataManager, SchemaMetadataManagerRef};
139use serde::de::DeserializeOwned;
140use serde::{Deserialize, Serialize};
141use snafu::{OptionExt, ResultExt, ensure};
142use store_api::storage::RegionNumber;
143use table::metadata::{TableId, TableInfo};
144use table::table_name::TableName;
145use table_info::{TableInfoKey, TableInfoManager, TableInfoValue};
146use table_name::{TableNameKey, TableNameManager, TableNameValue};
147use topic_name::TopicNameManager;
148use topic_region::{TopicRegionKey, TopicRegionManager};
149use view_info::{ViewInfoKey, ViewInfoManager, ViewInfoValue};
150
151use self::catalog_name::{CatalogManager, CatalogNameKey, CatalogNameValue};
152use self::datanode_table::RegionInfo;
153use self::flow::flow_info::FlowInfoValue;
154use self::flow::flow_name::FlowNameValue;
155use self::schema_name::{SchemaManager, SchemaNameKey, SchemaNameValue};
156use self::table_route::{TableRouteManager, TableRouteValue};
157use self::tombstone::TombstoneManager;
158use crate::DatanodeId;
159use crate::error::{self, Result, SerdeJsonSnafu};
160use crate::key::flow::flow_state::FlowStateValue;
161use crate::key::node_address::NodeAddressValue;
162use crate::key::table_repart::{TableRepartKey, TableRepartManager};
163use crate::key::table_route::TableRouteKey;
164use crate::key::topic_region::TopicRegionValue;
165use crate::key::txn_helper::TxnOpGetResponseSet;
166use crate::kv_backend::KvBackendRef;
167use crate::kv_backend::txn::{Txn, TxnOp};
168use crate::rpc::KeyValue;
169use crate::rpc::router::{LeaderState, RegionRoute, region_distribution};
170use crate::rpc::store::{BatchDeleteRequest, PutRequest};
171use crate::state_store::PoisonValue;
172use crate::wal_provider::RegionWalOptions;
173
174pub const TOPIC_NAME_PATTERN: &str = r"[a-zA-Z0-9_:-][a-zA-Z0-9_:\-\.@#]*";
175pub const LEGACY_MAINTENANCE_KEY: &str = "__maintenance";
176pub const MAINTENANCE_KEY: &str = "__switches/maintenance";
177pub const REPARTITION_GC_REQUIRED_KEY: &str = "__requirements/gc/repartition";
178pub const PAUSE_PROCEDURE_KEY: &str = "__switches/pause_procedure";
179pub const RECOVERY_MODE_KEY: &str = "__switches/recovery";
180
181pub const DATANODE_TABLE_KEY_PREFIX: &str = "__dn_table";
182pub const TABLE_INFO_KEY_PREFIX: &str = "__table_info";
183pub const VIEW_INFO_KEY_PREFIX: &str = "__view_info";
184pub const TABLE_NAME_KEY_PREFIX: &str = "__table_name";
185pub const CATALOG_NAME_KEY_PREFIX: &str = "__catalog_name";
186pub const SCHEMA_NAME_KEY_PREFIX: &str = "__schema_name";
187pub const TABLE_ROUTE_PREFIX: &str = "__table_route";
188pub const TABLE_REPART_PREFIX: &str = "__table_repart";
189pub const NODE_ADDRESS_PREFIX: &str = "__node_address";
190pub const KAFKA_TOPIC_KEY_PREFIX: &str = "__topic_name/kafka";
191pub const LEGACY_TOPIC_KEY_PREFIX: &str = "__created_wal_topics/kafka";
193pub const TOPIC_REGION_PREFIX: &str = "__topic_region";
194
195pub const ELECTION_KEY: &str = "__metasrv_election";
197pub const CANDIDATES_ROOT: &str = "__metasrv_election_candidates/";
199
200pub const CACHE_KEY_PREFIXES: [&str; 5] = [
202 TABLE_NAME_KEY_PREFIX,
203 CATALOG_NAME_KEY_PREFIX,
204 SCHEMA_NAME_KEY_PREFIX,
205 TABLE_ROUTE_PREFIX,
206 NODE_ADDRESS_PREFIX,
207];
208
209#[derive(Debug, Clone, PartialEq, Eq, Default, Serialize)]
211pub struct RegionRoleSet {
212 pub leader_regions: Vec<RegionNumber>,
214 pub follower_regions: Vec<RegionNumber>,
216}
217
218impl<'de> Deserialize<'de> for RegionRoleSet {
219 fn deserialize<D>(deserializer: D) -> std::result::Result<Self, D::Error>
220 where
221 D: serde::Deserializer<'de>,
222 {
223 #[derive(Deserialize)]
224 #[serde(untagged)]
225 enum RegionRoleSetOrLeaderOnly {
226 Full {
227 leader_regions: Vec<RegionNumber>,
228 follower_regions: Vec<RegionNumber>,
229 },
230 LeaderOnly(Vec<RegionNumber>),
231 }
232 match RegionRoleSetOrLeaderOnly::deserialize(deserializer)? {
233 RegionRoleSetOrLeaderOnly::Full {
234 leader_regions,
235 follower_regions,
236 } => Ok(RegionRoleSet::new(leader_regions, follower_regions)),
237 RegionRoleSetOrLeaderOnly::LeaderOnly(leader_regions) => {
238 Ok(RegionRoleSet::new(leader_regions, vec![]))
239 }
240 }
241 }
242}
243
244impl RegionRoleSet {
245 pub fn new(leader_regions: Vec<RegionNumber>, follower_regions: Vec<RegionNumber>) -> Self {
247 Self {
248 leader_regions,
249 follower_regions,
250 }
251 }
252
253 pub fn add_leader_region(&mut self, region_number: RegionNumber) {
255 self.leader_regions.push(region_number);
256 }
257
258 pub fn add_follower_region(&mut self, region_number: RegionNumber) {
260 self.follower_regions.push(region_number);
261 }
262
263 pub fn sort(&mut self) {
265 self.follower_regions.sort();
266 self.leader_regions.sort();
267 }
268}
269
270pub type RegionDistribution = BTreeMap<DatanodeId, RegionRoleSet>;
274
275pub type FlowId = u32;
277pub type FlowPartitionId = u32;
279
280lazy_static! {
281 pub static ref TOPIC_NAME_PATTERN_REGEX: Regex = Regex::new(TOPIC_NAME_PATTERN).unwrap();
282}
283
284lazy_static! {
285 static ref TABLE_INFO_KEY_PATTERN: Regex =
286 Regex::new(&format!("^{TABLE_INFO_KEY_PREFIX}/([0-9]+)$")).unwrap();
287}
288
289lazy_static! {
290 static ref VIEW_INFO_KEY_PATTERN: Regex =
291 Regex::new(&format!("^{VIEW_INFO_KEY_PREFIX}/([0-9]+)$")).unwrap();
292}
293
294lazy_static! {
295 static ref TABLE_ROUTE_KEY_PATTERN: Regex =
296 Regex::new(&format!("^{TABLE_ROUTE_PREFIX}/([0-9]+)$")).unwrap();
297}
298
299lazy_static! {
300 pub(crate) static ref TABLE_REPART_KEY_PATTERN: Regex =
301 Regex::new(&format!("^{TABLE_REPART_PREFIX}/([0-9]+)$")).unwrap();
302}
303
304lazy_static! {
305 static ref DATANODE_TABLE_KEY_PATTERN: Regex =
306 Regex::new(&format!("^{DATANODE_TABLE_KEY_PREFIX}/([0-9]+)/([0-9]+)$")).unwrap();
307}
308
309lazy_static! {
310 static ref TABLE_NAME_KEY_PATTERN: Regex = Regex::new(&format!(
311 "^{TABLE_NAME_KEY_PREFIX}/({NAME_PATTERN})/({NAME_PATTERN})/({NAME_PATTERN})$"
312 ))
313 .unwrap();
314}
315
316lazy_static! {
317 static ref CATALOG_NAME_KEY_PATTERN: Regex = Regex::new(&format!(
319 "^{CATALOG_NAME_KEY_PREFIX}/({NAME_PATTERN})$"
320 ))
321 .unwrap();
322}
323
324lazy_static! {
325 static ref SCHEMA_NAME_KEY_PATTERN:Regex=Regex::new(&format!(
327 "^{SCHEMA_NAME_KEY_PREFIX}/({NAME_PATTERN})/({NAME_PATTERN})$"
328 ))
329 .unwrap();
330}
331
332lazy_static! {
333 static ref NODE_ADDRESS_PATTERN: Regex =
334 Regex::new(&format!("^{NODE_ADDRESS_PREFIX}/([0-9]+)/([0-9]+)$")).unwrap();
335}
336
337lazy_static! {
338 pub static ref KAFKA_TOPIC_KEY_PATTERN: Regex =
339 Regex::new(&format!("^{KAFKA_TOPIC_KEY_PREFIX}/(.*)$")).unwrap();
340}
341
342lazy_static! {
343 pub static ref TOPIC_REGION_PATTERN: Regex = Regex::new(&format!(
344 "^{TOPIC_REGION_PREFIX}/({TOPIC_NAME_PATTERN})/([0-9]+)$"
345 ))
346 .unwrap();
347}
348
349pub trait MetadataKey<'a, T> {
351 fn to_bytes(&self) -> Vec<u8>;
352
353 fn from_bytes(bytes: &'a [u8]) -> Result<T>;
354}
355
356#[derive(Debug, Clone, PartialEq)]
357pub struct BytesAdapter(Vec<u8>);
358
359impl From<Vec<u8>> for BytesAdapter {
360 fn from(value: Vec<u8>) -> Self {
361 Self(value)
362 }
363}
364
365impl<'a> MetadataKey<'a, BytesAdapter> for BytesAdapter {
366 fn to_bytes(&self) -> Vec<u8> {
367 self.0.clone()
368 }
369
370 fn from_bytes(bytes: &'a [u8]) -> Result<BytesAdapter> {
371 Ok(BytesAdapter(bytes.to_vec()))
372 }
373}
374
375pub(crate) trait MetadataKeyGetTxnOp {
376 fn build_get_op(
377 &self,
378 ) -> (
379 TxnOp,
380 impl for<'a> FnMut(&'a mut TxnOpGetResponseSet) -> Option<Vec<u8>>,
381 );
382}
383
384pub trait MetadataValue {
385 fn try_from_raw_value(raw_value: &[u8]) -> Result<Self>
386 where
387 Self: Sized;
388
389 fn try_as_raw_value(&self) -> Result<Vec<u8>>;
390}
391
392pub type TableMetadataManagerRef = Arc<TableMetadataManager>;
393
394pub struct TableMetadataManager {
395 table_name_manager: TableNameManager,
396 table_info_manager: TableInfoManager,
397 view_info_manager: ViewInfoManager,
398 datanode_table_manager: DatanodeTableManager,
399 catalog_manager: CatalogManager,
400 schema_manager: SchemaManager,
401 table_route_manager: TableRouteManager,
402 table_repart_manager: TableRepartManager,
403 tombstone_manager: TombstoneManager,
404 topic_name_manager: TopicNameManager,
405 topic_region_manager: TopicRegionManager,
406 kv_backend: KvBackendRef,
407}
408
409#[macro_export]
410macro_rules! ensure_values {
411 ($got:expr, $expected_value:expr, $name:expr) => {
412 ensure!(
413 $got == $expected_value,
414 error::UnexpectedSnafu {
415 err_msg: format!(
416 "Reads the different value: {:?} during {}, expected: {:?}",
417 $got, $name, $expected_value
418 )
419 }
420 );
421 };
422}
423
424pub struct DeserializedValueWithBytes<T: DeserializeOwned + Serialize> {
434 bytes: Bytes,
436 inner: T,
438}
439
440#[derive(Debug, Clone, PartialEq, Eq)]
441pub struct DroppedTableName {
442 pub table_id: TableId,
444 pub table_name: TableName,
446 pub dropped_at: Option<i64>,
448 pub retention_expires_at: Option<i64>,
450 pub drop_generation: Option<String>,
452 pub purging: bool,
454}
455
456#[derive(Debug, Clone)]
457pub struct DroppedTableMetadata {
458 pub table_id: TableId,
460 pub table_name: TableName,
462 pub table_info_value: TableInfoValue,
464 pub table_route_value: TableRouteValue,
466 pub region_wal_options: HashMap<RegionNumber, WalOptions>,
468 pub dropped_at: Option<i64>,
470 pub retention_expires_at: Option<i64>,
472 pub drop_generation: Option<String>,
474}
475
476pub struct DroppedTableLifecycle<'a> {
478 pub dropped_at: Option<i64>,
479 pub retention_expires_at: Option<i64>,
480 pub drop_generation: Option<&'a str>,
481}
482
483const DROPPED_AT_KEY_PREFIX: &str = "__dropped_at";
484const RETENTION_EXPIRES_AT_KEY_PREFIX: &str = "__retention_expires_at";
485const DROP_GENERATION_KEY_PREFIX: &str = "__drop_generation";
486const PURGING_KEY_PREFIX: &str = "__purging";
487
488pub(crate) fn dropped_at_key(table_id: TableId) -> Vec<u8> {
489 format!("{DROPPED_AT_KEY_PREFIX}/{table_id}").into_bytes()
490}
491
492pub(crate) fn retention_expires_at_key(table_id: TableId) -> Vec<u8> {
493 format!("{RETENTION_EXPIRES_AT_KEY_PREFIX}/{table_id}").into_bytes()
494}
495
496pub(crate) fn drop_generation_key(table_id: TableId) -> Vec<u8> {
497 format!("{DROP_GENERATION_KEY_PREFIX}/{table_id}").into_bytes()
498}
499
500pub(crate) fn purging_key(table_id: TableId) -> Vec<u8> {
501 format!("{PURGING_KEY_PREFIX}/{table_id}").into_bytes()
502}
503
504impl<T: DeserializeOwned + Serialize> Deref for DeserializedValueWithBytes<T> {
505 type Target = T;
506
507 fn deref(&self) -> &Self::Target {
508 &self.inner
509 }
510}
511
512impl<T: DeserializeOwned + Serialize> DerefMut for DeserializedValueWithBytes<T> {
513 fn deref_mut(&mut self) -> &mut Self::Target {
514 &mut self.inner
515 }
516}
517
518impl<T: DeserializeOwned + Serialize + Debug> Debug for DeserializedValueWithBytes<T> {
519 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
520 write!(
521 f,
522 "DeserializedValueWithBytes(inner: {:?}, bytes: {:?})",
523 self.inner, self.bytes
524 )
525 }
526}
527
528impl<T: DeserializeOwned + Serialize> Serialize for DeserializedValueWithBytes<T> {
529 fn serialize<S>(&self, serializer: S) -> std::result::Result<S::Ok, S::Error>
533 where
534 S: serde::Serializer,
535 {
536 serializer.serialize_str(&String::from_utf8_lossy(&self.bytes))
539 }
540}
541
542impl<'de, T: DeserializeOwned + Serialize + MetadataValue> Deserialize<'de>
543 for DeserializedValueWithBytes<T>
544{
545 fn deserialize<D>(deserializer: D) -> std::result::Result<Self, D::Error>
549 where
550 D: serde::Deserializer<'de>,
551 {
552 let buf = String::deserialize(deserializer)?;
553 let bytes = Bytes::from(buf);
554
555 let value = DeserializedValueWithBytes::from_inner_bytes(bytes)
556 .map_err(|err| serde::de::Error::custom(err.to_string()))?;
557
558 Ok(value)
559 }
560}
561
562impl<T: Serialize + DeserializeOwned + Clone> Clone for DeserializedValueWithBytes<T> {
563 fn clone(&self) -> Self {
564 Self {
565 bytes: self.bytes.clone(),
566 inner: self.inner.clone(),
567 }
568 }
569}
570
571impl<T: Serialize + DeserializeOwned + MetadataValue> DeserializedValueWithBytes<T> {
572 pub fn from_inner_bytes(bytes: Bytes) -> Result<Self> {
575 let inner = T::try_from_raw_value(&bytes)?;
576 Ok(Self { bytes, inner })
577 }
578
579 pub fn from_inner_slice(bytes: &[u8]) -> Result<Self> {
582 Self::from_inner_bytes(Bytes::copy_from_slice(bytes))
583 }
584
585 pub fn into_inner(self) -> T {
586 self.inner
587 }
588
589 pub fn get_inner_ref(&self) -> &T {
590 &self.inner
591 }
592
593 pub fn get_raw_bytes(&self) -> Vec<u8> {
595 self.bytes.to_vec()
596 }
597
598 #[cfg(any(test, feature = "testing"))]
599 pub fn from_inner(inner: T) -> Self {
600 let bytes = serde_json::to_vec(&inner).unwrap();
601
602 Self {
603 bytes: Bytes::from(bytes),
604 inner,
605 }
606 }
607}
608
609impl TableMetadataManager {
610 pub fn new(kv_backend: KvBackendRef) -> Self {
611 TableMetadataManager {
612 table_name_manager: TableNameManager::new(kv_backend.clone()),
613 table_info_manager: TableInfoManager::new(kv_backend.clone()),
614 view_info_manager: ViewInfoManager::new(kv_backend.clone()),
615 datanode_table_manager: DatanodeTableManager::new(kv_backend.clone()),
616 catalog_manager: CatalogManager::new(kv_backend.clone()),
617 schema_manager: SchemaManager::new(kv_backend.clone()),
618 table_route_manager: TableRouteManager::new(kv_backend.clone()),
619 table_repart_manager: TableRepartManager::new(kv_backend.clone()),
620 tombstone_manager: TombstoneManager::new(kv_backend.clone()),
621 topic_name_manager: TopicNameManager::new(kv_backend.clone()),
622 topic_region_manager: TopicRegionManager::new(kv_backend.clone()),
623 kv_backend,
624 }
625 }
626
627 pub fn new_with_custom_tombstone_prefix(
629 kv_backend: KvBackendRef,
630 tombstone_prefix: &str,
631 ) -> Self {
632 Self {
633 table_name_manager: TableNameManager::new(kv_backend.clone()),
634 table_info_manager: TableInfoManager::new(kv_backend.clone()),
635 view_info_manager: ViewInfoManager::new(kv_backend.clone()),
636 datanode_table_manager: DatanodeTableManager::new(kv_backend.clone()),
637 catalog_manager: CatalogManager::new(kv_backend.clone()),
638 schema_manager: SchemaManager::new(kv_backend.clone()),
639 table_route_manager: TableRouteManager::new(kv_backend.clone()),
640 table_repart_manager: TableRepartManager::new(kv_backend.clone()),
641 tombstone_manager: TombstoneManager::new_with_prefix(
642 kv_backend.clone(),
643 tombstone_prefix,
644 ),
645 topic_name_manager: TopicNameManager::new(kv_backend.clone()),
646 topic_region_manager: TopicRegionManager::new(kv_backend.clone()),
647 kv_backend,
648 }
649 }
650
651 pub async fn init(&self) -> Result<()> {
652 let catalog_name = CatalogNameKey::new(DEFAULT_CATALOG_NAME);
653
654 self.catalog_manager().create(catalog_name, true).await?;
655
656 let internal_schemas = [
657 DEFAULT_SCHEMA_NAME,
658 INFORMATION_SCHEMA_NAME,
659 DEFAULT_PRIVATE_SCHEMA_NAME,
660 ];
661
662 for schema_name in internal_schemas {
663 let schema_key = SchemaNameKey::new(DEFAULT_CATALOG_NAME, schema_name);
664
665 self.schema_manager().create(schema_key, None, true).await?;
666 }
667
668 Ok(())
669 }
670
671 pub fn table_name_manager(&self) -> &TableNameManager {
672 &self.table_name_manager
673 }
674
675 pub fn table_info_manager(&self) -> &TableInfoManager {
676 &self.table_info_manager
677 }
678
679 pub fn view_info_manager(&self) -> &ViewInfoManager {
680 &self.view_info_manager
681 }
682
683 pub fn datanode_table_manager(&self) -> &DatanodeTableManager {
684 &self.datanode_table_manager
685 }
686
687 pub fn catalog_manager(&self) -> &CatalogManager {
688 &self.catalog_manager
689 }
690
691 pub fn schema_manager(&self) -> &SchemaManager {
692 &self.schema_manager
693 }
694
695 pub fn table_route_manager(&self) -> &TableRouteManager {
696 &self.table_route_manager
697 }
698
699 pub fn table_repart_manager(&self) -> &TableRepartManager {
700 &self.table_repart_manager
701 }
702
703 pub fn topic_name_manager(&self) -> &TopicNameManager {
704 &self.topic_name_manager
705 }
706
707 pub fn topic_region_manager(&self) -> &TopicRegionManager {
708 &self.topic_region_manager
709 }
710
711 pub fn kv_backend(&self) -> &KvBackendRef {
712 &self.kv_backend
713 }
714
715 pub async fn get_full_table_info(
716 &self,
717 table_id: TableId,
718 ) -> Result<(
719 Option<DeserializedValueWithBytes<TableInfoValue>>,
720 Option<DeserializedValueWithBytes<TableRouteValue>>,
721 )> {
722 let table_info_key = TableInfoKey::new(table_id);
723 let table_route_key = TableRouteKey::new(table_id);
724 let (table_info_txn, table_info_filter) = table_info_key.build_get_op();
725 let (table_route_txn, table_route_filter) = table_route_key.build_get_op();
726
727 let txn = Txn::new().and_then(vec![table_info_txn, table_route_txn]);
728 let mut res = self.kv_backend.txn(txn).await?;
729 let mut set = TxnOpGetResponseSet::from(&mut res.responses);
730 let table_info_value = TxnOpGetResponseSet::decode_with(table_info_filter)(&mut set)?;
731 let mut table_route_value = TxnOpGetResponseSet::decode_with(table_route_filter)(&mut set)?;
732 if let Some(table_route_value) = &mut table_route_value {
733 self.table_route_manager()
734 .table_route_storage()
735 .remap_table_route(table_route_value)
736 .await?;
737 }
738 Ok((table_info_value, table_route_value))
739 }
740
741 pub async fn create_view_metadata(
751 &self,
752 view_info: TableInfo,
753 raw_logical_plan: Vec<u8>,
754 table_names: HashSet<TableName>,
755 columns: Vec<String>,
756 plan_columns: Vec<String>,
757 definition: String,
758 ) -> Result<()> {
759 let view_id = view_info.ident.table_id;
760
761 let view_name = TableNameKey::new(
763 &view_info.catalog_name,
764 &view_info.schema_name,
765 &view_info.name,
766 );
767 let create_table_name_txn = self
768 .table_name_manager()
769 .build_create_txn(&view_name, view_id)?;
770
771 let table_info_value = TableInfoValue::new(view_info);
773
774 let (create_table_info_txn, on_create_table_info_failure) = self
775 .table_info_manager()
776 .build_create_txn(view_id, &table_info_value)?;
777
778 let view_info_value = ViewInfoValue::new(
780 raw_logical_plan.into(),
781 table_names,
782 columns,
783 plan_columns,
784 definition,
785 );
786 let (create_view_info_txn, on_create_view_info_failure) = self
787 .view_info_manager()
788 .build_create_txn(view_id, &view_info_value)?;
789
790 let txn = Txn::merge_all(vec![
791 create_table_name_txn,
792 create_table_info_txn,
793 create_view_info_txn,
794 ]);
795
796 let mut r = self.kv_backend.txn(txn).await?;
797
798 if !r.succeeded {
800 let mut set = TxnOpGetResponseSet::from(&mut r.responses);
801 let remote_table_info = on_create_table_info_failure(&mut set)?
802 .context(error::UnexpectedSnafu {
803 err_msg: "Reads the empty table info in comparing operation of creating table metadata",
804 })?
805 .into_inner();
806
807 let remote_view_info = on_create_view_info_failure(&mut set)?
808 .context(error::UnexpectedSnafu {
809 err_msg: "Reads the empty view info in comparing operation of creating view metadata",
810 })?
811 .into_inner();
812
813 let op_name = "the creating view metadata";
814 ensure_values!(remote_table_info, table_info_value, op_name);
815 ensure_values!(remote_view_info, view_info_value, op_name);
816 }
817
818 Ok(())
819 }
820
821 pub async fn create_table_metadata(
824 &self,
825 table_info: TableInfo,
826 table_route_value: TableRouteValue,
827 region_wal_options: RegionWalOptions,
828 ) -> Result<()> {
829 let table_id = table_info.ident.table_id;
830 let engine = table_info.meta.engine.clone();
831
832 let table_name = TableNameKey::new(
834 &table_info.catalog_name,
835 &table_info.schema_name,
836 &table_info.name,
837 );
838 let create_table_name_txn = self
839 .table_name_manager()
840 .build_create_txn(&table_name, table_id)?;
841
842 let region_options = table_info.to_region_options();
843 let table_info_value = TableInfoValue::new(table_info);
845 let (create_table_info_txn, on_create_table_info_failure) = self
846 .table_info_manager()
847 .build_create_txn(table_id, &table_info_value)?;
848
849 let (create_table_route_txn, on_create_table_route_failure) = self
850 .table_route_manager()
851 .table_route_storage()
852 .build_create_txn(table_id, &table_route_value)?;
853
854 let create_topic_region_txn = self
855 .topic_region_manager
856 .build_create_txn(table_id, ®ion_wal_options)?;
857
858 let mut txn = Txn::merge_all(vec![
859 create_table_name_txn,
860 create_table_info_txn,
861 create_table_route_txn,
862 create_topic_region_txn,
863 ]);
864
865 if let TableRouteValue::Physical(x) = &table_route_value {
866 let region_storage_path = table_info_value.region_storage_path();
867 let create_datanode_table_txn = self.datanode_table_manager().build_create_txn(
868 table_id,
869 &engine,
870 ®ion_storage_path,
871 region_options,
872 region_wal_options,
873 region_distribution(&x.region_routes),
874 )?;
875 txn = txn.merge(create_datanode_table_txn);
876 }
877
878 let mut r = self.kv_backend.txn(txn).await?;
879
880 if !r.succeeded {
882 let mut set = TxnOpGetResponseSet::from(&mut r.responses);
883 let remote_table_info = on_create_table_info_failure(&mut set)?
884 .context(error::UnexpectedSnafu {
885 err_msg: "Reads the empty table info in comparing operation of creating table metadata",
886 })?
887 .into_inner();
888
889 let remote_table_route = on_create_table_route_failure(&mut set)?
890 .context(error::UnexpectedSnafu {
891 err_msg: "Reads the empty table route in comparing operation of creating table metadata",
892 })?
893 .into_inner();
894
895 let op_name = "the creating table metadata";
896 ensure_values!(remote_table_info, table_info_value, op_name);
897 ensure_values!(remote_table_route, table_route_value, op_name);
898 }
899
900 Ok(())
901 }
902
903 pub fn create_logical_tables_metadata_chunk_size(&self) -> usize {
904 self.kv_backend.max_txn_ops() / 3
907 }
908
909 pub async fn create_logical_tables_metadata(
911 &self,
912 tables_data: Vec<(TableInfo, TableRouteValue)>,
913 ) -> Result<()> {
914 let len = tables_data.len();
915 let mut txns = Vec::with_capacity(3 * len);
916 struct OnFailure<F1, R1, F2, R2>
917 where
918 F1: FnOnce(&mut TxnOpGetResponseSet) -> R1,
919 F2: FnOnce(&mut TxnOpGetResponseSet) -> R2,
920 {
921 table_info_value: TableInfoValue,
922 on_create_table_info_failure: F1,
923 table_route_value: TableRouteValue,
924 on_create_table_route_failure: F2,
925 }
926 let mut on_failures = Vec::with_capacity(len);
927 for (table_info, table_route_value) in tables_data {
928 let table_id = table_info.ident.table_id;
929
930 let table_name = TableNameKey::new(
932 &table_info.catalog_name,
933 &table_info.schema_name,
934 &table_info.name,
935 );
936 let create_table_name_txn = self
937 .table_name_manager()
938 .build_create_txn(&table_name, table_id)?;
939 txns.push(create_table_name_txn);
940
941 let table_info_value = TableInfoValue::new(table_info);
943 let (create_table_info_txn, on_create_table_info_failure) =
944 self.table_info_manager()
945 .build_create_txn(table_id, &table_info_value)?;
946 txns.push(create_table_info_txn);
947
948 let (create_table_route_txn, on_create_table_route_failure) = self
949 .table_route_manager()
950 .table_route_storage()
951 .build_create_txn(table_id, &table_route_value)?;
952 txns.push(create_table_route_txn);
953
954 on_failures.push(OnFailure {
955 table_info_value,
956 on_create_table_info_failure,
957 table_route_value,
958 on_create_table_route_failure,
959 });
960 }
961
962 let txn = Txn::merge_all(txns);
963 let mut r = self.kv_backend.txn(txn).await?;
964
965 if !r.succeeded {
967 let mut set = TxnOpGetResponseSet::from(&mut r.responses);
968 for on_failure in on_failures {
969 let remote_table_info = (on_failure.on_create_table_info_failure)(&mut set)?
970 .context(error::UnexpectedSnafu {
971 err_msg: "Reads the empty table info in comparing operation of creating table metadata",
972 })?
973 .into_inner();
974
975 let remote_table_route = (on_failure.on_create_table_route_failure)(&mut set)?
976 .context(error::UnexpectedSnafu {
977 err_msg: "Reads the empty table route in comparing operation of creating table metadata",
978 })?
979 .into_inner();
980
981 let op_name = "the creating logical tables metadata";
982 ensure_values!(remote_table_info, on_failure.table_info_value, op_name);
983 ensure_values!(remote_table_route, on_failure.table_route_value, op_name);
984 }
985 }
986
987 Ok(())
988 }
989
990 fn table_metadata_keys(
991 &self,
992 table_id: TableId,
993 table_name: &TableName,
994 table_route_value: &TableRouteValue,
995 region_wal_options: &HashMap<RegionNumber, WalOptions>,
996 ) -> Result<Vec<Vec<u8>>> {
997 let datanode_ids = if table_route_value.is_physical() {
999 region_distribution(table_route_value.region_routes()?)
1000 .into_keys()
1001 .collect()
1002 } else {
1003 vec![]
1004 };
1005 let mut keys = Vec::with_capacity(3 + datanode_ids.len());
1006 let table_name = TableNameKey::new(
1007 &table_name.catalog_name,
1008 &table_name.schema_name,
1009 &table_name.table_name,
1010 );
1011 let table_info_key = TableInfoKey::new(table_id);
1012 let table_route_key = TableRouteKey::new(table_id);
1013 let table_repart_key = TableRepartKey::new(table_id);
1014 let datanode_table_keys = datanode_ids
1015 .into_iter()
1016 .map(|datanode_id| DatanodeTableKey::new(datanode_id, table_id))
1017 .collect::<HashSet<_>>();
1018 let topic_region_map = self
1019 .topic_region_manager
1020 .get_topic_region_mapping(table_id, region_wal_options);
1021 let topic_region_keys = topic_region_map
1022 .iter()
1023 .map(|(region_id, topic)| TopicRegionKey::new(*region_id, topic))
1024 .collect::<Vec<_>>();
1025 keys.push(table_name.to_bytes());
1026 keys.push(table_info_key.to_bytes());
1027 keys.push(table_route_key.to_bytes());
1028 keys.push(table_repart_key.to_bytes());
1029 for key in &datanode_table_keys {
1030 keys.push(key.to_bytes());
1031 }
1032 for key in topic_region_keys {
1033 keys.push(key.to_bytes());
1034 }
1035 Ok(keys)
1036 }
1037
1038 pub async fn delete_table_metadata(
1041 &self,
1042 table_id: TableId,
1043 table_name: &TableName,
1044 table_route_value: &TableRouteValue,
1045 region_wal_options: &HashMap<RegionNumber, WalOptions>,
1046 dropped_at: Option<i64>,
1047 ) -> Result<()> {
1048 self.delete_table_metadata_with_retention(
1049 table_id,
1050 table_name,
1051 table_route_value,
1052 region_wal_options,
1053 dropped_at,
1054 None,
1055 )
1056 .await
1057 }
1058
1059 pub async fn delete_table_metadata_with_retention(
1061 &self,
1062 table_id: TableId,
1063 table_name: &TableName,
1064 table_route_value: &TableRouteValue,
1065 region_wal_options: &HashMap<RegionNumber, WalOptions>,
1066 dropped_at: Option<i64>,
1067 retention_expires_at: Option<i64>,
1068 ) -> Result<()> {
1069 self.delete_table_metadata_with_retention_and_generation(
1070 table_id,
1071 table_name,
1072 table_route_value,
1073 region_wal_options,
1074 DroppedTableLifecycle {
1075 dropped_at,
1076 retention_expires_at,
1077 drop_generation: None,
1078 },
1079 )
1080 .await
1081 }
1082
1083 pub async fn delete_table_metadata_with_retention_and_generation(
1085 &self,
1086 table_id: TableId,
1087 table_name: &TableName,
1088 table_route_value: &TableRouteValue,
1089 region_wal_options: &HashMap<RegionNumber, WalOptions>,
1090 lifecycle: DroppedTableLifecycle<'_>,
1091 ) -> Result<()> {
1092 let keys =
1093 self.table_metadata_keys(table_id, table_name, table_route_value, region_wal_options)?;
1094 if let Some(dropped_at) = lifecycle.dropped_at {
1095 let mut markers = vec![(
1096 dropped_at_key(table_id),
1097 dropped_at.to_string().into_bytes(),
1098 )];
1099 if let Some(retention_expires_at) = lifecycle.retention_expires_at {
1100 markers.push((
1101 retention_expires_at_key(table_id),
1102 retention_expires_at.to_string().into_bytes(),
1103 ));
1104 }
1105 if let Some(drop_generation) = lifecycle.drop_generation {
1106 markers.push((
1107 drop_generation_key(table_id),
1108 drop_generation.as_bytes().to_vec(),
1109 ));
1110 }
1111 self.tombstone_manager
1112 .create_with_markers(keys, markers)
1113 .await
1114 .map(|_| ())
1115 } else {
1116 self.tombstone_manager.create(keys).await.map(|_| ())
1117 }
1118 }
1119
1120 pub async fn list_dropped_tables(&self) -> Result<Vec<DroppedTableName>> {
1122 self.collect_dropped_tables(self.tombstone_manager.tombstoned_table_names())
1123 .await
1124 }
1125
1126 pub async fn list_dropped_tables_by_catalog(
1128 &self,
1129 catalog: &str,
1130 ) -> Result<Vec<DroppedTableName>> {
1131 self.collect_dropped_tables(
1132 self.tombstone_manager
1133 .tombstoned_table_names_by_catalog(catalog),
1134 )
1135 .await
1136 }
1137
1138 async fn collect_dropped_tables(
1139 &self,
1140 mut stream: BoxStream<'static, Result<KeyValue>>,
1141 ) -> Result<Vec<DroppedTableName>> {
1142 let mut dropped_tables = Vec::new();
1143
1144 while let Some(kv) = stream.try_next().await? {
1145 let raw_key = self.tombstone_manager.strip_tombstone_prefix(&kv.key)?;
1146 let table_name = TableNameKey::from_bytes(raw_key)?.into();
1147 let table_id = TableNameValue::try_from_raw_value(&kv.value)?.table_id();
1148 dropped_tables.push(DroppedTableName {
1149 table_id,
1150 table_name,
1151 dropped_at: None,
1152 retention_expires_at: None,
1153 drop_generation: None,
1154 purging: false,
1155 });
1156 }
1157 if dropped_tables.is_empty() {
1158 return Ok(dropped_tables);
1159 }
1160
1161 let marker_keys = dropped_tables
1162 .iter()
1163 .flat_map(|table| {
1164 vec![
1165 dropped_at_key(table.table_id),
1166 retention_expires_at_key(table.table_id),
1167 drop_generation_key(table.table_id),
1168 purging_key(table.table_id),
1169 ]
1170 })
1171 .collect::<Vec<_>>();
1172 let marker_values = self.tombstone_manager.batch_get(&marker_keys).await?;
1173 for table in &mut dropped_tables {
1174 table.dropped_at = Self::parse_dropped_at(
1175 table.table_id,
1176 marker_values.get(&dropped_at_key(table.table_id)),
1177 )?;
1178 table.retention_expires_at = Self::parse_retention_expires_at(
1179 table.table_id,
1180 marker_values.get(&retention_expires_at_key(table.table_id)),
1181 )?;
1182 table.drop_generation = Self::parse_drop_generation(
1183 marker_values.get(&drop_generation_key(table.table_id)),
1184 );
1185 table.purging = marker_values.contains_key(&purging_key(table.table_id));
1186 }
1187
1188 Ok(dropped_tables)
1189 }
1190
1191 pub async fn get_dropped_table(
1193 &self,
1194 table_name: &TableName,
1195 ) -> Result<Option<DroppedTableMetadata>> {
1196 let table_name_key = TableNameKey::from(table_name);
1197 let Some(kv) = self
1198 .tombstone_manager
1199 .get(&table_name_key.to_bytes())
1200 .await?
1201 else {
1202 return Ok(None);
1203 };
1204
1205 let table_id = TableNameValue::try_from_raw_value(&kv.value)?.table_id();
1206 self.get_dropped_table_metadata(table_id, table_name.clone())
1207 .await
1208 }
1209
1210 pub async fn get_dropped_table_by_id(
1212 &self,
1213 table_id: TableId,
1214 ) -> Result<Option<DroppedTableMetadata>> {
1215 self.get_dropped_table_metadata(table_id, None).await
1216 }
1217
1218 pub async fn is_dropped_table_purging(&self, table_id: TableId) -> Result<bool> {
1220 self.dropped_table_purge_claim(table_id)
1221 .await
1222 .map(|claim| claim.is_some())
1223 }
1224
1225 pub async fn dropped_table_purge_claim(&self, table_id: TableId) -> Result<Option<String>> {
1227 self.tombstone_manager
1228 .get(&purging_key(table_id))
1229 .await
1230 .map(|marker| marker.map(|marker| String::from_utf8_lossy(&marker.value).into_owned()))
1231 }
1232
1233 pub async fn mark_dropped_table_purging(
1236 &self,
1237 table_id: TableId,
1238 drop_generation: Option<&str>,
1239 ) -> Result<()> {
1240 self.kv_backend
1241 .put(
1242 PutRequest::new()
1243 .with_key(self.tombstone_manager.to_tombstone(&purging_key(table_id)))
1244 .with_value(drop_generation.unwrap_or_default()),
1245 )
1246 .await?;
1247 Ok(())
1248 }
1249
1250 pub async fn delete_table_metadata_tombstone(
1253 &self,
1254 table_id: TableId,
1255 table_name: &TableName,
1256 table_route_value: &TableRouteValue,
1257 region_wal_options: &HashMap<RegionNumber, WalOptions>,
1258 ) -> Result<()> {
1259 let table_metadata_keys =
1260 self.table_metadata_keys(table_id, table_name, table_route_value, region_wal_options)?;
1261 self.tombstone_manager
1262 .delete_with_markers(
1263 table_metadata_keys,
1264 vec![
1265 dropped_at_key(table_id),
1266 retention_expires_at_key(table_id),
1267 drop_generation_key(table_id),
1268 purging_key(table_id),
1269 ],
1270 )
1271 .await
1272 .map(|_| ())
1273 }
1274
1275 pub async fn restore_table_metadata(
1278 &self,
1279 table_id: TableId,
1280 table_name: &TableName,
1281 table_route_value: &TableRouteValue,
1282 region_wal_options: &HashMap<RegionNumber, WalOptions>,
1283 ) -> Result<()> {
1284 let keys =
1285 self.table_metadata_keys(table_id, table_name, table_route_value, region_wal_options)?;
1286 self.tombstone_manager
1287 .restore_with_markers(
1288 keys,
1289 vec![
1290 dropped_at_key(table_id),
1291 retention_expires_at_key(table_id),
1292 drop_generation_key(table_id),
1293 ],
1294 )
1295 .await
1296 .map(|_| ())
1297 }
1298
1299 pub async fn destroy_table_metadata(
1302 &self,
1303 table_id: TableId,
1304 table_name: &TableName,
1305 table_route_value: &TableRouteValue,
1306 region_wal_options: &HashMap<RegionNumber, WalOptions>,
1307 ) -> Result<()> {
1308 let keys =
1309 self.table_metadata_keys(table_id, table_name, table_route_value, region_wal_options)?;
1310 let _ = self
1311 .kv_backend
1312 .batch_delete(BatchDeleteRequest::new().with_keys(keys))
1313 .await?;
1314 Ok(())
1315 }
1316
1317 async fn get_dropped_table_metadata<T>(
1319 &self,
1320 table_id: TableId,
1321 table_name: T,
1322 ) -> Result<Option<DroppedTableMetadata>>
1323 where
1324 T: Into<Option<TableName>>,
1325 {
1326 let table_info_key = TableInfoKey::new(table_id);
1327 let Some(table_info_kv) = self
1328 .tombstone_manager
1329 .get(&table_info_key.to_bytes())
1330 .await?
1331 else {
1332 return Ok(None);
1333 };
1334
1335 let table_info_value = TableInfoValue::try_from_raw_value(&table_info_kv.value)?;
1336 let table_name = table_name
1337 .into()
1338 .unwrap_or_else(|| table_info_value.table_name());
1339
1340 let table_route_key = TableRouteKey::new(table_id);
1341 let table_route_kv = self
1342 .tombstone_manager
1343 .get(&table_route_key.to_bytes())
1344 .await?
1345 .with_context(|| error::UnexpectedSnafu {
1346 err_msg: format!("Missing tombstoned table route metadata for table id {table_id}"),
1347 })?;
1348 let mut table_route_value = TableRouteValue::try_from_raw_value(&table_route_kv.value)?;
1349 self.table_route_manager
1350 .table_route_storage()
1351 .remap_table_route(&mut table_route_value)
1352 .await?;
1353
1354 let region_wal_options = self
1355 .dropped_region_wal_options(table_id, &table_route_value)
1356 .await?;
1357 let dropped_at_key = dropped_at_key(table_id);
1358 let retention_expires_at_key = retention_expires_at_key(table_id);
1359 let drop_generation_key = drop_generation_key(table_id);
1360 let marker_values = self
1361 .tombstone_manager
1362 .batch_get(&[
1363 dropped_at_key.clone(),
1364 retention_expires_at_key.clone(),
1365 drop_generation_key.clone(),
1366 ])
1367 .await?;
1368 let dropped_at = Self::parse_dropped_at(table_id, marker_values.get(&dropped_at_key))?;
1369 let retention_expires_at = Self::parse_retention_expires_at(
1370 table_id,
1371 marker_values.get(&retention_expires_at_key),
1372 )?;
1373 let drop_generation = Self::parse_drop_generation(marker_values.get(&drop_generation_key));
1374
1375 Ok(Some(DroppedTableMetadata {
1376 table_id,
1377 table_name,
1378 table_info_value,
1379 table_route_value,
1380 region_wal_options,
1381 dropped_at,
1382 retention_expires_at,
1383 drop_generation,
1384 }))
1385 }
1386
1387 fn parse_dropped_at(table_id: TableId, kv: Option<&KeyValue>) -> Result<Option<i64>> {
1388 Self::parse_timestamp_marker(table_id, "dropped timestamp", kv)
1389 }
1390
1391 fn parse_retention_expires_at(table_id: TableId, kv: Option<&KeyValue>) -> Result<Option<i64>> {
1392 Self::parse_timestamp_marker(table_id, "retention deadline", kv)
1393 }
1394
1395 fn parse_drop_generation(kv: Option<&KeyValue>) -> Option<String> {
1396 kv.map(|kv| String::from_utf8_lossy(&kv.value).into_owned())
1397 }
1398
1399 fn parse_timestamp_marker(
1400 table_id: TableId,
1401 description: &str,
1402 kv: Option<&KeyValue>,
1403 ) -> Result<Option<i64>> {
1404 let Some(kv) = kv else {
1405 return Ok(None);
1406 };
1407 let value = String::from_utf8_lossy(&kv.value);
1408 value.parse().map(Some).map_err(|err| {
1409 error::UnexpectedSnafu {
1410 err_msg: format!("Invalid {description} '{value}' for table id {table_id}: {err}"),
1411 }
1412 .build()
1413 })
1414 }
1415
1416 async fn dropped_region_wal_options(
1418 &self,
1419 table_id: TableId,
1420 table_route_value: &TableRouteValue,
1421 ) -> Result<HashMap<RegionNumber, WalOptions>> {
1422 let mut region_wal_options = HashMap::new();
1423 let Some(region_routes) = table_route_value.region_routes().ok() else {
1424 return Ok(region_wal_options);
1425 };
1426 let datanode_table_keys = region_distribution(region_routes)
1427 .into_keys()
1428 .map(|datanode_id| DatanodeTableKey::new(datanode_id, table_id))
1429 .collect::<Vec<_>>();
1430 let datanode_table_key_bytes = datanode_table_keys
1431 .iter()
1432 .map(|key| key.to_bytes())
1433 .collect::<Vec<_>>();
1434 let datanode_table_values = self
1435 .tombstone_manager
1436 .batch_get(&datanode_table_key_bytes)
1437 .await?;
1438
1439 for datanode_table_key in datanode_table_keys {
1440 let Some(kv) = datanode_table_values.get(&datanode_table_key.to_bytes()) else {
1441 continue;
1442 };
1443
1444 let datanode_table_value = DatanodeTableValue::try_from_raw_value(&kv.value)?;
1445 for (region_number, wal_options) in &datanode_table_value.region_info.region_wal_options
1446 {
1447 region_wal_options.insert(*region_number, wal_options.clone());
1448 }
1449 }
1450
1451 Ok(region_wal_options)
1452 }
1453
1454 fn view_info_keys(&self, view_id: TableId, view_name: &TableName) -> Result<Vec<Vec<u8>>> {
1455 let mut keys = Vec::with_capacity(3);
1456 let view_name = TableNameKey::new(
1457 &view_name.catalog_name,
1458 &view_name.schema_name,
1459 &view_name.table_name,
1460 );
1461 let table_info_key = TableInfoKey::new(view_id);
1462 let view_info_key = ViewInfoKey::new(view_id);
1463 keys.push(view_name.to_bytes());
1464 keys.push(table_info_key.to_bytes());
1465 keys.push(view_info_key.to_bytes());
1466
1467 Ok(keys)
1468 }
1469
1470 pub async fn destroy_view_info(&self, view_id: TableId, view_name: &TableName) -> Result<()> {
1473 let keys = self.view_info_keys(view_id, view_name)?;
1474 let _ = self
1475 .kv_backend
1476 .batch_delete(BatchDeleteRequest::new().with_keys(keys))
1477 .await?;
1478 Ok(())
1479 }
1480
1481 pub async fn rename_table(
1485 &self,
1486 current_table_info_value: &DeserializedValueWithBytes<TableInfoValue>,
1487 new_table_name: String,
1488 ) -> Result<()> {
1489 let current_table_info = ¤t_table_info_value.table_info;
1490 let table_id = current_table_info.ident.table_id;
1491
1492 let table_name_key = TableNameKey::new(
1493 ¤t_table_info.catalog_name,
1494 ¤t_table_info.schema_name,
1495 ¤t_table_info.name,
1496 );
1497
1498 let new_table_name_key = TableNameKey::new(
1499 ¤t_table_info.catalog_name,
1500 ¤t_table_info.schema_name,
1501 &new_table_name,
1502 );
1503
1504 let update_table_name_txn = self.table_name_manager().build_update_txn(
1506 &table_name_key,
1507 &new_table_name_key,
1508 table_id,
1509 )?;
1510
1511 let new_table_info_value = current_table_info_value
1512 .inner
1513 .with_update(move |table_info| {
1514 table_info.name = new_table_name;
1515 });
1516
1517 let (update_table_info_txn, on_update_table_info_failure) = self
1519 .table_info_manager()
1520 .build_update_txn(table_id, current_table_info_value, &new_table_info_value)?;
1521
1522 let txn = Txn::merge_all(vec![update_table_name_txn, update_table_info_txn]);
1523
1524 let mut r = self.kv_backend.txn(txn).await?;
1525
1526 if !r.succeeded {
1528 let mut set = TxnOpGetResponseSet::from(&mut r.responses);
1529 let remote_table_info = on_update_table_info_failure(&mut set)?
1530 .context(error::UnexpectedSnafu {
1531 err_msg: "Reads the empty table info in comparing operation of the rename table metadata",
1532 })?
1533 .into_inner();
1534
1535 let op_name = "the renaming table metadata";
1536 ensure_values!(remote_table_info, new_table_info_value, op_name);
1537 }
1538
1539 Ok(())
1540 }
1541
1542 pub async fn update_table_info(
1546 &self,
1547 current_table_info_value: &DeserializedValueWithBytes<TableInfoValue>,
1548 region_distribution: Option<RegionDistribution>,
1549 new_table_info: TableInfo,
1550 ) -> Result<()> {
1551 let table_id = current_table_info_value.table_info.ident.table_id;
1552 let new_table_info_value = current_table_info_value.update(new_table_info);
1553
1554 let (update_table_info_txn, on_update_table_info_failure) = self
1556 .table_info_manager()
1557 .build_update_txn(table_id, current_table_info_value, &new_table_info_value)?;
1558
1559 let txn = if let Some(region_distribution) = region_distribution {
1560 let new_region_options = new_table_info_value.table_info.to_region_options();
1562 let update_datanode_table_options_txn = self
1563 .datanode_table_manager
1564 .build_update_table_options_txn(table_id, region_distribution, new_region_options)
1565 .await?;
1566 Txn::merge_all([update_table_info_txn, update_datanode_table_options_txn])
1567 } else {
1568 update_table_info_txn
1569 };
1570
1571 let mut r = self.kv_backend.txn(txn).await?;
1572 if !r.succeeded {
1574 let mut set = TxnOpGetResponseSet::from(&mut r.responses);
1575 let remote_table_info = on_update_table_info_failure(&mut set)?
1576 .context(error::UnexpectedSnafu {
1577 err_msg: "Reads the empty table info in comparing operation of the updating table info",
1578 })?
1579 .into_inner();
1580
1581 let op_name = "the updating table info";
1582 ensure_values!(remote_table_info, new_table_info_value, op_name);
1583 }
1584 Ok(())
1585 }
1586
1587 #[allow(clippy::too_many_arguments)]
1598 pub async fn update_view_info(
1599 &self,
1600 view_id: TableId,
1601 current_view_info_value: &DeserializedValueWithBytes<ViewInfoValue>,
1602 new_view_info: Vec<u8>,
1603 table_names: HashSet<TableName>,
1604 columns: Vec<String>,
1605 plan_columns: Vec<String>,
1606 definition: String,
1607 ) -> Result<()> {
1608 let new_view_info_value = current_view_info_value.update(
1609 new_view_info.into(),
1610 table_names,
1611 columns,
1612 plan_columns,
1613 definition,
1614 );
1615
1616 let (update_view_info_txn, on_update_view_info_failure) = self
1618 .view_info_manager()
1619 .build_update_txn(view_id, current_view_info_value, &new_view_info_value)?;
1620
1621 let mut r = self.kv_backend.txn(update_view_info_txn).await?;
1622
1623 if !r.succeeded {
1625 let mut set = TxnOpGetResponseSet::from(&mut r.responses);
1626 let remote_view_info = on_update_view_info_failure(&mut set)?
1627 .context(error::UnexpectedSnafu {
1628 err_msg: "Reads the empty view info in comparing operation of the updating view info",
1629 })?
1630 .into_inner();
1631
1632 let op_name = "the updating view info";
1633 ensure_values!(remote_view_info, new_view_info_value, op_name);
1634 }
1635 Ok(())
1636 }
1637
1638 pub fn batch_update_table_info_value_chunk_size(&self) -> usize {
1639 self.kv_backend.max_txn_ops()
1640 }
1641
1642 pub async fn batch_update_table_info_values(
1643 &self,
1644 table_info_value_pairs: Vec<(DeserializedValueWithBytes<TableInfoValue>, TableInfo)>,
1645 ) -> Result<()> {
1646 let len = table_info_value_pairs.len();
1647 let mut txns = Vec::with_capacity(len);
1648 struct OnFailure<F, R>
1649 where
1650 F: FnOnce(&mut TxnOpGetResponseSet) -> R,
1651 {
1652 table_info_value: TableInfoValue,
1653 on_update_table_info_failure: F,
1654 }
1655 let mut on_failures = Vec::with_capacity(len);
1656
1657 for (table_info_value, new_table_info) in table_info_value_pairs {
1658 let table_id = table_info_value.table_info.ident.table_id;
1659
1660 let new_table_info_value = table_info_value.update(new_table_info);
1661
1662 let (update_table_info_txn, on_update_table_info_failure) =
1663 self.table_info_manager().build_update_txn(
1664 table_id,
1665 &table_info_value,
1666 &new_table_info_value,
1667 )?;
1668
1669 txns.push(update_table_info_txn);
1670
1671 on_failures.push(OnFailure {
1672 table_info_value: new_table_info_value,
1673 on_update_table_info_failure,
1674 });
1675 }
1676
1677 let txn = Txn::merge_all(txns);
1678 let mut r = self.kv_backend.txn(txn).await?;
1679
1680 if !r.succeeded {
1681 let mut set = TxnOpGetResponseSet::from(&mut r.responses);
1682 for on_failure in on_failures {
1683 let remote_table_info = (on_failure.on_update_table_info_failure)(&mut set)?
1684 .context(error::UnexpectedSnafu {
1685 err_msg: "Reads the empty table info in comparing operation of the updating table info",
1686 })?
1687 .into_inner();
1688
1689 let op_name = "the batch updating table info";
1690 ensure_values!(remote_table_info, on_failure.table_info_value, op_name);
1691 }
1692 }
1693
1694 Ok(())
1695 }
1696
1697 pub async fn update_table_route(
1698 &self,
1699 table_id: TableId,
1700 region_info: RegionInfo,
1701 current_table_route_value: &DeserializedValueWithBytes<TableRouteValue>,
1702 new_region_routes: Vec<RegionRoute>,
1703 new_region_options: &HashMap<String, String>,
1704 new_region_wal_options: &RegionWalOptions,
1705 ) -> Result<()> {
1706 let current_region_distribution =
1708 region_distribution(current_table_route_value.region_routes()?);
1709 let new_region_distribution = region_distribution(&new_region_routes);
1710
1711 let update_topic_region_txn = self.topic_region_manager.build_update_txn(
1712 table_id,
1713 ®ion_info.region_wal_options,
1714 new_region_wal_options,
1715 )?;
1716 let update_datanode_table_txn = self.datanode_table_manager().build_update_txn(
1717 table_id,
1718 region_info,
1719 current_region_distribution,
1720 new_region_distribution,
1721 new_region_options,
1722 new_region_wal_options,
1723 )?;
1724
1725 let new_table_route_value = current_table_route_value.update(new_region_routes)?;
1727 let (update_table_route_txn, on_update_table_route_failure) = self
1728 .table_route_manager()
1729 .table_route_storage()
1730 .build_update_txn(table_id, current_table_route_value, &new_table_route_value)?;
1731
1732 let txn = Txn::merge_all(vec![
1733 update_datanode_table_txn,
1734 update_table_route_txn,
1735 update_topic_region_txn,
1736 ]);
1737
1738 let mut r = self.kv_backend.txn(txn).await?;
1739
1740 if !r.succeeded {
1742 let mut set = TxnOpGetResponseSet::from(&mut r.responses);
1743 let remote_table_route = on_update_table_route_failure(&mut set)?
1744 .context(error::UnexpectedSnafu {
1745 err_msg: "Reads the empty table route in comparing operation of the updating table route",
1746 })?
1747 .into_inner();
1748
1749 let op_name = "the updating table route";
1750 ensure_values!(remote_table_route, new_table_route_value, op_name);
1751 }
1752
1753 Ok(())
1754 }
1755
1756 pub async fn update_leader_region_status<F>(
1758 &self,
1759 table_id: TableId,
1760 current_table_route_value: &DeserializedValueWithBytes<TableRouteValue>,
1761 next_region_route_status: F,
1762 ) -> Result<()>
1763 where
1764 F: Fn(&RegionRoute) -> Option<Option<LeaderState>>,
1765 {
1766 let mut new_region_routes = current_table_route_value.region_routes()?.clone();
1767
1768 let mut updated = 0;
1769 for route in &mut new_region_routes {
1770 if let Some(state) = next_region_route_status(route)
1771 && route.set_leader_state(state)
1772 {
1773 updated += 1;
1774 }
1775 }
1776
1777 if updated == 0 {
1778 warn!("No leader status updated");
1779 return Ok(());
1780 }
1781
1782 let new_table_route_value = current_table_route_value.update(new_region_routes)?;
1784
1785 let (update_table_route_txn, on_update_table_route_failure) = self
1786 .table_route_manager()
1787 .table_route_storage()
1788 .build_update_txn(table_id, current_table_route_value, &new_table_route_value)?;
1789
1790 let mut r = self.kv_backend.txn(update_table_route_txn).await?;
1791
1792 if !r.succeeded {
1794 let mut set = TxnOpGetResponseSet::from(&mut r.responses);
1795 let remote_table_route = on_update_table_route_failure(&mut set)?
1796 .context(error::UnexpectedSnafu {
1797 err_msg: "Reads the empty table route in comparing operation of the updating leader region status",
1798 })?
1799 .into_inner();
1800
1801 let op_name = "the updating leader region status";
1802 ensure_values!(remote_table_route, new_table_route_value, op_name);
1803 }
1804
1805 Ok(())
1806 }
1807}
1808
1809#[macro_export]
1810macro_rules! impl_metadata_value {
1811 ($($val_ty: ty), *) => {
1812 $(
1813 impl $crate::key::MetadataValue for $val_ty {
1814 fn try_from_raw_value(raw_value: &[u8]) -> Result<Self> {
1815 serde_json::from_slice(raw_value).context(SerdeJsonSnafu)
1816 }
1817
1818 fn try_as_raw_value(&self) -> Result<Vec<u8>> {
1819 serde_json::to_vec(self).context(SerdeJsonSnafu)
1820 }
1821 }
1822 )*
1823 }
1824}
1825
1826macro_rules! impl_metadata_key_get_txn_op {
1827 ($($key: ty), *) => {
1828 $(
1829 impl $crate::key::MetadataKeyGetTxnOp for $key {
1830 fn build_get_op(
1833 &self,
1834 ) -> (
1835 TxnOp,
1836 impl for<'a> FnMut(
1837 &'a mut TxnOpGetResponseSet,
1838 ) -> Option<Vec<u8>>,
1839 ) {
1840 let raw_key = self.to_bytes();
1841 (
1842 TxnOp::Get(raw_key.clone()),
1843 TxnOpGetResponseSet::filter(raw_key),
1844 )
1845 }
1846 }
1847 )*
1848 }
1849}
1850
1851impl_metadata_key_get_txn_op! {
1852 TableNameKey<'_>,
1853 TableInfoKey,
1854 ViewInfoKey,
1855 TableRouteKey,
1856 DatanodeTableKey
1857}
1858
1859#[macro_export]
1860macro_rules! impl_optional_metadata_value {
1861 ($($val_ty: ty), *) => {
1862 $(
1863 impl $val_ty {
1864 pub fn try_from_raw_value(raw_value: &[u8]) -> Result<Option<Self>> {
1865 serde_json::from_slice(raw_value).context(SerdeJsonSnafu)
1866 }
1867
1868 pub fn try_as_raw_value(&self) -> Result<Vec<u8>> {
1869 serde_json::to_vec(self).context(SerdeJsonSnafu)
1870 }
1871 }
1872 )*
1873 }
1874}
1875
1876impl_metadata_value! {
1877 TableNameValue,
1878 TableInfoValue,
1879 ViewInfoValue,
1880 DatanodeTableValue,
1881 FlowInfoValue,
1882 FlowNameValue,
1883 FlowRouteValue,
1884 TableFlowValue,
1885 NodeAddressValue,
1886 SchemaNameValue,
1887 FlowStateValue,
1888 PoisonValue,
1889 TopicRegionValue
1890}
1891
1892impl_optional_metadata_value! {
1893 CatalogNameValue,
1894 SchemaNameValue
1895}
1896
1897#[cfg(test)]
1898mod tests {
1899 use std::collections::{BTreeMap, HashMap, HashSet};
1900 use std::sync::Arc;
1901 use std::sync::atomic::{AtomicUsize, Ordering};
1902
1903 use bytes::Bytes;
1904 use common_catalog::consts::{DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME};
1905 use common_time::util::current_time_millis;
1906 use common_wal::options::{KafkaWalOptions, WalOptions};
1907 use futures::TryStreamExt;
1908 use store_api::storage::{RegionId, RegionNumber};
1909 use table::metadata::{TableId, TableInfo};
1910 use table::table_name::TableName;
1911
1912 use super::datanode_table::DatanodeTableKey;
1913 use super::test_utils;
1914 use crate::ddl::allocator::wal_options::WalOptionsAllocator;
1915 use crate::ddl::test_util::create_table::test_create_table_task;
1916 use crate::ddl::utils::region_storage_path;
1917 use crate::error::{Error, Result};
1918 use crate::key::datanode_table::RegionInfo;
1919 use crate::key::node_address::{NodeAddressKey, NodeAddressValue};
1920 use crate::key::table_info::TableInfoValue;
1921 use crate::key::table_name::TableNameKey;
1922 use crate::key::table_route::TableRouteValue;
1923 use crate::key::topic_region::TopicRegionKey;
1924 use crate::key::{
1925 DeserializedValueWithBytes, DroppedTableMetadata, MetadataValue, RegionDistribution,
1926 RegionRoleSet, TOPIC_REGION_PREFIX, TableMetadataManager, ViewInfoValue,
1927 };
1928 use crate::kv_backend::KvBackend;
1929 use crate::kv_backend::memory::MemoryKvBackend;
1930 use crate::kv_backend::read_only::ReadOnlyKvBackend;
1931 use crate::kv_backend::test_util::MockKvBackend;
1932 use crate::peer::Peer;
1933 use crate::rpc::router::{LeaderState, Region, RegionRoute, region_distribution};
1934 use crate::rpc::store::{BatchGetResponse, PutRequest, RangeRequest, RangeResponse};
1935 use crate::wal_provider::{RegionWalOptions, WalProvider};
1936
1937 #[test]
1938 fn test_deserialized_value_with_bytes() {
1939 let region_route = new_test_region_route();
1940 let region_routes = vec![region_route.clone()];
1941
1942 let expected_region_routes =
1943 TableRouteValue::physical(vec![region_route.clone(), region_route.clone()]);
1944 let expected = serde_json::to_vec(&expected_region_routes).unwrap();
1945
1946 let value = DeserializedValueWithBytes {
1949 inner: TableRouteValue::physical(region_routes.clone()),
1951 bytes: Bytes::from(expected.clone()),
1952 };
1953
1954 let encoded = serde_json::to_vec(&value).unwrap();
1955
1956 let decoded: DeserializedValueWithBytes<TableRouteValue> =
1959 serde_json::from_slice(&encoded).unwrap();
1960
1961 assert_eq!(decoded.inner, expected_region_routes);
1962 assert_eq!(decoded.bytes, expected);
1963 }
1964
1965 fn new_test_region_route() -> RegionRoute {
1966 new_region_route(1, 2)
1967 }
1968
1969 fn new_region_route(region_id: u64, datanode: u64) -> RegionRoute {
1970 RegionRoute {
1971 region: Region {
1972 id: region_id.into(),
1973 name: "r1".to_string(),
1974 attrs: BTreeMap::new(),
1975 partition_expr: Default::default(),
1976 },
1977 leader_peer: Some(Peer::new(datanode, "a2")),
1978 follower_peers: vec![],
1979 leader_state: None,
1980 leader_down_since: None,
1981 write_route_policy: None,
1982 }
1983 }
1984
1985 fn new_test_table_info() -> TableInfo {
1986 test_utils::new_test_table_info(10)
1987 }
1988
1989 fn new_test_table_names() -> HashSet<TableName> {
1990 let mut set = HashSet::new();
1991 set.insert(TableName {
1992 catalog_name: "greptime".to_string(),
1993 schema_name: "public".to_string(),
1994 table_name: "a_table".to_string(),
1995 });
1996 set.insert(TableName {
1997 catalog_name: "greptime".to_string(),
1998 schema_name: "public".to_string(),
1999 table_name: "b_table".to_string(),
2000 });
2001 set
2002 }
2003
2004 async fn create_physical_table_metadata(
2005 table_metadata_manager: &TableMetadataManager,
2006 table_info: TableInfo,
2007 region_routes: Vec<RegionRoute>,
2008 region_wal_options: RegionWalOptions,
2009 ) -> Result<()> {
2010 table_metadata_manager
2011 .create_table_metadata(
2012 table_info,
2013 TableRouteValue::physical(region_routes),
2014 region_wal_options,
2015 )
2016 .await
2017 }
2018
2019 fn create_mock_region_wal_options() -> HashMap<RegionNumber, WalOptions> {
2020 let topics = (0..2)
2021 .map(|i| format!("greptimedb_topic{}", i))
2022 .collect::<Vec<_>>();
2023 let wal_options = topics
2024 .iter()
2025 .map(|topic| WalOptions::Kafka(KafkaWalOptions::new(topic.clone())))
2026 .collect::<Vec<_>>();
2027
2028 (0..16)
2029 .enumerate()
2030 .map(|(i, region_number)| (region_number, wal_options[i % wal_options.len()].clone()))
2031 .collect()
2032 }
2033
2034 fn create_mixed_region_wal_options() -> HashMap<RegionNumber, WalOptions> {
2035 HashMap::from([
2036 (
2037 0,
2038 WalOptions::Kafka(KafkaWalOptions::new("greptimedb_topic0".to_string())),
2039 ),
2040 (1, WalOptions::RaftEngine),
2041 (2, WalOptions::Noop),
2042 (
2043 3,
2044 WalOptions::Kafka(KafkaWalOptions::new("greptimedb_topic1".to_string())),
2045 ),
2046 ])
2047 }
2048
2049 fn test_physical_region_route(
2050 table_id: TableId,
2051 region_number: RegionNumber,
2052 leader_peer: u64,
2053 follower_peers: Vec<u64>,
2054 ) -> RegionRoute {
2055 RegionRoute {
2056 region: Region::new_test(RegionId::new(table_id, region_number)),
2057 leader_peer: Some(Peer::empty(leader_peer)),
2058 follower_peers: follower_peers.into_iter().map(Peer::empty).collect(),
2059 leader_state: None,
2060 leader_down_since: None,
2061 write_route_policy: None,
2062 }
2063 }
2064
2065 async fn create_dropped_physical_table_metadata(
2066 table_id: TableId,
2067 table_name: &str,
2068 region_routes: Vec<RegionRoute>,
2069 region_wal_options: HashMap<RegionNumber, WalOptions>,
2070 ) -> (
2071 Arc<MemoryKvBackend<Error>>,
2072 TableMetadataManager,
2073 TableName,
2074 TableInfo,
2075 Vec<RegionRoute>,
2076 HashMap<RegionNumber, WalOptions>,
2077 ) {
2078 let mem_kv = Arc::new(MemoryKvBackend::default());
2079 let table_metadata_manager = TableMetadataManager::new(mem_kv.clone());
2080 let task = test_create_table_task(table_name, table_id);
2081 let table_info = task.table_info.clone();
2082 table_metadata_manager
2083 .create_table_metadata(
2084 table_info.clone(),
2085 TableRouteValue::physical(region_routes),
2086 region_wal_options.clone(),
2087 )
2088 .await
2089 .unwrap();
2090
2091 let table_route_value = table_metadata_manager
2092 .table_route_manager
2093 .table_route_storage()
2094 .get_with_raw_bytes(table_id)
2095 .await
2096 .unwrap()
2097 .unwrap();
2098 let region_routes = table_route_value.region_routes().unwrap().clone();
2099 let table_name = TableName::new(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME, table_name);
2100 let table_route_value = TableRouteValue::physical(region_routes.clone());
2101 table_metadata_manager
2102 .delete_table_metadata(
2103 table_id,
2104 &table_name,
2105 &table_route_value,
2106 ®ion_wal_options,
2107 None,
2108 )
2109 .await
2110 .unwrap();
2111
2112 (
2113 mem_kv,
2114 table_metadata_manager,
2115 table_name,
2116 table_info,
2117 region_routes,
2118 region_wal_options,
2119 )
2120 }
2121
2122 fn assert_dropped_table_metadata(
2123 dropped_table: &DroppedTableMetadata,
2124 table_id: TableId,
2125 table_name: &TableName,
2126 table_info: &TableInfo,
2127 region_routes: &[RegionRoute],
2128 region_wal_options: &HashMap<RegionNumber, WalOptions>,
2129 ) {
2130 assert_eq!(dropped_table.table_id, table_id);
2131 assert_eq!(&dropped_table.table_name, table_name);
2132 assert_eq!(&dropped_table.table_info_value.table_info, table_info);
2133 assert_eq!(
2134 dropped_table.table_route_value.region_routes().unwrap(),
2135 region_routes
2136 );
2137 assert_eq!(&dropped_table.region_wal_options, region_wal_options);
2138 }
2139
2140 #[tokio::test]
2141 async fn test_raft_engine_topic_region_map() {
2142 let mem_kv = Arc::new(MemoryKvBackend::default());
2143 let table_metadata_manager = TableMetadataManager::new(mem_kv.clone());
2144 let region_route = new_test_region_route();
2145 let region_routes = &vec![region_route.clone()];
2146 let table_info = new_test_table_info();
2147 let wal_provider = WalProvider::RaftEngine;
2148 let regions: Vec<_> = (0..16).collect();
2149 let region_wal_options = wal_provider.allocate(®ions, false).await.unwrap();
2150 create_physical_table_metadata(
2151 &table_metadata_manager,
2152 table_info.clone(),
2153 region_routes.clone(),
2154 region_wal_options.clone(),
2155 )
2156 .await
2157 .unwrap();
2158
2159 let topic_region_key = TOPIC_REGION_PREFIX.to_string();
2160 let range_req = RangeRequest::new().with_prefix(topic_region_key);
2161 let resp = mem_kv.range(range_req).await.unwrap();
2162 assert!(resp.kvs.is_empty());
2164 }
2165
2166 #[tokio::test]
2167 async fn test_create_table_metadata() {
2168 let mem_kv = Arc::new(MemoryKvBackend::default());
2169 let table_metadata_manager = TableMetadataManager::new(mem_kv);
2170 let region_route = new_test_region_route();
2171 let region_routes = &vec![region_route.clone()];
2172 let table_info = new_test_table_info();
2173 let region_wal_options = create_mock_region_wal_options();
2174
2175 create_physical_table_metadata(
2177 &table_metadata_manager,
2178 table_info.clone(),
2179 region_routes.clone(),
2180 region_wal_options.clone(),
2181 )
2182 .await
2183 .unwrap();
2184
2185 assert!(
2187 create_physical_table_metadata(
2188 &table_metadata_manager,
2189 table_info.clone(),
2190 region_routes.clone(),
2191 region_wal_options.clone(),
2192 )
2193 .await
2194 .is_ok()
2195 );
2196
2197 let mut modified_region_routes = region_routes.clone();
2198 modified_region_routes.push(region_route.clone());
2199 assert!(
2201 create_physical_table_metadata(
2202 &table_metadata_manager,
2203 table_info.clone(),
2204 modified_region_routes,
2205 region_wal_options.clone(),
2206 )
2207 .await
2208 .is_err()
2209 );
2210
2211 let (remote_table_info, remote_table_route) = table_metadata_manager
2212 .get_full_table_info(10)
2213 .await
2214 .unwrap();
2215
2216 assert_eq!(
2217 remote_table_info.unwrap().into_inner().table_info,
2218 table_info
2219 );
2220 assert_eq!(
2221 remote_table_route
2222 .unwrap()
2223 .into_inner()
2224 .region_routes()
2225 .unwrap(),
2226 region_routes
2227 );
2228
2229 for i in 0..2 {
2230 let region_number = i as u32;
2231 let region_id = RegionId::new(table_info.ident.table_id, region_number);
2232 let topic = format!("greptimedb_topic{}", i);
2233 let regions = table_metadata_manager
2234 .topic_region_manager
2235 .regions(&topic)
2236 .await
2237 .unwrap()
2238 .into_keys()
2239 .collect::<Vec<_>>();
2240 assert_eq!(regions.len(), 8);
2241 assert!(regions.contains(®ion_id));
2242 }
2243 }
2244
2245 #[tokio::test]
2246 async fn test_get_full_table_info_remaps_route_address() {
2247 let mem_kv = Arc::new(MemoryKvBackend::default());
2248 let table_metadata_manager = TableMetadataManager::new(mem_kv.clone());
2249
2250 let mut region_route = new_test_region_route();
2251 region_route.follower_peers = vec![Peer::empty(3)];
2252 let region_routes = vec![region_route];
2253 let table_info = new_test_table_info();
2254 let table_id = table_info.ident.table_id;
2255
2256 create_physical_table_metadata(
2257 &table_metadata_manager,
2258 table_info,
2259 region_routes,
2260 HashMap::new(),
2261 )
2262 .await
2263 .unwrap();
2264
2265 mem_kv
2266 .put(PutRequest {
2267 key: NodeAddressKey::with_datanode(2).to_string().into_bytes(),
2268 value: NodeAddressValue::new(Peer::new(2, "new-a2"))
2269 .try_as_raw_value()
2270 .unwrap(),
2271 ..Default::default()
2272 })
2273 .await
2274 .unwrap();
2275 mem_kv
2276 .put(PutRequest {
2277 key: NodeAddressKey::with_datanode(3).to_string().into_bytes(),
2278 value: NodeAddressValue::new(Peer::new(3, "new-a3"))
2279 .try_as_raw_value()
2280 .unwrap(),
2281 ..Default::default()
2282 })
2283 .await
2284 .unwrap();
2285
2286 let (_, table_route) = table_metadata_manager
2287 .get_full_table_info(table_id)
2288 .await
2289 .unwrap();
2290 let table_route = table_route.unwrap().into_inner();
2291 let region_routes = table_route.region_routes().unwrap();
2292
2293 assert_eq!(
2294 region_routes[0].leader_peer.as_ref().unwrap().addr,
2295 "new-a2"
2296 );
2297 assert_eq!(region_routes[0].follower_peers[0].addr, "new-a3");
2298 }
2299
2300 #[tokio::test]
2301 async fn test_get_full_table_info_with_read_only_kv_backend() {
2302 let mem_kv = Arc::new(MemoryKvBackend::default());
2303 let writable_manager = TableMetadataManager::new(mem_kv.clone());
2304
2305 let region_routes = vec![new_test_region_route()];
2306 let table_info = new_test_table_info();
2307 let table_id = table_info.ident.table_id;
2308
2309 create_physical_table_metadata(
2310 &writable_manager,
2311 table_info.clone(),
2312 region_routes.clone(),
2313 HashMap::new(),
2314 )
2315 .await
2316 .unwrap();
2317
2318 let read_only_kv = Arc::new(ReadOnlyKvBackend::new(mem_kv));
2319 let read_only_manager = TableMetadataManager::new(read_only_kv);
2320
2321 let (remote_table_info, remote_table_route) = read_only_manager
2322 .get_full_table_info(table_id)
2323 .await
2324 .unwrap();
2325
2326 assert_eq!(
2327 remote_table_info.unwrap().into_inner().table_info,
2328 table_info
2329 );
2330 assert_eq!(
2331 remote_table_route
2332 .unwrap()
2333 .into_inner()
2334 .region_routes()
2335 .unwrap(),
2336 ®ion_routes
2337 );
2338 }
2339
2340 #[tokio::test]
2341 async fn test_create_logic_tables_metadata() {
2342 let mem_kv = Arc::new(MemoryKvBackend::default());
2343 let table_metadata_manager = TableMetadataManager::new(mem_kv);
2344 let region_route = new_test_region_route();
2345 let region_routes = vec![region_route.clone()];
2346 let table_info = new_test_table_info();
2347 let table_id = table_info.ident.table_id;
2348 let table_route_value = TableRouteValue::physical(region_routes.clone());
2349
2350 let tables_data = vec![(table_info.clone(), table_route_value.clone())];
2351 table_metadata_manager
2353 .create_logical_tables_metadata(tables_data.clone())
2354 .await
2355 .unwrap();
2356
2357 assert!(
2359 table_metadata_manager
2360 .create_logical_tables_metadata(tables_data)
2361 .await
2362 .is_ok()
2363 );
2364
2365 let mut modified_region_routes = region_routes.clone();
2366 modified_region_routes.push(new_region_route(2, 3));
2367 let modified_table_route_value = TableRouteValue::physical(modified_region_routes.clone());
2368 let modified_tables_data = vec![(table_info.clone(), modified_table_route_value)];
2369 assert!(
2371 table_metadata_manager
2372 .create_logical_tables_metadata(modified_tables_data)
2373 .await
2374 .is_err()
2375 );
2376
2377 let (remote_table_info, remote_table_route) = table_metadata_manager
2378 .get_full_table_info(table_id)
2379 .await
2380 .unwrap();
2381
2382 assert_eq!(
2383 remote_table_info.unwrap().into_inner().table_info,
2384 table_info
2385 );
2386 assert_eq!(
2387 remote_table_route
2388 .unwrap()
2389 .into_inner()
2390 .region_routes()
2391 .unwrap(),
2392 ®ion_routes
2393 );
2394 }
2395
2396 #[tokio::test]
2397 async fn test_create_many_logical_tables_metadata() {
2398 let kv_backend = Arc::new(MemoryKvBackend::default());
2399 let table_metadata_manager = TableMetadataManager::new(kv_backend);
2400
2401 let mut tables_data = vec![];
2402 for i in 0..128 {
2403 let table_id = i + 1;
2404 let regin_number = table_id * 3;
2405 let region_id = RegionId::new(table_id, regin_number);
2406 let region_route = new_region_route(region_id.as_u64(), 2);
2407 let region_routes = vec![region_route.clone()];
2408 let table_info = test_utils::new_test_table_info_with_name(
2409 table_id,
2410 &format!("my_table_{}", table_id),
2411 );
2412 let table_route_value = TableRouteValue::physical(region_routes.clone());
2413
2414 tables_data.push((table_info, table_route_value));
2415 }
2416
2417 table_metadata_manager
2419 .create_logical_tables_metadata(tables_data)
2420 .await
2421 .unwrap();
2422 }
2423
2424 #[tokio::test]
2425 async fn test_delete_table_metadata() {
2426 let mem_kv = Arc::new(MemoryKvBackend::default());
2427 let table_metadata_manager = TableMetadataManager::new(mem_kv);
2428 let region_route = new_test_region_route();
2429 let region_routes = &vec![region_route.clone()];
2430 let table_info = new_test_table_info();
2431 let table_id = table_info.ident.table_id;
2432 let datanode_id = 2;
2433 let region_wal_options = create_mock_region_wal_options();
2434 let serialized_region_wal_options = region_wal_options.clone();
2435
2436 create_physical_table_metadata(
2438 &table_metadata_manager,
2439 table_info.clone(),
2440 region_routes.clone(),
2441 serialized_region_wal_options,
2442 )
2443 .await
2444 .unwrap();
2445
2446 let table_name = TableName::new(
2447 table_info.catalog_name,
2448 table_info.schema_name,
2449 table_info.name,
2450 );
2451 let table_route_value = &TableRouteValue::physical(region_routes.clone());
2452 table_metadata_manager
2454 .delete_table_metadata(
2455 table_id,
2456 &table_name,
2457 table_route_value,
2458 ®ion_wal_options,
2459 None,
2460 )
2461 .await
2462 .unwrap();
2463 table_metadata_manager
2465 .delete_table_metadata(
2466 table_id,
2467 &table_name,
2468 table_route_value,
2469 ®ion_wal_options,
2470 None,
2471 )
2472 .await
2473 .unwrap();
2474 assert!(
2475 table_metadata_manager
2476 .table_info_manager()
2477 .get(table_id)
2478 .await
2479 .unwrap()
2480 .is_none()
2481 );
2482 assert!(
2483 table_metadata_manager
2484 .table_route_manager()
2485 .table_route_storage()
2486 .get(table_id)
2487 .await
2488 .unwrap()
2489 .is_none()
2490 );
2491 assert!(
2492 table_metadata_manager
2493 .datanode_table_manager()
2494 .tables(datanode_id)
2495 .try_collect::<Vec<_>>()
2496 .await
2497 .unwrap()
2498 .is_empty()
2499 );
2500 let table_info = table_metadata_manager
2502 .table_info_manager()
2503 .get(table_id)
2504 .await
2505 .unwrap();
2506 assert!(table_info.is_none());
2507 let table_route = table_metadata_manager
2508 .table_route_manager()
2509 .table_route_storage()
2510 .get(table_id)
2511 .await
2512 .unwrap();
2513 assert!(table_route.is_none());
2514 let regions = table_metadata_manager
2516 .topic_region_manager
2517 .regions("greptimedb_topic0")
2518 .await
2519 .unwrap();
2520 assert_eq!(regions.len(), 0);
2521 let regions = table_metadata_manager
2522 .topic_region_manager
2523 .regions("greptimedb_topic1")
2524 .await
2525 .unwrap();
2526 assert_eq!(regions.len(), 0);
2527 }
2528
2529 #[tokio::test]
2530 async fn test_rename_table() {
2531 let mem_kv = Arc::new(MemoryKvBackend::default());
2532 let table_metadata_manager = TableMetadataManager::new(mem_kv);
2533 let region_route = new_test_region_route();
2534 let region_routes = vec![region_route.clone()];
2535 let table_info = new_test_table_info();
2536 let table_id = table_info.ident.table_id;
2537 create_physical_table_metadata(
2539 &table_metadata_manager,
2540 table_info.clone(),
2541 region_routes.clone(),
2542 HashMap::new(),
2543 )
2544 .await
2545 .unwrap();
2546
2547 let new_table_name = "another_name".to_string();
2548 let table_info_value =
2549 DeserializedValueWithBytes::from_inner(TableInfoValue::new(table_info.clone()));
2550
2551 table_metadata_manager
2552 .rename_table(&table_info_value, new_table_name.clone())
2553 .await
2554 .unwrap();
2555 table_metadata_manager
2557 .rename_table(&table_info_value, new_table_name.clone())
2558 .await
2559 .unwrap();
2560 let mut modified_table_info = table_info.clone();
2561 modified_table_info.name = "hi".to_string();
2562 let modified_table_info_value =
2563 DeserializedValueWithBytes::from_inner(table_info_value.update(modified_table_info));
2564 assert!(
2567 table_metadata_manager
2568 .rename_table(&modified_table_info_value, new_table_name.clone())
2569 .await
2570 .is_err()
2571 );
2572
2573 let old_table_name = TableNameKey::new(
2574 &table_info.catalog_name,
2575 &table_info.schema_name,
2576 &table_info.name,
2577 );
2578 let new_table_name = TableNameKey::new(
2579 &table_info.catalog_name,
2580 &table_info.schema_name,
2581 &new_table_name,
2582 );
2583
2584 assert!(
2585 table_metadata_manager
2586 .table_name_manager()
2587 .get(old_table_name)
2588 .await
2589 .unwrap()
2590 .is_none()
2591 );
2592
2593 assert_eq!(
2594 table_metadata_manager
2595 .table_name_manager()
2596 .get(new_table_name)
2597 .await
2598 .unwrap()
2599 .unwrap()
2600 .table_id(),
2601 table_id
2602 );
2603 }
2604
2605 #[tokio::test]
2606 async fn test_update_table_info() {
2607 let mem_kv = Arc::new(MemoryKvBackend::default());
2608 let table_metadata_manager = TableMetadataManager::new(mem_kv);
2609 let region_route = new_test_region_route();
2610 let region_routes = vec![region_route.clone()];
2611 let table_info = new_test_table_info();
2612 let table_id = table_info.ident.table_id;
2613 create_physical_table_metadata(
2615 &table_metadata_manager,
2616 table_info.clone(),
2617 region_routes.clone(),
2618 HashMap::new(),
2619 )
2620 .await
2621 .unwrap();
2622
2623 let mut new_table_info = table_info.clone();
2624 new_table_info.name = "hi".to_string();
2625 let current_table_info_value =
2626 DeserializedValueWithBytes::from_inner(TableInfoValue::new(table_info.clone()));
2627 table_metadata_manager
2629 .update_table_info(¤t_table_info_value, None, new_table_info.clone())
2630 .await
2631 .unwrap();
2632 table_metadata_manager
2634 .update_table_info(¤t_table_info_value, None, new_table_info.clone())
2635 .await
2636 .unwrap();
2637
2638 let updated_table_info = table_metadata_manager
2640 .table_info_manager()
2641 .get(table_id)
2642 .await
2643 .unwrap()
2644 .unwrap()
2645 .into_inner();
2646 assert_eq!(updated_table_info.table_info, new_table_info);
2647
2648 let mut wrong_table_info = table_info.clone();
2649 wrong_table_info.name = "wrong".to_string();
2650 let wrong_table_info_value = DeserializedValueWithBytes::from_inner(
2651 current_table_info_value.update(wrong_table_info),
2652 );
2653 assert!(
2656 table_metadata_manager
2657 .update_table_info(&wrong_table_info_value, None, new_table_info)
2658 .await
2659 .is_err()
2660 )
2661 }
2662
2663 #[tokio::test]
2664 async fn test_update_table_leader_region_status() {
2665 let mem_kv = Arc::new(MemoryKvBackend::default());
2666 let table_metadata_manager = TableMetadataManager::new(mem_kv);
2667 let datanode = 1;
2668 let region_routes = vec![
2669 RegionRoute {
2670 region: Region {
2671 id: 1.into(),
2672 name: "r1".to_string(),
2673 attrs: BTreeMap::new(),
2674 partition_expr: Default::default(),
2675 },
2676 leader_peer: Some(Peer::new(datanode, "a2")),
2677 leader_state: Some(LeaderState::Downgrading),
2678 follower_peers: vec![],
2679 leader_down_since: Some(current_time_millis()),
2680 write_route_policy: None,
2681 },
2682 RegionRoute {
2683 region: Region {
2684 id: 2.into(),
2685 name: "r2".to_string(),
2686 attrs: BTreeMap::new(),
2687 partition_expr: Default::default(),
2688 },
2689 leader_peer: Some(Peer::new(datanode, "a1")),
2690 leader_state: None,
2691 follower_peers: vec![],
2692 leader_down_since: None,
2693 write_route_policy: None,
2694 },
2695 ];
2696 let table_info = new_test_table_info();
2697 let table_id = table_info.ident.table_id;
2698 let current_table_route_value = DeserializedValueWithBytes::from_inner(
2699 TableRouteValue::physical(region_routes.clone()),
2700 );
2701
2702 create_physical_table_metadata(
2704 &table_metadata_manager,
2705 table_info.clone(),
2706 region_routes.clone(),
2707 HashMap::new(),
2708 )
2709 .await
2710 .unwrap();
2711
2712 table_metadata_manager
2713 .update_leader_region_status(table_id, ¤t_table_route_value, |region_route| {
2714 if region_route.leader_state.is_some() {
2715 None
2716 } else {
2717 Some(Some(LeaderState::Downgrading))
2718 }
2719 })
2720 .await
2721 .unwrap();
2722
2723 let updated_route_value = table_metadata_manager
2724 .table_route_manager()
2725 .table_route_storage()
2726 .get(table_id)
2727 .await
2728 .unwrap()
2729 .unwrap();
2730
2731 assert_eq!(
2732 updated_route_value.region_routes().unwrap()[0].leader_state,
2733 Some(LeaderState::Downgrading)
2734 );
2735
2736 assert!(
2737 updated_route_value.region_routes().unwrap()[0]
2738 .leader_down_since
2739 .is_some()
2740 );
2741
2742 assert_eq!(
2743 updated_route_value.region_routes().unwrap()[1].leader_state,
2744 Some(LeaderState::Downgrading)
2745 );
2746 assert!(
2747 updated_route_value.region_routes().unwrap()[1]
2748 .leader_down_since
2749 .is_some()
2750 );
2751 }
2752
2753 async fn assert_datanode_table(
2754 table_metadata_manager: &TableMetadataManager,
2755 table_id: u32,
2756 region_routes: &[RegionRoute],
2757 ) {
2758 let region_distribution = region_distribution(region_routes);
2759 for (datanode, regions) in region_distribution {
2760 let got = table_metadata_manager
2761 .datanode_table_manager()
2762 .get(&DatanodeTableKey::new(datanode, table_id))
2763 .await
2764 .unwrap()
2765 .unwrap();
2766
2767 assert_eq!(got.regions, regions.leader_regions);
2768 assert_eq!(got.follower_regions, regions.follower_regions);
2769 }
2770 }
2771
2772 #[tokio::test]
2773 async fn test_update_table_route() {
2774 let mem_kv = Arc::new(MemoryKvBackend::default());
2775 let table_metadata_manager = TableMetadataManager::new(mem_kv);
2776 let region_route = new_test_region_route();
2777 let region_routes = vec![region_route.clone()];
2778 let table_info = new_test_table_info();
2779 let table_id = table_info.ident.table_id;
2780 let engine = table_info.meta.engine.as_str();
2781 let region_storage_path =
2782 region_storage_path(&table_info.catalog_name, &table_info.schema_name);
2783 let current_table_route_value = DeserializedValueWithBytes::from_inner(
2784 TableRouteValue::physical(region_routes.clone()),
2785 );
2786
2787 create_physical_table_metadata(
2789 &table_metadata_manager,
2790 table_info.clone(),
2791 region_routes.clone(),
2792 HashMap::new(),
2793 )
2794 .await
2795 .unwrap();
2796
2797 assert_datanode_table(&table_metadata_manager, table_id, ®ion_routes).await;
2798 let new_region_routes = vec![
2799 new_region_route(1, 1),
2800 new_region_route(2, 2),
2801 new_region_route(3, 3),
2802 ];
2803 table_metadata_manager
2805 .update_table_route(
2806 table_id,
2807 RegionInfo {
2808 engine: engine.to_string(),
2809 region_storage_path: region_storage_path.clone(),
2810 region_options: HashMap::new(),
2811 region_wal_options: HashMap::new(),
2812 },
2813 ¤t_table_route_value,
2814 new_region_routes.clone(),
2815 &HashMap::new(),
2816 &HashMap::new(),
2817 )
2818 .await
2819 .unwrap();
2820 assert_datanode_table(&table_metadata_manager, table_id, &new_region_routes).await;
2821
2822 table_metadata_manager
2824 .update_table_route(
2825 table_id,
2826 RegionInfo {
2827 engine: engine.to_string(),
2828 region_storage_path: region_storage_path.clone(),
2829 region_options: HashMap::new(),
2830 region_wal_options: HashMap::new(),
2831 },
2832 ¤t_table_route_value,
2833 new_region_routes.clone(),
2834 &HashMap::new(),
2835 &HashMap::new(),
2836 )
2837 .await
2838 .unwrap();
2839
2840 let current_table_route_value = DeserializedValueWithBytes::from_inner(
2841 current_table_route_value
2842 .inner
2843 .update(new_region_routes.clone())
2844 .unwrap(),
2845 );
2846 let new_region_routes = vec![new_region_route(2, 4), new_region_route(5, 5)];
2847 table_metadata_manager
2849 .update_table_route(
2850 table_id,
2851 RegionInfo {
2852 engine: engine.to_string(),
2853 region_storage_path: region_storage_path.clone(),
2854 region_options: HashMap::new(),
2855 region_wal_options: HashMap::new(),
2856 },
2857 ¤t_table_route_value,
2858 new_region_routes.clone(),
2859 &HashMap::new(),
2860 &HashMap::new(),
2861 )
2862 .await
2863 .unwrap();
2864 assert_datanode_table(&table_metadata_manager, table_id, &new_region_routes).await;
2865
2866 let wrong_table_route_value = DeserializedValueWithBytes::from_inner(
2869 current_table_route_value
2870 .update(vec![
2871 new_region_route(1, 1),
2872 new_region_route(2, 2),
2873 new_region_route(3, 3),
2874 new_region_route(4, 4),
2875 ])
2876 .unwrap(),
2877 );
2878 assert!(
2879 table_metadata_manager
2880 .update_table_route(
2881 table_id,
2882 RegionInfo {
2883 engine: engine.to_string(),
2884 region_storage_path: region_storage_path.clone(),
2885 region_options: HashMap::new(),
2886 region_wal_options: HashMap::new(),
2887 },
2888 &wrong_table_route_value,
2889 new_region_routes,
2890 &HashMap::new(),
2891 &HashMap::new(),
2892 )
2893 .await
2894 .is_err()
2895 );
2896 }
2897
2898 #[tokio::test]
2899 async fn test_update_table_route_with_topic_region_mapping() {
2900 let mem_kv = Arc::new(MemoryKvBackend::default());
2901 let table_metadata_manager = TableMetadataManager::new(mem_kv.clone());
2902 let region_route = new_test_region_route();
2903 let region_routes = vec![region_route.clone()];
2904 let table_info = new_test_table_info();
2905 let table_id = table_info.ident.table_id;
2906 let engine = table_info.meta.engine.as_str();
2907 let region_storage_path =
2908 region_storage_path(&table_info.catalog_name, &table_info.schema_name);
2909
2910 let old_region_wal_options: RegionWalOptions = vec![
2912 (
2913 1,
2914 WalOptions::Kafka(KafkaWalOptions::new("topic_1".to_string())),
2915 ),
2916 (
2917 2,
2918 WalOptions::Kafka(KafkaWalOptions::new("topic_2".to_string())),
2919 ),
2920 ]
2921 .into_iter()
2922 .collect();
2923
2924 create_physical_table_metadata(
2925 &table_metadata_manager,
2926 table_info.clone(),
2927 region_routes.clone(),
2928 old_region_wal_options.clone(),
2929 )
2930 .await
2931 .unwrap();
2932
2933 let current_table_route_value = DeserializedValueWithBytes::from_inner(
2934 TableRouteValue::physical(region_routes.clone()),
2935 );
2936
2937 let region_id_1 = RegionId::new(table_id, 1);
2939 let region_id_2 = RegionId::new(table_id, 2);
2940 let topic_1_key = TopicRegionKey::new(region_id_1, "topic_1");
2941 let topic_2_key = TopicRegionKey::new(region_id_2, "topic_2");
2942 assert!(
2943 table_metadata_manager
2944 .topic_region_manager
2945 .get(topic_1_key.clone())
2946 .await
2947 .unwrap()
2948 .is_some()
2949 );
2950 assert!(
2951 table_metadata_manager
2952 .topic_region_manager
2953 .get(topic_2_key.clone())
2954 .await
2955 .unwrap()
2956 .is_some()
2957 );
2958
2959 let new_region_routes = vec![
2961 new_region_route(1, 1),
2962 new_region_route(2, 2),
2963 new_region_route(3, 3), ];
2965 let new_region_wal_options: RegionWalOptions = vec![
2966 (
2967 1,
2968 WalOptions::Kafka(KafkaWalOptions::new("topic_1".to_string())), ),
2970 (
2971 2,
2972 WalOptions::Kafka(KafkaWalOptions::new("topic_2".to_string())), ),
2974 (
2975 3,
2976 WalOptions::Kafka(KafkaWalOptions::new("topic_3".to_string())), ),
2978 ]
2979 .into_iter()
2980 .collect();
2981 let current_table_route_value_updated = DeserializedValueWithBytes::from_inner(
2982 current_table_route_value
2983 .inner
2984 .update(new_region_routes.clone())
2985 .unwrap(),
2986 );
2987 table_metadata_manager
2988 .update_table_route(
2989 table_id,
2990 RegionInfo {
2991 engine: engine.to_string(),
2992 region_storage_path: region_storage_path.clone(),
2993 region_options: HashMap::new(),
2994 region_wal_options: old_region_wal_options.clone(),
2995 },
2996 ¤t_table_route_value,
2997 new_region_routes.clone(),
2998 &HashMap::new(),
2999 &new_region_wal_options,
3000 )
3001 .await
3002 .unwrap();
3003 let region_id_3 = RegionId::new(table_id, 3);
3005 let topic_3_key = TopicRegionKey::new(region_id_3, "topic_3");
3006 assert!(
3007 table_metadata_manager
3008 .topic_region_manager
3009 .get(topic_3_key)
3010 .await
3011 .unwrap()
3012 .is_some()
3013 );
3014 let newer_region_routes = vec![
3016 new_region_route(1, 1),
3017 ];
3020 let newer_region_wal_options: RegionWalOptions = vec![
3021 (
3022 1,
3023 WalOptions::Kafka(KafkaWalOptions::new("topic_1".to_string())), ),
3025 (
3026 3,
3027 WalOptions::Kafka(KafkaWalOptions::new("topic_3_new".to_string())), ),
3029 ]
3030 .into_iter()
3031 .collect();
3032 table_metadata_manager
3033 .update_table_route(
3034 table_id,
3035 RegionInfo {
3036 engine: engine.to_string(),
3037 region_storage_path: region_storage_path.clone(),
3038 region_options: HashMap::new(),
3039 region_wal_options: new_region_wal_options.clone(),
3040 },
3041 ¤t_table_route_value_updated,
3042 newer_region_routes.clone(),
3043 &HashMap::new(),
3044 &newer_region_wal_options,
3045 )
3046 .await
3047 .unwrap();
3048 let topic_2_key_new = TopicRegionKey::new(region_id_2, "topic_2");
3050 assert!(
3051 table_metadata_manager
3052 .topic_region_manager
3053 .get(topic_2_key_new)
3054 .await
3055 .unwrap()
3056 .is_none()
3057 );
3058 let topic_3_key_old = TopicRegionKey::new(region_id_3, "topic_3");
3060 assert!(
3061 table_metadata_manager
3062 .topic_region_manager
3063 .get(topic_3_key_old)
3064 .await
3065 .unwrap()
3066 .is_none()
3067 );
3068 let topic_3_key_new = TopicRegionKey::new(region_id_3, "topic_3_new");
3070 assert!(
3071 table_metadata_manager
3072 .topic_region_manager
3073 .get(topic_3_key_new)
3074 .await
3075 .unwrap()
3076 .is_some()
3077 );
3078 assert!(
3080 table_metadata_manager
3081 .topic_region_manager
3082 .get(topic_1_key)
3083 .await
3084 .unwrap()
3085 .is_some()
3086 );
3087 }
3088
3089 #[tokio::test]
3090 async fn test_destroy_table_metadata() {
3091 let mem_kv = Arc::new(MemoryKvBackend::default());
3092 let table_metadata_manager = TableMetadataManager::new(mem_kv.clone());
3093 let table_id = 1025;
3094 let table_name = "foo";
3095 let task = test_create_table_task(table_name, table_id);
3096 let options = create_mixed_region_wal_options();
3097 let serialized_options = options.clone();
3098 table_metadata_manager
3099 .create_table_metadata(
3100 task.table_info,
3101 TableRouteValue::physical(vec![
3102 RegionRoute {
3103 region: Region::new_test(RegionId::new(table_id, 1)),
3104 leader_peer: Some(Peer::empty(1)),
3105 follower_peers: vec![Peer::empty(5)],
3106 leader_state: None,
3107 leader_down_since: None,
3108 write_route_policy: None,
3109 },
3110 RegionRoute {
3111 region: Region::new_test(RegionId::new(table_id, 2)),
3112 leader_peer: Some(Peer::empty(2)),
3113 follower_peers: vec![Peer::empty(4)],
3114 leader_state: None,
3115 leader_down_since: None,
3116 write_route_policy: None,
3117 },
3118 RegionRoute {
3119 region: Region::new_test(RegionId::new(table_id, 3)),
3120 leader_peer: Some(Peer::empty(3)),
3121 follower_peers: vec![],
3122 leader_state: None,
3123 leader_down_since: None,
3124 write_route_policy: None,
3125 },
3126 ]),
3127 serialized_options,
3128 )
3129 .await
3130 .unwrap();
3131 let table_name = TableName::new(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME, table_name);
3132 let table_route_value = table_metadata_manager
3133 .table_route_manager
3134 .table_route_storage()
3135 .get_with_raw_bytes(table_id)
3136 .await
3137 .unwrap()
3138 .unwrap();
3139 table_metadata_manager
3140 .destroy_table_metadata(table_id, &table_name, &table_route_value, &options)
3141 .await
3142 .unwrap();
3143 assert!(mem_kv.is_empty());
3144 }
3145
3146 #[tokio::test]
3147 async fn test_restore_table_metadata() {
3148 let mem_kv = Arc::new(MemoryKvBackend::default());
3149 let table_metadata_manager = TableMetadataManager::new(mem_kv.clone());
3150 let table_id = 1025;
3151 let table_name = "foo";
3152 let task = test_create_table_task(table_name, table_id);
3153 let options = create_mixed_region_wal_options();
3154 let serialized_options = options.clone();
3155 table_metadata_manager
3156 .create_table_metadata(
3157 task.table_info,
3158 TableRouteValue::physical(vec![
3159 RegionRoute {
3160 region: Region::new_test(RegionId::new(table_id, 1)),
3161 leader_peer: Some(Peer::empty(1)),
3162 follower_peers: vec![Peer::empty(5)],
3163 leader_state: None,
3164 leader_down_since: None,
3165 write_route_policy: None,
3166 },
3167 RegionRoute {
3168 region: Region::new_test(RegionId::new(table_id, 2)),
3169 leader_peer: Some(Peer::empty(2)),
3170 follower_peers: vec![Peer::empty(4)],
3171 leader_state: None,
3172 leader_down_since: None,
3173 write_route_policy: None,
3174 },
3175 RegionRoute {
3176 region: Region::new_test(RegionId::new(table_id, 3)),
3177 leader_peer: Some(Peer::empty(3)),
3178 follower_peers: vec![],
3179 leader_state: None,
3180 leader_down_since: None,
3181 write_route_policy: None,
3182 },
3183 ]),
3184 serialized_options,
3185 )
3186 .await
3187 .unwrap();
3188 let expected_result = mem_kv.dump();
3189 let table_route_value = table_metadata_manager
3190 .table_route_manager
3191 .table_route_storage()
3192 .get_with_raw_bytes(table_id)
3193 .await
3194 .unwrap()
3195 .unwrap();
3196 let region_routes = table_route_value.region_routes().unwrap();
3197 let table_name = TableName::new(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME, table_name);
3198 let table_route_value = TableRouteValue::physical(region_routes.clone());
3199 table_metadata_manager
3200 .delete_table_metadata(table_id, &table_name, &table_route_value, &options, None)
3201 .await
3202 .unwrap();
3203 table_metadata_manager
3204 .restore_table_metadata(table_id, &table_name, &table_route_value, &options)
3205 .await
3206 .unwrap();
3207 let kvs = mem_kv.dump();
3208 assert_eq!(kvs, expected_result);
3209 table_metadata_manager
3211 .restore_table_metadata(table_id, &table_name, &table_route_value, &options)
3212 .await
3213 .unwrap();
3214 let kvs = mem_kv.dump();
3215 assert_eq!(kvs, expected_result);
3216 }
3217
3218 #[tokio::test]
3219 async fn test_dropped_table_metadata_enumeration_and_lookup() {
3220 let table_id = 1025;
3221 let table_name = "foo";
3222 let (_, table_metadata_manager, table_name, table_info, region_routes, options) =
3223 create_dropped_physical_table_metadata(
3224 table_id,
3225 table_name,
3226 vec![
3227 test_physical_region_route(table_id, 1, 1, vec![5]),
3228 test_physical_region_route(table_id, 2, 2, vec![4]),
3229 test_physical_region_route(table_id, 3, 3, vec![]),
3230 ],
3231 create_mixed_region_wal_options(),
3232 )
3233 .await;
3234
3235 let dropped_tables = table_metadata_manager.list_dropped_tables().await.unwrap();
3236 assert_eq!(dropped_tables.len(), 1);
3237 assert_eq!(dropped_tables[0].table_id, table_id);
3238 assert_eq!(dropped_tables[0].table_name, table_name);
3239
3240 let dropped_table = table_metadata_manager
3241 .get_dropped_table(&table_name)
3242 .await
3243 .unwrap()
3244 .unwrap();
3245 assert_dropped_table_metadata(
3246 &dropped_table,
3247 table_id,
3248 &table_name,
3249 &table_info,
3250 ®ion_routes,
3251 &options,
3252 );
3253
3254 let dropped_table_by_id = table_metadata_manager
3255 .get_dropped_table_by_id(table_id)
3256 .await
3257 .unwrap()
3258 .unwrap();
3259 assert_dropped_table_metadata(
3260 &dropped_table_by_id,
3261 table_id,
3262 &table_name,
3263 &table_info,
3264 ®ion_routes,
3265 &options,
3266 );
3267 }
3268
3269 #[tokio::test]
3270 async fn test_list_dropped_tables_batches_timestamp_lookup() {
3271 let mem_kv = Arc::new(MemoryKvBackend::default());
3272 let table_metadata_manager = TableMetadataManager::new(mem_kv.clone());
3273 for (table_id, table_name, dropped_at) in
3274 [(1025, "legacy", None), (1026, "marked", Some(1234))]
3275 {
3276 let task = test_create_table_task(table_name, table_id);
3277 let table_name = task.table_name();
3278 let table_route = TableRouteValue::physical(vec![]);
3279 table_metadata_manager
3280 .create_table_metadata(task.table_info, table_route.clone(), HashMap::new())
3281 .await
3282 .unwrap();
3283 table_metadata_manager
3284 .delete_table_metadata(
3285 table_id,
3286 &table_name,
3287 &table_route,
3288 &HashMap::new(),
3289 dropped_at,
3290 )
3291 .await
3292 .unwrap();
3293 }
3294 let all_tombstones = mem_kv
3295 .range(RangeRequest::new().with_prefix("__tombstone/"))
3296 .await
3297 .unwrap()
3298 .kvs;
3299 let range_calls = Arc::new(AtomicUsize::new(0));
3300 let batch_get_calls = Arc::new(AtomicUsize::new(0));
3301 let mock = MockKvBackend {
3302 range_fn: Some({
3303 let all_tombstones = all_tombstones.clone();
3304 let range_calls = range_calls.clone();
3305 Arc::new(move |req| {
3306 range_calls.fetch_add(1, Ordering::Relaxed);
3307 let kvs = all_tombstones
3308 .iter()
3309 .filter(|kv| {
3310 if req.range_end.is_empty() {
3311 kv.key == req.key
3312 } else {
3313 kv.key >= req.key && kv.key < req.range_end
3314 }
3315 })
3316 .cloned()
3317 .collect();
3318 Ok(RangeResponse { kvs, more: false })
3319 })
3320 }),
3321 batch_get_fn: Some({
3322 let batch_get_calls = batch_get_calls.clone();
3323 Arc::new(move |req| {
3324 batch_get_calls.fetch_add(1, Ordering::Relaxed);
3325 let kvs = all_tombstones
3326 .iter()
3327 .filter(|kv| req.keys.contains(&kv.key))
3328 .cloned()
3329 .collect();
3330 Ok(BatchGetResponse { kvs })
3331 })
3332 }),
3333 put_fn: None,
3334 batch_put_fn: None,
3335 delete_range_fn: None,
3336 batch_delete_fn: None,
3337 txn: None,
3338 max_txn_ops: None,
3339 };
3340 let table_metadata_manager = TableMetadataManager::new(Arc::new(mock));
3341
3342 let dropped_tables = table_metadata_manager.list_dropped_tables().await.unwrap();
3343
3344 assert_eq!(
3345 dropped_tables
3346 .iter()
3347 .map(|table| (table.table_id, table.dropped_at))
3348 .collect::<Vec<_>>(),
3349 vec![(1025, None), (1026, Some(1234))]
3350 );
3351 assert_eq!(range_calls.load(Ordering::Relaxed), 1);
3352 assert_eq!(batch_get_calls.load(Ordering::Relaxed), 1);
3353 }
3354
3355 #[tokio::test]
3356 async fn test_dropped_table_lookup_survives_live_name_recreation() {
3357 let dropped_table_id = 1025;
3358 let recreated_table_id = 1026;
3359 let table_name = "foo";
3360 let (
3361 _,
3362 table_metadata_manager,
3363 dropped_table_name,
3364 dropped_table_info,
3365 region_routes,
3366 options,
3367 ) = create_dropped_physical_table_metadata(
3368 dropped_table_id,
3369 table_name,
3370 vec![
3371 test_physical_region_route(dropped_table_id, 1, 1, vec![5]),
3372 test_physical_region_route(dropped_table_id, 2, 2, vec![4]),
3373 ],
3374 create_mock_region_wal_options(),
3375 )
3376 .await;
3377
3378 let recreated_task = test_create_table_task(table_name, recreated_table_id);
3379 table_metadata_manager
3380 .create_table_metadata(
3381 recreated_task.table_info,
3382 TableRouteValue::physical(vec![test_physical_region_route(
3383 recreated_table_id,
3384 1,
3385 4,
3386 vec![],
3387 )]),
3388 HashMap::new(),
3389 )
3390 .await
3391 .unwrap();
3392
3393 assert_eq!(
3394 table_metadata_manager
3395 .table_name_manager()
3396 .get(TableNameKey::from(&dropped_table_name))
3397 .await
3398 .unwrap()
3399 .unwrap()
3400 .table_id(),
3401 recreated_table_id
3402 );
3403
3404 let dropped_table = table_metadata_manager
3405 .get_dropped_table(&dropped_table_name)
3406 .await
3407 .unwrap()
3408 .unwrap();
3409 assert_dropped_table_metadata(
3410 &dropped_table,
3411 dropped_table_id,
3412 &dropped_table_name,
3413 &dropped_table_info,
3414 ®ion_routes,
3415 &options,
3416 );
3417
3418 let dropped_tables = table_metadata_manager.list_dropped_tables().await.unwrap();
3419 assert_eq!(dropped_tables.len(), 1);
3420 assert_eq!(dropped_tables[0].table_id, dropped_table_id);
3421 assert_eq!(dropped_tables[0].table_name, dropped_table_name);
3422 }
3423
3424 #[tokio::test]
3425 async fn test_dropped_table_exact_lookup_ignores_unrelated_malformed_table_name_tombstone() {
3426 let table_id = 1025;
3427 let table_name = "foo";
3428 let (mem_kv, table_metadata_manager, table_name, table_info, region_routes, options) =
3429 create_dropped_physical_table_metadata(
3430 table_id,
3431 table_name,
3432 vec![
3433 test_physical_region_route(table_id, 1, 1, vec![5]),
3434 test_physical_region_route(table_id, 2, 2, vec![4]),
3435 ],
3436 create_mixed_region_wal_options(),
3437 )
3438 .await;
3439
3440 mem_kv
3441 .put(
3442 PutRequest::new()
3443 .with_key("__tombstone/__table_name/not-a-table-name-key")
3444 .with_value("malformed"),
3445 )
3446 .await
3447 .unwrap();
3448
3449 let dropped_table = table_metadata_manager
3450 .get_dropped_table(&table_name)
3451 .await
3452 .unwrap()
3453 .unwrap();
3454 assert_dropped_table_metadata(
3455 &dropped_table,
3456 table_id,
3457 &table_name,
3458 &table_info,
3459 ®ion_routes,
3460 &options,
3461 );
3462 }
3463
3464 #[tokio::test]
3465 async fn test_dropped_table_exact_lookup_tracks_reused_name_after_purge() {
3466 let old_table_id = 1025;
3467 let new_table_id = 1026;
3468 let (_, manager, table_name, _, old_routes, old_options) =
3469 create_dropped_physical_table_metadata(
3470 old_table_id,
3471 "foo",
3472 vec![test_physical_region_route(old_table_id, 1, 1, vec![])],
3473 HashMap::new(),
3474 )
3475 .await;
3476
3477 manager
3478 .delete_table_metadata_tombstone(
3479 old_table_id,
3480 &table_name,
3481 &TableRouteValue::physical(old_routes),
3482 &old_options,
3483 )
3484 .await
3485 .unwrap();
3486 let new_task = test_create_table_task("foo", new_table_id);
3487 let new_info = new_task.table_info.clone();
3488 let new_route =
3489 TableRouteValue::physical(vec![test_physical_region_route(new_table_id, 1, 1, vec![])]);
3490 manager
3491 .create_table_metadata(new_task.table_info, new_route.clone(), HashMap::new())
3492 .await
3493 .unwrap();
3494 manager
3495 .delete_table_metadata(new_table_id, &table_name, &new_route, &HashMap::new(), None)
3496 .await
3497 .unwrap();
3498
3499 let dropped = manager
3500 .get_dropped_table(&table_name)
3501 .await
3502 .unwrap()
3503 .unwrap();
3504 assert_eq!(new_table_id, dropped.table_id);
3505 assert_eq!(new_info, dropped.table_info_value.table_info);
3506 }
3507
3508 #[tokio::test]
3509 async fn test_dropped_table_exact_lookup_missing() {
3510 let manager = TableMetadataManager::new(Arc::new(MemoryKvBackend::default()));
3511
3512 assert!(
3513 manager
3514 .get_dropped_table(&TableName::new("greptime", "public", "missing"))
3515 .await
3516 .unwrap()
3517 .is_none()
3518 );
3519 }
3520
3521 #[tokio::test]
3522 async fn test_create_update_view_info() {
3523 let mem_kv = Arc::new(MemoryKvBackend::default());
3524 let table_metadata_manager = TableMetadataManager::new(mem_kv);
3525
3526 let view_info = new_test_table_info();
3527
3528 let view_id = view_info.ident.table_id;
3529
3530 let logical_plan: Vec<u8> = vec![1, 2, 3];
3531 let columns = vec!["a".to_string()];
3532 let plan_columns = vec!["number".to_string()];
3533 let table_names = new_test_table_names();
3534 let definition = "CREATE VIEW test AS SELECT * FROM numbers";
3535
3536 table_metadata_manager
3538 .create_view_metadata(
3539 view_info.clone(),
3540 logical_plan.clone(),
3541 table_names.clone(),
3542 columns.clone(),
3543 plan_columns.clone(),
3544 definition.to_string(),
3545 )
3546 .await
3547 .unwrap();
3548
3549 {
3550 let current_view_info = table_metadata_manager
3552 .view_info_manager()
3553 .get(view_id)
3554 .await
3555 .unwrap()
3556 .unwrap()
3557 .into_inner();
3558 assert_eq!(current_view_info.view_info, logical_plan);
3559 assert_eq!(current_view_info.table_names, table_names);
3560 assert_eq!(current_view_info.definition, definition);
3561 assert_eq!(current_view_info.columns, columns);
3562 assert_eq!(current_view_info.plan_columns, plan_columns);
3563 let current_table_info = table_metadata_manager
3565 .table_info_manager()
3566 .get(view_id)
3567 .await
3568 .unwrap()
3569 .unwrap()
3570 .into_inner();
3571 assert_eq!(current_table_info.table_info, view_info);
3572 }
3573
3574 let new_logical_plan: Vec<u8> = vec![4, 5, 6];
3575 let new_table_names = {
3576 let mut set = HashSet::new();
3577 set.insert(TableName {
3578 catalog_name: "greptime".to_string(),
3579 schema_name: "public".to_string(),
3580 table_name: "b_table".to_string(),
3581 });
3582 set.insert(TableName {
3583 catalog_name: "greptime".to_string(),
3584 schema_name: "public".to_string(),
3585 table_name: "c_table".to_string(),
3586 });
3587 set
3588 };
3589 let new_columns = vec!["b".to_string()];
3590 let new_plan_columns = vec!["number2".to_string()];
3591 let new_definition = "CREATE VIEW test AS SELECT * FROM b_table join c_table";
3592
3593 let current_view_info_value = DeserializedValueWithBytes::from_inner(ViewInfoValue::new(
3594 logical_plan.clone().into(),
3595 table_names,
3596 columns,
3597 plan_columns,
3598 definition.to_string(),
3599 ));
3600 table_metadata_manager
3602 .update_view_info(
3603 view_id,
3604 ¤t_view_info_value,
3605 new_logical_plan.clone(),
3606 new_table_names.clone(),
3607 new_columns.clone(),
3608 new_plan_columns.clone(),
3609 new_definition.to_string(),
3610 )
3611 .await
3612 .unwrap();
3613 table_metadata_manager
3615 .update_view_info(
3616 view_id,
3617 ¤t_view_info_value,
3618 new_logical_plan.clone(),
3619 new_table_names.clone(),
3620 new_columns.clone(),
3621 new_plan_columns.clone(),
3622 new_definition.to_string(),
3623 )
3624 .await
3625 .unwrap();
3626
3627 let updated_view_info = table_metadata_manager
3629 .view_info_manager()
3630 .get(view_id)
3631 .await
3632 .unwrap()
3633 .unwrap()
3634 .into_inner();
3635 assert_eq!(updated_view_info.view_info, new_logical_plan);
3636 assert_eq!(updated_view_info.table_names, new_table_names);
3637 assert_eq!(updated_view_info.definition, new_definition);
3638 assert_eq!(updated_view_info.columns, new_columns);
3639 assert_eq!(updated_view_info.plan_columns, new_plan_columns);
3640
3641 let wrong_view_info = logical_plan.clone();
3642 let wrong_definition = "wrong_definition";
3643 let wrong_view_info_value =
3644 DeserializedValueWithBytes::from_inner(current_view_info_value.update(
3645 wrong_view_info.into(),
3646 new_table_names.clone(),
3647 new_columns.clone(),
3648 new_plan_columns.clone(),
3649 wrong_definition.to_string(),
3650 ));
3651 assert!(
3654 table_metadata_manager
3655 .update_view_info(
3656 view_id,
3657 &wrong_view_info_value,
3658 new_logical_plan.clone(),
3659 new_table_names.clone(),
3660 vec!["c".to_string()],
3661 vec!["number3".to_string()],
3662 wrong_definition.to_string(),
3663 )
3664 .await
3665 .is_err()
3666 );
3667
3668 let current_view_info = table_metadata_manager
3670 .view_info_manager()
3671 .get(view_id)
3672 .await
3673 .unwrap()
3674 .unwrap()
3675 .into_inner();
3676 assert_eq!(current_view_info.view_info, new_logical_plan);
3677 assert_eq!(current_view_info.table_names, new_table_names);
3678 assert_eq!(current_view_info.definition, new_definition);
3679 assert_eq!(current_view_info.columns, new_columns);
3680 assert_eq!(current_view_info.plan_columns, new_plan_columns);
3681 }
3682
3683 #[test]
3684 fn test_region_role_set_deserialize() {
3685 let s = r#"{"leader_regions": [1, 2, 3], "follower_regions": [4, 5, 6]}"#;
3686 let region_role_set: RegionRoleSet = serde_json::from_str(s).unwrap();
3687 assert_eq!(region_role_set.leader_regions, vec![1, 2, 3]);
3688 assert_eq!(region_role_set.follower_regions, vec![4, 5, 6]);
3689
3690 let s = r#"[1, 2, 3]"#;
3691 let region_role_set: RegionRoleSet = serde_json::from_str(s).unwrap();
3692 assert_eq!(region_role_set.leader_regions, vec![1, 2, 3]);
3693 assert!(region_role_set.follower_regions.is_empty());
3694 }
3695
3696 #[test]
3697 fn test_region_distribution_deserialize() {
3698 let s = r#"{"1": [1,2,3], "2": {"leader_regions": [7, 8, 9], "follower_regions": [10, 11, 12]}}"#;
3699 let region_distribution: RegionDistribution = serde_json::from_str(s).unwrap();
3700 assert_eq!(region_distribution.len(), 2);
3701 assert_eq!(region_distribution[&1].leader_regions, vec![1, 2, 3]);
3702 assert!(region_distribution[&1].follower_regions.is_empty());
3703 assert_eq!(region_distribution[&2].leader_regions, vec![7, 8, 9]);
3704 assert_eq!(region_distribution[&2].follower_regions, vec![10, 11, 12]);
3705 }
3706}