Skip to main content

common_meta/
key.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! This mod defines all the keys used in the metadata store (Metasrv).
16//! Specifically, there are these kinds of keys:
17//!
18//! 1. Datanode table key: `__dn_table/{datanode_id}/{table_id}`
19//!     - The value is a [DatanodeTableValue] struct; it contains `table_id` and the regions that
20//!       belong to this Datanode.
21//!     - This key is primary used in the startup of Datanode, to let Datanode know which tables
22//!       and regions it should open.
23//!
24//! 2. Table info key: `__table_info/{table_id}`
25//!     - The value is a [TableInfoValue] struct; it contains the whole table info (like column
26//!       schemas).
27//!     - This key is mainly used in constructing the table in Datanode and Frontend.
28//!
29//! 3. Catalog name key: `__catalog_name/{catalog_name}`
30//!     - Indices all catalog names
31//!
32//! 4. Schema name key: `__schema_name/{catalog_name}/{schema_name}`
33//!     - Indices all schema names belong to the {catalog_name}
34//!
35//! 5. Table name key: `__table_name/{catalog_name}/{schema_name}/{table_name}`
36//!     - The value is a [TableNameValue] struct; it contains the table id.
37//!     - Used in the table name to table id lookup.
38//!
39//! 6. Flow info key: `__flow/info/{flow_id}`
40//!     - Stores metadata of the flow.
41//!
42//! 7. Flow route key: `__flow/route/{flow_id}/{partition_id}`
43//!     - Stores route of the flow.
44//!
45//! 8. Flow name key: `__flow/name/{catalog}/{flow_name}`
46//!     - Mapping {catalog}/{flow_name} to {flow_id}
47//!
48//! 9. Flownode flow key: `__flow/flownode/{flownode_id}/{flow_id}/{partition_id}`
49//!     - Mapping {flownode_id} to {flow_id}
50//!
51//! 10. Table flow key: `__flow/source_table/{table_id}/{flownode_id}/{flow_id}/{partition_id}`
52//!     - Mapping source table's {table_id} to {flownode_id}
53//!     - Used in `Flownode` booting.
54//!
55//! 11. View info key: `__view_info/{view_id}`
56//!     - The value is a [ViewInfoValue] struct; it contains the encoded logical plan.
57//!     - This key is mainly used in constructing the view in Datanode and Frontend.
58//!
59//! 12. Kafka topic key: `__topic_name/kafka/{topic_name}`
60//!     - The key is used to track existing topics in Kafka.
61//!     - The value is a [TopicNameValue](crate::key::topic_name::TopicNameValue) struct; it contains the `pruned_entry_id` which represents
62//!       the highest entry id that has been pruned from the remote WAL.
63//!     - When a region uses this topic, it should start replaying entries from `pruned_entry_id + 1` (minimum available entry id).
64//!
65//! 13. Topic name to region map key `__topic_region/{topic_name}/{region_id}`
66//!     - Mapping {topic_name} to {region_id}
67//!
68//! All keys have related managers. The managers take care of the serialization and deserialization
69//! of keys and values, and the interaction with the underlying KV store backend.
70//!
71//! To simplify the managers used in struct fields and function parameters, we define "unify"
72//! table metadata manager: [TableMetadataManager]
73//! and flow metadata manager: [FlowMetadataManager](crate::key::flow::FlowMetadataManager).
74//! It contains all the managers defined above. It's recommended to just use this manager only.
75//!
76//! The whole picture of flow keys will be like this:
77//!
78//! __flow/
79//!   info/
80//!     {flow_id}
81//!   route/
82//!     {flow_id}/
83//!      {partition_id}
84//!
85//!    name/
86//!      {catalog_name}
87//!        {flow_name}
88//!
89//!    flownode/
90//!      {flownode_id}/
91//!        {flow_id}/
92//!          {partition_id}
93//!
94//!    source_table/
95//!      {table_id}/
96//!        {flownode_id}/
97//!          {flow_id}/
98//!            {partition_id}
99
100pub 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";
191// The legacy topic key prefix is used to store the topic name in previous versions.
192pub const LEGACY_TOPIC_KEY_PREFIX: &str = "__created_wal_topics/kafka";
193pub const TOPIC_REGION_PREFIX: &str = "__topic_region";
194
195/// The election key.
196pub const ELECTION_KEY: &str = "__metasrv_election";
197/// The root key of metasrv election candidates.
198pub const CANDIDATES_ROOT: &str = "__metasrv_election_candidates/";
199
200/// The keys with these prefixes will be loaded into the cache when the leader starts.
201pub 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/// A set of regions with the same role.
210#[derive(Debug, Clone, PartialEq, Eq, Default, Serialize)]
211pub struct RegionRoleSet {
212    /// Leader regions.
213    pub leader_regions: Vec<RegionNumber>,
214    /// Follower regions.
215    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    /// Create a new region role set.
246    pub fn new(leader_regions: Vec<RegionNumber>, follower_regions: Vec<RegionNumber>) -> Self {
247        Self {
248            leader_regions,
249            follower_regions,
250        }
251    }
252
253    /// Add a leader region to the set.
254    pub fn add_leader_region(&mut self, region_number: RegionNumber) {
255        self.leader_regions.push(region_number);
256    }
257
258    /// Add a follower region to the set.
259    pub fn add_follower_region(&mut self, region_number: RegionNumber) {
260        self.follower_regions.push(region_number);
261    }
262
263    /// Sort the regions.
264    pub fn sort(&mut self) {
265        self.follower_regions.sort();
266        self.leader_regions.sort();
267    }
268}
269
270/// The distribution of regions.
271///
272/// The key is the datanode id, the value is the region role set.
273pub type RegionDistribution = BTreeMap<DatanodeId, RegionRoleSet>;
274
275/// The id of flow.
276pub type FlowId = u32;
277/// The partition of flow.
278pub 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    /// CATALOG_NAME_KEY: {CATALOG_NAME_KEY_PREFIX}/{catalog_name}
318    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    /// SCHEMA_NAME_KEY: {SCHEMA_NAME_KEY_PREFIX}/{catalog_name}/{schema_name}
326    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
349/// The key of metadata.
350pub 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
424/// A struct containing a deserialized value(`inner`) and an original bytes.
425///
426/// - Serialize behaviors:
427///
428/// The `inner` field will be ignored.
429///
430/// - Deserialize behaviors:
431///
432/// The `inner` field will be deserialized from the `bytes` field.
433pub struct DeserializedValueWithBytes<T: DeserializeOwned + Serialize> {
434    // The original bytes of the inner.
435    bytes: Bytes,
436    // The value was deserialized from the original bytes.
437    inner: T,
438}
439
440#[derive(Debug, Clone, PartialEq, Eq)]
441pub struct DroppedTableName {
442    /// Table id stored in the tombstoned table-name mapping.
443    pub table_id: TableId,
444    /// Original fully qualified table name.
445    pub table_name: TableName,
446    /// Unix timestamp in milliseconds when this table was soft-dropped.
447    pub dropped_at: Option<i64>,
448    /// Fixed automatic-GC deadline in Unix milliseconds.
449    pub retention_expires_at: Option<i64>,
450    /// Unique identity of this soft-drop generation.
451    pub drop_generation: Option<String>,
452    /// Whether an automatic purge has durably claimed this dropped table.
453    pub purging: bool,
454}
455
456#[derive(Debug, Clone)]
457pub struct DroppedTableMetadata {
458    /// Table id of the dropped table.
459    pub table_id: TableId,
460    /// Original fully qualified table name.
461    pub table_name: TableName,
462    /// Tombstoned table info value.
463    pub table_info_value: TableInfoValue,
464    /// Tombstoned table route value.
465    pub table_route_value: TableRouteValue,
466    /// Per-region WAL options recovered from tombstoned datanode metadata.
467    pub region_wal_options: HashMap<RegionNumber, WalOptions>,
468    /// Unix timestamp in milliseconds when this table was soft-dropped.
469    pub dropped_at: Option<i64>,
470    /// Fixed automatic-GC deadline in Unix milliseconds.
471    pub retention_expires_at: Option<i64>,
472    /// Unique identity of this soft-drop generation.
473    pub drop_generation: Option<String>,
474}
475
476/// Lifecycle markers stored with a logically deleted table.
477pub 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    /// - Serialize behaviors:
530    ///
531    /// The `inner` field will be ignored.
532    fn serialize<S>(&self, serializer: S) -> std::result::Result<S::Ok, S::Error>
533    where
534        S: serde::Serializer,
535    {
536        // Safety: The original bytes are always JSON encoded.
537        // It's more efficiently than `serialize_bytes`.
538        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    /// - Deserialize behaviors:
546    ///
547    /// The `inner` field will be deserialized from the `bytes` field.
548    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    /// Returns a struct containing a deserialized value and an original `bytes`.
573    /// It accepts original bytes of inner.
574    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    /// Returns a struct containing a deserialized value and an original `bytes`.
580    /// It accepts original bytes of inner.
581    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    /// Returns original `bytes`
594    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    /// Creates a new `TableMetadataManager` with a custom tombstone prefix.
628    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    /// Creates metadata for view and returns an error if different metadata exists.
742    /// The caller MUST ensure it has the exclusive access to `TableNameKey`.
743    /// Parameters include:
744    /// - `view_info`: the encoded logical plan
745    /// - `table_names`: the resolved fully table names in logical plan
746    /// - `columns`: the view columns
747    /// - `plan_columns`: the original plan columns
748    /// - `definition`: The SQL to create the view
749    ///
750    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        // Creates view name.
762        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        // Creates table info.
772        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        // Creates view info
779        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        // Checks whether metadata was already created.
799        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    /// Creates metadata for table and returns an error if different metadata exists.
822    /// The caller MUST ensure it has the exclusive access to `TableNameKey`.
823    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        // Creates table name.
833        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        // Creates table info.
844        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, &region_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                &region_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        // Checks whether metadata was already created.
881        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        // The batch size is max_txn_size / 3 because the size of the `tables_data`
905        // is 3 times the size of the `tables_data`.
906        self.kv_backend.max_txn_ops() / 3
907    }
908
909    /// Creates metadata for multiple logical tables and return an error if different metadata exists.
910    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            // Creates table name.
931            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            // Creates table info.
942            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        // Checks whether metadata was already created.
966        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        // Builds keys
998        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    /// Deletes metadata for table **logically**.
1039    /// The caller MUST ensure it has the exclusive access to `TableNameKey`.
1040    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    /// Deletes metadata logically and stores fixed soft-drop lifecycle markers.
1060    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    /// Deletes metadata logically and stores fixed soft-drop lifecycle markers.
1084    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    /// Lists dropped tables from tombstoned table-name entries.
1121    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    /// Lists dropped tables from tombstoned table-name entries in the provided catalog.
1127    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    /// Gets dropped table metadata by its original full table name.
1192    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    /// Gets dropped table metadata by table id.
1211    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    /// Returns whether an automatic purge has durably claimed the dropped table.
1219    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    /// Returns the soft-drop generation claimed by an automatic purge.
1226    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    /// Durably claims a dropped table before automatic purge starts cleaning its regions.
1234    /// The caller MUST hold the table lock and revalidate the tombstone before calling this.
1235    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    /// Deletes metadata tombstone for table **permanently**.
1251    /// The caller MUST ensure it has the exclusive access to `TableNameKey`.
1252    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    /// Restores metadata for table.
1276    /// The caller MUST ensure it has the exclusive access to `TableNameKey`.
1277    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    /// Deletes metadata for table **permanently**.
1300    /// The caller MUST ensure it has the exclusive access to `TableNameKey`.
1301    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    /// Rebuilds dropped table metadata from tombstoned keys.
1318    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    /// Rebuilds region WAL options from tombstoned datanode-table entries.
1417    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    /// Deletes metadata for view **permanently**.
1471    /// The caller MUST ensure it has the exclusive access to `ViewNameKey`.
1472    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    /// Renames the table name and returns an error if different metadata exists.
1482    /// The caller MUST ensure it has the exclusive access to old and new `TableNameKey`s,
1483    /// and the new `TableNameKey` MUST be empty.
1484    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 = &current_table_info_value.table_info;
1490        let table_id = current_table_info.ident.table_id;
1491
1492        let table_name_key = TableNameKey::new(
1493            &current_table_info.catalog_name,
1494            &current_table_info.schema_name,
1495            &current_table_info.name,
1496        );
1497
1498        let new_table_name_key = TableNameKey::new(
1499            &current_table_info.catalog_name,
1500            &current_table_info.schema_name,
1501            &new_table_name,
1502        );
1503
1504        // Updates table name.
1505        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        // Updates table info.
1518        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        // Checks whether metadata was already updated.
1527        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    /// Updates table info and returns an error if different metadata exists.
1543    /// And cascade-ly update all redundant table options for each region
1544    /// if region_distribution is present.
1545    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        // Updates table info.
1555        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            // region options induced from table info.
1561            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        // Checks whether metadata was already updated.
1573        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    /// Updates view info and returns an error if different metadata exists.
1588    /// Parameters include:
1589    /// - `view_id`: the view id
1590    /// - `current_view_info_value`: the current view info for CAS checking
1591    /// - `new_view_info`: the encoded logical plan
1592    /// - `table_names`: the resolved fully table names in logical plan
1593    /// - `columns`: the view columns
1594    /// - `plan_columns`: the original plan columns
1595    /// - `definition`: The SQL to create the view
1596    ///
1597    #[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        // Updates view info.
1617        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        // Checks whether metadata was already updated.
1624        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        // Updates the datanode table key value pairs.
1707        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            &region_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        // Updates the table_route.
1726        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        // Checks whether metadata was already updated.
1741        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    /// Updates the leader status of the [RegionRoute].
1757    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        // Updates the table_route.
1783        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        // Checks whether metadata was already updated.
1793        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                /// Returns a [TxnOp] to retrieve the corresponding value
1831                /// and a filter to retrieve the value from the [TxnOpGetResponseSet]
1832                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        // Serialize behaviors:
1947        // The inner field will be ignored.
1948        let value = DeserializedValueWithBytes {
1949            // ignored
1950            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        // Deserialize behaviors:
1957        // The inner field will be deserialized from the bytes field.
1958        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                &region_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(&regions, 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        // Should be empty because the topic region map is empty for raft engine.
2163        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        // creates metadata.
2176        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        // if metadata was already created, it should be ok.
2186        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        // if remote metadata was exists, it should return an error.
2200        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(&region_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            &region_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        // creates metadata.
2352        table_metadata_manager
2353            .create_logical_tables_metadata(tables_data.clone())
2354            .await
2355            .unwrap();
2356
2357        // if metadata was already created, it should be ok.
2358        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        // if remote metadata was exists, it should return an error.
2370        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            &region_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        // creates metadata.
2418        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        // creates metadata.
2437        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        // deletes metadata.
2453        table_metadata_manager
2454            .delete_table_metadata(
2455                table_id,
2456                &table_name,
2457                table_route_value,
2458                &region_wal_options,
2459                None,
2460            )
2461            .await
2462            .unwrap();
2463        // Should be ignored.
2464        table_metadata_manager
2465            .delete_table_metadata(
2466                table_id,
2467                &table_name,
2468                table_route_value,
2469                &region_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        // Checks removed values
2501        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        // Logical delete removes the topic region mapping as well.
2515        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        // creates metadata.
2538        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        // if remote metadata was updated, it should be ok.
2556        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        // if the table_info_value is wrong, it should return an error.
2565        // The ABA problem.
2566        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        // creates metadata.
2614        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        // should be ok.
2628        table_metadata_manager
2629            .update_table_info(&current_table_info_value, None, new_table_info.clone())
2630            .await
2631            .unwrap();
2632        // if table info was updated, it should be ok.
2633        table_metadata_manager
2634            .update_table_info(&current_table_info_value, None, new_table_info.clone())
2635            .await
2636            .unwrap();
2637
2638        // updated table_info should equal the `new_table_info`
2639        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        // if the current_table_info_value is wrong, it should return an error.
2654        // The ABA problem.
2655        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        // creates metadata.
2703        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, &current_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        // creates metadata.
2788        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, &region_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        // it should be ok.
2804        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                &current_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        // if the table route was updated. it should be ok.
2823        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                &current_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        // it should be ok.
2848        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                &current_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        // if the current_table_route_value is wrong, it should return an error.
2867        // The ABA problem.
2868        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        // Create initial metadata with Kafka WAL options
2911        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        // Verify initial topic region mappings exist
2938        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        // Test 1: Add new region with new topic
2960        let new_region_routes = vec![
2961            new_region_route(1, 1),
2962            new_region_route(2, 2),
2963            new_region_route(3, 3), // New region
2964        ];
2965        let new_region_wal_options: RegionWalOptions = vec![
2966            (
2967                1,
2968                WalOptions::Kafka(KafkaWalOptions::new("topic_1".to_string())), // Unchanged
2969            ),
2970            (
2971                2,
2972                WalOptions::Kafka(KafkaWalOptions::new("topic_2".to_string())), // Unchanged
2973            ),
2974            (
2975                3,
2976                WalOptions::Kafka(KafkaWalOptions::new("topic_3".to_string())), // New topic
2977            ),
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                &current_table_route_value,
2997                new_region_routes.clone(),
2998                &HashMap::new(),
2999                &new_region_wal_options,
3000            )
3001            .await
3002            .unwrap();
3003        // Verify new topic region mapping was created
3004        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        // Test 2: Remove a region and change topic for another
3015        let newer_region_routes = vec![
3016            new_region_route(1, 1),
3017            // Region 2 removed
3018            // Region 3 now has different topic
3019        ];
3020        let newer_region_wal_options: RegionWalOptions = vec![
3021            (
3022                1,
3023                WalOptions::Kafka(KafkaWalOptions::new("topic_1".to_string())), // Unchanged
3024            ),
3025            (
3026                3,
3027                WalOptions::Kafka(KafkaWalOptions::new("topic_3_new".to_string())), // Changed topic
3028            ),
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                &current_table_route_value_updated,
3042                newer_region_routes.clone(),
3043                &HashMap::new(),
3044                &newer_region_wal_options,
3045            )
3046            .await
3047            .unwrap();
3048        // Verify region 2 mapping was deleted
3049        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        // Verify region 3 old topic mapping was deleted
3059        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        // Verify region 3 new topic mapping was created
3069        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        // Verify region 1 mapping still exists (unchanged)
3079        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        // Should be ignored.
3210        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            &region_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            &region_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            &region_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            &region_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        // Create metadata
3537        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            // assert view info
3551            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            // assert table info
3564            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        // should be ok.
3601        table_metadata_manager
3602            .update_view_info(
3603                view_id,
3604                &current_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        // if table info was updated, it should be ok.
3614        table_metadata_manager
3615            .update_view_info(
3616                view_id,
3617                &current_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        // updated view_info should equal the `new_logical_plan`
3628        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        // if the current_view_info_value is wrong, it should return an error.
3652        // The ABA problem.
3653        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        // The view_info is not changed.
3669        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}