Skip to main content

metric_engine/
metadata_region.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
15use std::collections::hash_map::Entry;
16use std::collections::{BTreeMap, HashMap};
17use std::sync::{Arc, Mutex, Weak};
18use std::time::Duration;
19
20use api::v1::helper::row;
21use api::v1::value::ValueData;
22use api::v1::{ColumnDataType, ColumnSchema, Rows, SemanticType};
23use async_stream::try_stream;
24use base64::Engine;
25use base64::engine::general_purpose::STANDARD_NO_PAD;
26use common_base::readable_size::ReadableSize;
27use common_recordbatch::{RecordBatch, SendableRecordBatchStream};
28use common_telemetry::{debug, info, warn};
29use datafusion::prelude::{col, lit};
30use futures_util::TryStreamExt;
31use futures_util::stream::BoxStream;
32use mito2::engine::MitoEngine;
33use moka::future::Cache;
34use moka::policy::EvictionPolicy;
35use snafu::{OptionExt, ResultExt};
36use store_api::metadata::ColumnMetadata;
37use store_api::metric_engine_consts::{
38    METADATA_SCHEMA_KEY_COLUMN_INDEX, METADATA_SCHEMA_KEY_COLUMN_NAME,
39    METADATA_SCHEMA_TIMESTAMP_COLUMN_NAME, METADATA_SCHEMA_VALUE_COLUMN_INDEX,
40    METADATA_SCHEMA_VALUE_COLUMN_NAME,
41};
42use store_api::region_engine::RegionEngine;
43use store_api::region_request::{RegionDeleteRequest, RegionPutRequest, RegionRequest};
44use store_api::storage::{RegionId, ScanRequest};
45use tokio::sync::{
46    OwnedRwLockReadGuard, OwnedRwLockWriteGuard, RwLock, RwLockReadGuard, RwLockWriteGuard,
47};
48
49use crate::error::{
50    CacheGetSnafu, CollectRecordBatchStreamSnafu, DecodeColumnValueSnafu,
51    DeserializeColumnMetadataSnafu, LogicalRegionNotFoundSnafu, MitoReadOperationSnafu,
52    MitoWriteOperationSnafu, ParseRegionIdSnafu, Result,
53};
54use crate::utils;
55
56const REGION_PREFIX: &str = "__region_";
57const COLUMN_PREFIX: &str = "__column_";
58
59/// The other two fields key and value will be used as a k-v storage.
60/// It contains two group of key:
61/// - `__region_<LOGICAL_REGION_ID>` is used for marking table existence. It doesn't have value.
62/// - `__column_<LOGICAL_REGION_ID>_<COLUMN_NAME>` is used for marking column existence,
63///   the value is column's semantic type. To avoid the key conflict, this column key
64///   will be encoded by base64([STANDARD_NO_PAD]).
65///
66/// This is a generic handler like [MetricEngine](crate::engine::MetricEngine). It
67/// will handle all the metadata related operations across physical tables. Thus
68/// every operation should be associated to a [RegionId], which is the physical
69/// table id + region sequence. This handler will transform the region group by
70/// itself.
71pub struct MetadataRegion {
72    pub(crate) mito: MitoEngine,
73    /// The cache for contents(key-value pairs) of region metadata.
74    ///
75    /// The cache should be invalidated when any new values are put into the metadata region or any
76    /// values are deleted from the metadata region.
77    cache: Cache<RegionId, RegionMetadataCacheEntry>,
78    /// Serializes cache fills with metadata writes and invalidation per metadata region.
79    ///
80    /// Holds weak references only; strong references live in [`CacheAccessLockLease`]s
81    /// returned by [`Self::cache_access_lock`]. The last lease dropped for a region
82    /// removes its entry, so the map self-cleans on success, error, or cancellation
83    /// without coupling cleanup to region drop.
84    cache_access_locks: CacheAccessLockRegistry,
85    /// Logical lock for operations that need to be serialized. Like update & read region columns.
86    ///
87    /// Region entry will be registered on creating and opening logical region, and deregistered on
88    /// removing logical region.
89    logical_region_lock: RwLock<HashMap<RegionId, Arc<RwLock<()>>>>,
90}
91
92#[derive(Clone)]
93struct RegionMetadataCacheEntry {
94    key_values: Arc<BTreeMap<String, String>>,
95    size: usize,
96}
97
98/// Weak index of per-region cache access locks.
99///
100/// The mutex is never held across `.await`.
101type CacheAccessLockRegistry = Arc<Mutex<HashMap<RegionId, Weak<RwLock<()>>>>>;
102
103/// Lease that keeps a per-region cache access lock alive.
104///
105/// Dropping the last lease for a region removes the matching registry entry.
106struct CacheAccessLockLease {
107    registry: CacheAccessLockRegistry,
108    region_id: RegionId,
109    /// Always `Some` while the lease is alive.
110    ///
111    /// `Option` lets `Drop` release the strong `Arc` before pruning the weak entry.
112    lock: Option<Arc<RwLock<()>>>,
113}
114
115impl CacheAccessLockLease {
116    async fn read(&self) -> RwLockReadGuard<'_, ()> {
117        self.lock().read().await
118    }
119
120    async fn write(&self) -> RwLockWriteGuard<'_, ()> {
121        self.lock().write().await
122    }
123
124    fn lock(&self) -> &RwLock<()> {
125        self.lock
126            .as_deref()
127            // Safety: `lock` is initialized when the lease is created and is taken only
128            // by `Drop`. `read` and `write` borrow `self`, so `Drop` cannot run while
129            // this reference is in use.
130            .expect("cache access lock lease must hold a lock")
131    }
132}
133
134impl Drop for CacheAccessLockLease {
135    fn drop(&mut self) {
136        let Some(lock) = self.lock.take() else {
137            return;
138        };
139        let weak = Arc::downgrade(&lock);
140        // Release this lease before pruning; concurrent drops must not observe
141        // each other's strong refs.
142        drop(lock);
143
144        let mut registry = self
145            .registry
146            .lock()
147            .unwrap_or_else(|poisoned| poisoned.into_inner());
148        let Some(current) = registry.get(&self.region_id) else {
149            return;
150        };
151        // `ptr_eq` protects a newer lock for the same region; `upgrade` ensures it is dead.
152        if current.ptr_eq(&weak) && current.upgrade().is_none() {
153            registry.remove(&self.region_id);
154        }
155    }
156}
157
158/// The max size of the region metadata cache.
159const MAX_CACHE_SIZE: u64 = ReadableSize::mb(128).as_bytes();
160/// The TTL of the region metadata cache.
161const CACHE_TTL: Duration = Duration::from_secs(5 * 60);
162
163impl MetadataRegion {
164    pub fn new(mito: MitoEngine) -> Self {
165        let cache = Cache::builder()
166            .max_capacity(MAX_CACHE_SIZE)
167            // Use the LRU eviction policy to minimize frequent mito scans.
168            // Recently accessed items are retained longer in the cache.
169            .eviction_policy(EvictionPolicy::lru())
170            .time_to_live(CACHE_TTL)
171            .weigher(|_, v: &RegionMetadataCacheEntry| v.size as u32)
172            .build();
173        Self {
174            mito,
175            cache,
176            cache_access_locks: CacheAccessLockRegistry::default(),
177            logical_region_lock: RwLock::new(HashMap::new()),
178        }
179    }
180
181    /// Open a logical region.
182    ///
183    /// Returns true if the logical region is opened for the first time.
184    pub async fn open_logical_region(&self, logical_region_id: RegionId) -> bool {
185        match self
186            .logical_region_lock
187            .write()
188            .await
189            .entry(logical_region_id)
190        {
191            Entry::Occupied(_) => false,
192            Entry::Vacant(vacant_entry) => {
193                vacant_entry.insert(Arc::new(RwLock::new(())));
194                true
195            }
196        }
197    }
198
199    /// Retrieve a read lock guard of given logical region id.
200    pub async fn read_lock_logical_region(
201        &self,
202        logical_region_id: RegionId,
203    ) -> Result<OwnedRwLockReadGuard<()>> {
204        let lock = self
205            .logical_region_lock
206            .read()
207            .await
208            .get(&logical_region_id)
209            .context(LogicalRegionNotFoundSnafu {
210                region_id: logical_region_id,
211            })?
212            .clone();
213        Ok(RwLock::read_owned(lock).await)
214    }
215
216    /// Retrieve a write lock guard of given logical region id.
217    pub async fn write_lock_logical_region(
218        &self,
219        logical_region_id: RegionId,
220    ) -> Result<OwnedRwLockWriteGuard<()>> {
221        let lock = self
222            .logical_region_lock
223            .read()
224            .await
225            .get(&logical_region_id)
226            .context(LogicalRegionNotFoundSnafu {
227                region_id: logical_region_id,
228            })?
229            .clone();
230        Ok(RwLock::write_owned(lock).await)
231    }
232
233    /// Remove a registered logical region from metadata.
234    ///
235    /// This method doesn't check if the previous key exists.
236    pub async fn remove_logical_region(
237        &self,
238        physical_region_id: RegionId,
239        logical_region_id: RegionId,
240    ) -> Result<()> {
241        // concat region key
242        let region_id = utils::to_metadata_region_id(physical_region_id);
243        let region_key = Self::concat_region_key(logical_region_id);
244
245        // concat column keys
246        let logical_columns = self
247            .logical_columns(physical_region_id, logical_region_id)
248            .await?;
249        let mut column_keys = logical_columns
250            .into_iter()
251            .map(|(col, _)| Self::concat_column_key(logical_region_id, &col))
252            .collect::<Vec<_>>();
253
254        // remove region key and column keys
255        column_keys.push(region_key);
256        self.delete(region_id, &column_keys).await?;
257
258        self.logical_region_lock
259            .write()
260            .await
261            .remove(&logical_region_id);
262
263        Ok(())
264    }
265
266    // TODO(ruihang): avoid using `get_all`
267    /// Get all the columns of a given logical region.
268    /// Return a list of (column_name, column_metadata).
269    pub async fn logical_columns(
270        &self,
271        physical_region_id: RegionId,
272        logical_region_id: RegionId,
273    ) -> Result<Vec<(String, ColumnMetadata)>> {
274        let metadata_region_id = utils::to_metadata_region_id(physical_region_id);
275        let region_column_prefix = Self::concat_column_key_prefix(logical_region_id);
276
277        let mut columns = vec![];
278        for (k, v) in self
279            .get_all_with_prefix(metadata_region_id, &region_column_prefix)
280            .await?
281        {
282            if !k.starts_with(&region_column_prefix) {
283                continue;
284            }
285            // Safety: we have checked the prefix
286            let (_, column_name) = Self::parse_column_key(&k)?.unwrap();
287            let column_metadata = Self::deserialize_column_metadata(&v)?;
288            columns.push((column_name, column_metadata));
289        }
290
291        Ok(columns)
292    }
293
294    /// Return all logical regions associated with the physical region.
295    pub async fn logical_regions(&self, physical_region_id: RegionId) -> Result<Vec<RegionId>> {
296        let metadata_region_id = utils::to_metadata_region_id(physical_region_id);
297
298        let mut regions = vec![];
299        for k in self
300            .get_all_key_with_prefix(metadata_region_id, REGION_PREFIX)
301            .await?
302        {
303            if !k.starts_with(REGION_PREFIX) {
304                continue;
305            }
306            // Safety: we have checked the prefix
307            let region_id = Self::parse_region_key(&k).unwrap();
308            let region_id = region_id.parse::<u64>().unwrap().into();
309            regions.push(region_id);
310        }
311
312        Ok(regions)
313    }
314}
315
316// utils to concat and parse key/value
317impl MetadataRegion {
318    pub fn concat_region_key(region_id: RegionId) -> String {
319        format!("{REGION_PREFIX}{}", region_id.as_u64())
320    }
321
322    /// Column name will be encoded by base64([STANDARD_NO_PAD])
323    pub fn concat_column_key(region_id: RegionId, column_name: &str) -> String {
324        let encoded_column_name = STANDARD_NO_PAD.encode(column_name);
325        format!(
326            "{COLUMN_PREFIX}{}_{}",
327            region_id.as_u64(),
328            encoded_column_name
329        )
330    }
331
332    /// Concat a column key prefix without column name
333    pub fn concat_column_key_prefix(region_id: RegionId) -> String {
334        format!("{COLUMN_PREFIX}{}_", region_id.as_u64())
335    }
336
337    pub fn parse_region_key(key: &str) -> Option<&str> {
338        key.strip_prefix(REGION_PREFIX)
339    }
340
341    /// Parse column key to (logical_region_id, column_name)
342    pub fn parse_column_key(key: &str) -> Result<Option<(RegionId, String)>> {
343        if let Some(stripped) = key.strip_prefix(COLUMN_PREFIX) {
344            let mut iter = stripped.split('_');
345
346            let region_id_raw = iter.next().unwrap();
347            let region_id = region_id_raw
348                .parse::<u64>()
349                .with_context(|_| ParseRegionIdSnafu { raw: region_id_raw })?
350                .into();
351
352            let encoded_column_name = iter.next().unwrap();
353            let column_name = STANDARD_NO_PAD
354                .decode(encoded_column_name)
355                .context(DecodeColumnValueSnafu)?;
356
357            Ok(Some((region_id, String::from_utf8(column_name).unwrap())))
358        } else {
359            Ok(None)
360        }
361    }
362
363    pub fn serialize_column_metadata(column_metadata: &ColumnMetadata) -> String {
364        serde_json::to_string(column_metadata).unwrap()
365    }
366
367    pub fn deserialize_column_metadata(column_metadata: &str) -> Result<ColumnMetadata> {
368        serde_json::from_str(column_metadata).with_context(|_| DeserializeColumnMetadataSnafu {
369            raw: column_metadata,
370        })
371    }
372}
373
374/// Decode a record batch stream to a stream of items.
375pub fn decode_batch_stream<T: Send + 'static>(
376    mut record_batch_stream: SendableRecordBatchStream,
377    decode: fn(RecordBatch) -> Vec<T>,
378) -> BoxStream<'static, Result<T>> {
379    let stream = try_stream! {
380        while let Some(batch) = record_batch_stream.try_next().await.context(CollectRecordBatchStreamSnafu)? {
381            for item in decode(batch) {
382                yield item;
383            }
384        }
385    };
386    Box::pin(stream)
387}
388
389/// Decode a record batch to a list of key and value.
390fn decode_record_batch_to_key_and_value(batch: RecordBatch) -> Vec<(String, String)> {
391    let keys = batch.iter_column_as_string(0);
392    let values = batch.iter_column_as_string(1);
393    keys.zip(values)
394        .filter_map(|(k, v)| match (k, v) {
395            (Some(k), Some(v)) => Some((k, v)),
396            (Some(k), None) => Some((k, "".to_string())),
397            (None, _) => None,
398        })
399        .collect::<Vec<_>>()
400}
401
402/// Decode a record batch to a list of key.
403fn decode_record_batch_to_key(batch: RecordBatch) -> Vec<String> {
404    batch.iter_column_as_string(0).flatten().collect::<Vec<_>>()
405}
406
407// simulate to `KvBackend`
408//
409// methods in this block assume the given region id is transformed.
410impl MetadataRegion {
411    fn build_prefix_read_request(prefix: &str, key_only: bool) -> ScanRequest {
412        let filter_expr = col(METADATA_SCHEMA_KEY_COLUMN_NAME).like(lit(prefix));
413
414        let projection = if key_only {
415            vec![METADATA_SCHEMA_KEY_COLUMN_INDEX]
416        } else {
417            vec![
418                METADATA_SCHEMA_KEY_COLUMN_INDEX,
419                METADATA_SCHEMA_VALUE_COLUMN_INDEX,
420            ]
421        };
422        ScanRequest {
423            projection: Some(projection),
424            filters: vec![filter_expr],
425            ..Default::default()
426        }
427    }
428
429    fn build_read_request() -> ScanRequest {
430        let projection = vec![
431            METADATA_SCHEMA_KEY_COLUMN_INDEX,
432            METADATA_SCHEMA_VALUE_COLUMN_INDEX,
433        ];
434        ScanRequest {
435            projection: Some(projection),
436            ..Default::default()
437        }
438    }
439
440    async fn load_all(&self, metadata_region_id: RegionId) -> Result<RegionMetadataCacheEntry> {
441        let scan_req = MetadataRegion::build_read_request();
442        let record_batch_stream = self
443            .mito
444            .scan_to_stream(metadata_region_id, scan_req)
445            .await
446            .context(MitoReadOperationSnafu)?;
447
448        let kv = decode_batch_stream(record_batch_stream, decode_record_batch_to_key_and_value)
449            .try_collect::<BTreeMap<_, _>>()
450            .await?;
451        let mut size = 0;
452        for (k, v) in kv.iter() {
453            size += k.len();
454            size += v.len();
455        }
456        let kv = Arc::new(kv);
457        Ok(RegionMetadataCacheEntry {
458            key_values: kv,
459            size,
460        })
461    }
462
463    /// Acquires the cache access lock lease for `metadata_region_id`.
464    ///
465    /// The lease must be kept alive for as long as any guard taken from it.
466    fn cache_access_lock(&self, metadata_region_id: RegionId) -> CacheAccessLockLease {
467        let mut registry = self
468            .cache_access_locks
469            .lock()
470            .unwrap_or_else(|poisoned| poisoned.into_inner());
471        let lock = match registry.get(&metadata_region_id).and_then(Weak::upgrade) {
472            Some(lock) => lock,
473            None => {
474                let lock = Arc::new(RwLock::new(()));
475                registry.insert(metadata_region_id, Arc::downgrade(&lock));
476                lock
477            }
478        };
479
480        CacheAccessLockLease {
481            registry: Arc::clone(&self.cache_access_locks),
482            region_id: metadata_region_id,
483            lock: Some(lock),
484        }
485    }
486
487    async fn get_all_with_prefix(
488        &self,
489        metadata_region_id: RegionId,
490        prefix: &str,
491    ) -> Result<HashMap<String, String>> {
492        let cache_access_lock = self.cache_access_lock(metadata_region_id);
493        let _cache_guard = cache_access_lock.read().await;
494        let region_metadata = self
495            .cache
496            .try_get_with(metadata_region_id, self.load_all(metadata_region_id))
497            .await
498            .context(CacheGetSnafu)?;
499
500        let mut result = HashMap::new();
501        get_all_with_prefix(&region_metadata, prefix, |k, v| {
502            result.insert(k.to_string(), v.to_string());
503            Ok(())
504        })?;
505        Ok(result)
506    }
507
508    pub async fn get_all_key_with_prefix(
509        &self,
510        region_id: RegionId,
511        prefix: &str,
512    ) -> Result<Vec<String>> {
513        let scan_req = MetadataRegion::build_prefix_read_request(prefix, true);
514        let record_batch_stream = self
515            .mito
516            .scan_to_stream(region_id, scan_req)
517            .await
518            .context(MitoReadOperationSnafu)?;
519
520        decode_batch_stream(record_batch_stream, decode_record_batch_to_key)
521            .try_collect::<Vec<_>>()
522            .await
523    }
524
525    /// Delete the given keys. For performance consideration, this method
526    /// doesn't check if those keys exist or not.
527    async fn delete(&self, metadata_region_id: RegionId, keys: &[String]) -> Result<()> {
528        let delete_request = Self::build_delete_request(keys);
529        self.write_metadata(metadata_region_id, RegionRequest::Delete(delete_request))
530            .await
531    }
532
533    /// Writes metadata and invalidates the corresponding cache entry.
534    async fn write_metadata(
535        &self,
536        metadata_region_id: RegionId,
537        request: RegionRequest,
538    ) -> Result<()> {
539        let cache_access_lock = self.cache_access_lock(metadata_region_id);
540        let _cache_guard = cache_access_lock.write().await;
541        self.mito
542            .handle_request(metadata_region_id, request)
543            .await
544            .context(MitoWriteOperationSnafu)?;
545        self.cache.invalidate(&metadata_region_id).await;
546
547        Ok(())
548    }
549
550    pub(crate) fn build_put_request_from_iter(
551        kv: impl Iterator<Item = (String, String)>,
552    ) -> RegionPutRequest {
553        let cols = vec![
554            ColumnSchema {
555                column_name: METADATA_SCHEMA_TIMESTAMP_COLUMN_NAME.to_string(),
556                datatype: ColumnDataType::TimestampMillisecond as _,
557                semantic_type: SemanticType::Timestamp as _,
558                ..Default::default()
559            },
560            ColumnSchema {
561                column_name: METADATA_SCHEMA_KEY_COLUMN_NAME.to_string(),
562                datatype: ColumnDataType::String as _,
563                semantic_type: SemanticType::Tag as _,
564                ..Default::default()
565            },
566            ColumnSchema {
567                column_name: METADATA_SCHEMA_VALUE_COLUMN_NAME.to_string(),
568                datatype: ColumnDataType::String as _,
569                semantic_type: SemanticType::Field as _,
570                ..Default::default()
571            },
572        ];
573        let rows = Rows {
574            schema: cols,
575            rows: kv
576                .into_iter()
577                .map(|(key, value)| {
578                    row(vec![
579                        ValueData::TimestampMillisecondValue(0),
580                        ValueData::StringValue(key),
581                        ValueData::StringValue(value),
582                    ])
583                })
584                .collect(),
585        };
586
587        RegionPutRequest {
588            // Metadata must remain recoverable regardless of the user write policy.
589            skip_wal: false,
590            rows,
591            hint: None,
592            partition_expr_version: None,
593        }
594    }
595
596    fn build_delete_request(keys: &[String]) -> RegionDeleteRequest {
597        let cols = vec![
598            ColumnSchema {
599                column_name: METADATA_SCHEMA_TIMESTAMP_COLUMN_NAME.to_string(),
600                datatype: ColumnDataType::TimestampMillisecond as _,
601                semantic_type: SemanticType::Timestamp as _,
602                ..Default::default()
603            },
604            ColumnSchema {
605                column_name: METADATA_SCHEMA_KEY_COLUMN_NAME.to_string(),
606                datatype: ColumnDataType::String as _,
607                semantic_type: SemanticType::Tag as _,
608                ..Default::default()
609            },
610        ];
611        let rows = keys
612            .iter()
613            .map(|key| {
614                row(vec![
615                    ValueData::TimestampMillisecondValue(0),
616                    ValueData::StringValue(key.clone()),
617                ])
618            })
619            .collect();
620        let rows = Rows { schema: cols, rows };
621
622        RegionDeleteRequest {
623            rows,
624            hint: None,
625            partition_expr_version: None,
626        }
627    }
628
629    /// Add logical regions to the metadata region.
630    pub async fn add_logical_regions(
631        &self,
632        physical_region_id: RegionId,
633        write_region_id: bool,
634        logical_regions: impl Iterator<Item = (RegionId, HashMap<&str, &ColumnMetadata>)>,
635    ) -> Result<()> {
636        let metadata_region_id = utils::to_metadata_region_id(physical_region_id);
637        let iter = logical_regions
638            .into_iter()
639            .flat_map(|(logical_region_id, column_metadatas)| {
640                if write_region_id {
641                    Some((
642                        MetadataRegion::concat_region_key(logical_region_id),
643                        String::new(),
644                    ))
645                } else {
646                    None
647                }
648                .into_iter()
649                .chain(column_metadatas.into_iter().map(
650                    move |(name, column_metadata)| {
651                        (
652                            MetadataRegion::concat_column_key(logical_region_id, name),
653                            MetadataRegion::serialize_column_metadata(column_metadata),
654                        )
655                    },
656                ))
657            })
658            .collect::<Vec<_>>();
659
660        let put_request = MetadataRegion::build_put_request_from_iter(iter.into_iter());
661        self.write_metadata(metadata_region_id, RegionRequest::Put(put_request))
662            .await
663    }
664
665    /// Updates logical region metadata so that any entries previously referencing
666    /// `source_region_id` are modified to reference the data region of `physical_region_id`.
667    ///
668    /// This method should be called after copying files from `source_region_id`
669    /// into the target region. It scans the metadata for the target physical
670    /// region, finds logical regions with the same region number as the source,
671    /// and reinserts region and column entries updated to use the target's
672    /// region number.
673    pub async fn transform_logical_region_metadata(
674        &self,
675        physical_region_id: RegionId,
676        source_region_id: RegionId,
677    ) -> Result<()> {
678        let metadata_region_id = utils::to_metadata_region_id(physical_region_id);
679        let data_region_id = utils::to_data_region_id(physical_region_id);
680        let logical_regions = self
681            .logical_regions(data_region_id)
682            .await?
683            .into_iter()
684            .filter(|r| r.region_number() == source_region_id.region_number())
685            .collect::<Vec<_>>();
686        if logical_regions.is_empty() {
687            info!(
688                "No logical regions found from source region {}, physical region id: {}",
689                source_region_id, physical_region_id,
690            );
691            return Ok(());
692        }
693
694        let metadata = self.load_all(metadata_region_id).await?;
695        let mut output = Vec::new();
696        for logical_region_id in &logical_regions {
697            let prefix = MetadataRegion::concat_column_key_prefix(*logical_region_id);
698            get_all_with_prefix(&metadata, &prefix, |k, v| {
699                // Safety: we have checked the prefix
700                let (src_logical_region_id, column_name) = Self::parse_column_key(k)?.unwrap();
701                // Change the region number to the data region number.
702                let new_key = MetadataRegion::concat_column_key(
703                    RegionId::new(
704                        src_logical_region_id.table_id(),
705                        data_region_id.region_number(),
706                    ),
707                    &column_name,
708                );
709                output.push((new_key, v.to_string()));
710                Ok(())
711            })?;
712
713            let new_key = MetadataRegion::concat_region_key(RegionId::new(
714                logical_region_id.table_id(),
715                data_region_id.region_number(),
716            ));
717            output.push((new_key, String::new()));
718        }
719
720        if output.is_empty() {
721            warn!(
722                "No logical regions metadata found from source region {}, physical region id: {}",
723                source_region_id, physical_region_id
724            );
725            return Ok(());
726        }
727
728        debug!(
729            "Transform logical regions metadata to physical region {}, source region: {}, transformed metadata: {}",
730            data_region_id,
731            source_region_id,
732            output.len(),
733        );
734
735        let put_request = MetadataRegion::build_put_request_from_iter(output.into_iter());
736        self.write_metadata(metadata_region_id, RegionRequest::Put(put_request))
737            .await?;
738        info!(
739            "Transformed {} logical regions metadata to physical region {}, source region: {}",
740            logical_regions.len(),
741            data_region_id,
742            source_region_id
743        );
744        Ok(())
745    }
746}
747
748fn get_all_with_prefix(
749    region_metadata: &RegionMetadataCacheEntry,
750    prefix: &str,
751    mut callback: impl FnMut(&str, &str) -> Result<()>,
752) -> Result<()> {
753    let range = region_metadata.key_values.range(prefix.to_string()..);
754    for (k, v) in range {
755        if !k.starts_with(prefix) {
756            break;
757        }
758        callback(k, v)?;
759    }
760    Ok(())
761}
762
763#[cfg(test)]
764impl MetadataRegion {
765    /// Retrieves the value associated with the given key in the specified region.
766    /// Returns `Ok(None)` if the key is not found.
767    pub async fn get(&self, region_id: RegionId, key: &str) -> Result<Option<String>> {
768        use datatypes::arrow::array::{Array, AsArray};
769
770        let filter_expr = datafusion::prelude::col(METADATA_SCHEMA_KEY_COLUMN_NAME)
771            .eq(datafusion::prelude::lit(key));
772
773        let projection = Some(vec![METADATA_SCHEMA_VALUE_COLUMN_INDEX]);
774        let scan_req = ScanRequest {
775            projection,
776            filters: vec![filter_expr],
777            ..Default::default()
778        };
779        let record_batch_stream = self
780            .mito
781            .scan_to_stream(region_id, scan_req)
782            .await
783            .context(MitoReadOperationSnafu)?;
784        let scan_result = common_recordbatch::util::collect(record_batch_stream)
785            .await
786            .context(CollectRecordBatchStreamSnafu)?;
787
788        let Some(first_batch) = scan_result.first() else {
789            return Ok(None);
790        };
791
792        let column = first_batch.column(0);
793        let column = column.as_string::<i32>();
794        let val = column.is_valid(0).then(|| column.value(0).to_string());
795
796        Ok(val)
797    }
798
799    /// Check if the given column exists. Return the semantic type if exists.
800    pub async fn column_semantic_type(
801        &self,
802        physical_region_id: RegionId,
803        logical_region_id: RegionId,
804        column_name: &str,
805    ) -> Result<Option<SemanticType>> {
806        let region_id = utils::to_metadata_region_id(physical_region_id);
807        let column_key = Self::concat_column_key(logical_region_id, column_name);
808        let semantic_type = self.get(region_id, &column_key).await?;
809        semantic_type
810            .map(|s| Self::deserialize_column_metadata(&s).map(|c| c.semantic_type))
811            .transpose()
812    }
813}
814
815#[cfg(test)]
816mod test {
817    use datatypes::data_type::ConcreteDataType;
818    use datatypes::schema::ColumnSchema;
819
820    use super::*;
821    use crate::test_util::TestEnv;
822    use crate::utils::to_metadata_region_id;
823
824    #[test]
825    fn test_metadata_put_always_writes_wal() {
826        // Both metadata put paths use this constructor rather than forwarding
827        // a user insert request, so its WAL policy must always be independent.
828        for entries in [
829            vec![],
830            vec![("region", "")],
831            vec![("region", ""), ("column", "metadata")],
832        ] {
833            let request = MetadataRegion::build_put_request_from_iter(
834                entries
835                    .iter()
836                    .map(|(key, value)| (key.to_string(), value.to_string())),
837            );
838            assert!(!request.skip_wal);
839            assert!(request.hint.is_none());
840            assert!(request.partition_expr_version.is_none());
841            let expected_rows = entries
842                .into_iter()
843                .map(|(key, value)| {
844                    row(vec![
845                        ValueData::TimestampMillisecondValue(0),
846                        ValueData::StringValue(key.to_string()),
847                        ValueData::StringValue(value.to_string()),
848                    ])
849                })
850                .collect::<Vec<_>>();
851            assert_eq!(request.rows.rows, expected_rows);
852        }
853    }
854
855    #[test]
856    fn test_concat_table_key() {
857        let region_id = RegionId::new(1234, 7844);
858        let expected = "__region_5299989651108".to_string();
859        assert_eq!(MetadataRegion::concat_region_key(region_id), expected);
860    }
861
862    #[test]
863    fn test_concat_column_key() {
864        let region_id = RegionId::new(8489, 9184);
865        let column_name = "my_column";
866        let expected = "__column_36459977384928_bXlfY29sdW1u".to_string();
867        assert_eq!(
868            MetadataRegion::concat_column_key(region_id, column_name),
869            expected
870        );
871    }
872
873    #[test]
874    fn test_parse_table_key() {
875        let region_id = RegionId::new(87474, 10607);
876        let encoded = MetadataRegion::concat_column_key(region_id, "my_column");
877        assert_eq!(encoded, "__column_375697969260911_bXlfY29sdW1u");
878
879        let decoded = MetadataRegion::parse_column_key(&encoded).unwrap();
880        assert_eq!(decoded, Some((region_id, "my_column".to_string())));
881    }
882
883    #[test]
884    fn test_parse_valid_column_key() {
885        let region_id = RegionId::new(176, 910);
886        let encoded = MetadataRegion::concat_column_key(region_id, "my_column");
887        assert_eq!(encoded, "__column_755914245006_bXlfY29sdW1u");
888
889        let decoded = MetadataRegion::parse_column_key(&encoded).unwrap();
890        assert_eq!(decoded, Some((region_id, "my_column".to_string())));
891    }
892
893    #[test]
894    fn test_parse_invalid_column_key() {
895        let key = "__column_asdfasd_????";
896        let result = MetadataRegion::parse_column_key(key);
897        assert!(result.is_err());
898    }
899
900    #[test]
901    fn test_serialize_column_metadata() {
902        let semantic_type = SemanticType::Tag;
903        let column_metadata = ColumnMetadata {
904            column_schema: ColumnSchema::new("blabla", ConcreteDataType::string_datatype(), false),
905            semantic_type,
906            column_id: 5,
907        };
908        let old_fmt = "{\"column_schema\":{\"name\":\"blabla\",\"data_type\":{\"String\":null},\"is_nullable\":false,\"is_time_index\":false,\"default_constraint\":null,\"metadata\":{}},\"semantic_type\":\"Tag\",\"column_id\":5}".to_string();
909        let new_fmt = "{\"column_schema\":{\"name\":\"blabla\",\"data_type\":{\"String\":{\"size_type\":\"Utf8\"}},\"is_nullable\":false,\"is_time_index\":false,\"default_constraint\":null,\"metadata\":{}},\"semantic_type\":\"Tag\",\"column_id\":5}".to_string();
910        assert_eq!(
911            MetadataRegion::serialize_column_metadata(&column_metadata),
912            new_fmt
913        );
914        // Ensure both old and new formats can be deserialized.
915        assert_eq!(
916            MetadataRegion::deserialize_column_metadata(&old_fmt).unwrap(),
917            column_metadata
918        );
919        assert_eq!(
920            MetadataRegion::deserialize_column_metadata(&new_fmt).unwrap(),
921            column_metadata
922        );
923
924        let semantic_type = "\"Invalid Column Metadata\"";
925        assert!(MetadataRegion::deserialize_column_metadata(semantic_type).is_err());
926    }
927
928    fn test_column_metadatas() -> HashMap<String, ColumnMetadata> {
929        HashMap::from([
930            (
931                "label1".to_string(),
932                ColumnMetadata {
933                    column_schema: ColumnSchema::new(
934                        "label1".to_string(),
935                        ConcreteDataType::string_datatype(),
936                        false,
937                    ),
938                    semantic_type: SemanticType::Tag,
939                    column_id: 5,
940                },
941            ),
942            (
943                "label2".to_string(),
944                ColumnMetadata {
945                    column_schema: ColumnSchema::new(
946                        "label2".to_string(),
947                        ConcreteDataType::string_datatype(),
948                        false,
949                    ),
950                    semantic_type: SemanticType::Tag,
951                    column_id: 5,
952                },
953            ),
954        ])
955    }
956
957    async fn add_test_logical_region(
958        metadata_region: &MetadataRegion,
959        physical_region_id: RegionId,
960        logical_region_id: RegionId,
961    ) -> Result<()> {
962        let column_metadatas = test_column_metadatas();
963        let logical_regions = std::iter::once((
964            logical_region_id,
965            column_metadatas
966                .iter()
967                .map(|(name, metadata)| (name.as_str(), metadata))
968                .collect::<HashMap<_, _>>(),
969        ));
970        metadata_region
971            .add_logical_regions(physical_region_id, true, logical_regions)
972            .await
973    }
974
975    #[tokio::test]
976    async fn add_logical_regions_to_meta_region() {
977        let env = TestEnv::new().await;
978        env.init_metric_region().await;
979        let metadata_region = env.metadata_region();
980        let physical_region_id = to_metadata_region_id(env.default_physical_region_id());
981        let column_metadatas = test_column_metadatas();
982        let logical_region_id = RegionId::new(1024, 1);
983
984        let iter = vec![(
985            logical_region_id,
986            column_metadatas
987                .iter()
988                .map(|(k, v)| (k.as_str(), v))
989                .collect::<HashMap<_, _>>(),
990        )];
991        metadata_region
992            .add_logical_regions(physical_region_id, true, iter.into_iter())
993            .await
994            .unwrap();
995        // Add logical region again.
996        let iter = vec![(
997            logical_region_id,
998            column_metadatas
999                .iter()
1000                .map(|(k, v)| (k.as_str(), v))
1001                .collect::<HashMap<_, _>>(),
1002        )];
1003        metadata_region
1004            .add_logical_regions(physical_region_id, true, iter.into_iter())
1005            .await
1006            .unwrap();
1007
1008        // Check if the logical region is added.
1009        let logical_regions = metadata_region
1010            .logical_regions(physical_region_id)
1011            .await
1012            .unwrap();
1013        assert_eq!(logical_regions.len(), 2);
1014
1015        // Check if the logical region columns are added.
1016        let logical_columns = metadata_region
1017            .logical_columns(physical_region_id, logical_region_id)
1018            .await
1019            .unwrap()
1020            .into_iter()
1021            .collect::<HashMap<_, _>>();
1022        assert_eq!(logical_columns.len(), 2);
1023        assert_eq!(column_metadatas, logical_columns);
1024    }
1025
1026    #[tokio::test]
1027    async fn metadata_writes_are_synchronized_per_region() {
1028        let env = TestEnv::new().await;
1029        env.init_metric_region().await;
1030        let other_physical_region_id = RegionId::new(2, 2);
1031        env.create_physical_region(other_physical_region_id, "/test_dir2", vec![])
1032            .await;
1033        let metadata_region = Arc::new(env.metadata_region());
1034        let physical_region_id = env.default_physical_region_id();
1035        let metadata_region_id = to_metadata_region_id(physical_region_id);
1036
1037        let (snapshot_loaded_tx, snapshot_loaded_rx) = tokio::sync::oneshot::channel();
1038        let (release_snapshot_tx, release_snapshot_rx) = tokio::sync::oneshot::channel();
1039        let cache_fill = {
1040            let metadata_region = Arc::clone(&metadata_region);
1041            tokio::spawn(async move {
1042                let cache_access_lock = metadata_region.cache_access_lock(metadata_region_id);
1043                let _cache_guard = cache_access_lock.read().await;
1044                metadata_region
1045                    .cache
1046                    .try_get_with(metadata_region_id, async {
1047                        let snapshot = metadata_region.load_all(metadata_region_id).await?;
1048                        snapshot_loaded_tx.send(()).unwrap();
1049                        release_snapshot_rx.await.unwrap();
1050                        Ok::<_, crate::error::Error>(snapshot)
1051                    })
1052                    .await
1053                    .unwrap();
1054            })
1055        };
1056        snapshot_loaded_rx.await.unwrap();
1057
1058        let logical_region_id = RegionId::new(1024, 1);
1059        let metadata_write = {
1060            let metadata_region = Arc::clone(&metadata_region);
1061            tokio::spawn(async move {
1062                add_test_logical_region(&metadata_region, physical_region_id, logical_region_id)
1063                    .await
1064            })
1065        };
1066        let mut metadata_write = metadata_write;
1067        assert!(
1068            tokio::time::timeout(std::time::Duration::from_millis(50), &mut metadata_write)
1069                .await
1070                .is_err()
1071        );
1072
1073        let other_logical_region_id = RegionId::new(2048, 2);
1074        tokio::time::timeout(
1075            std::time::Duration::from_secs(5),
1076            add_test_logical_region(
1077                &metadata_region,
1078                other_physical_region_id,
1079                other_logical_region_id,
1080            ),
1081        )
1082        .await
1083        .unwrap()
1084        .unwrap();
1085
1086        release_snapshot_tx.send(()).unwrap();
1087        cache_fill.await.unwrap();
1088        metadata_write.await.unwrap().unwrap();
1089
1090        let logical_columns = metadata_region
1091            .logical_columns(physical_region_id, logical_region_id)
1092            .await
1093            .unwrap();
1094        assert_eq!(logical_columns.len(), 2);
1095        assert!(
1096            metadata_region
1097                .cache_access_locks
1098                .lock()
1099                .unwrap()
1100                .is_empty()
1101        );
1102    }
1103
1104    #[tokio::test]
1105    async fn cache_access_locks_self_clean() {
1106        let env = TestEnv::new().await;
1107        env.init_metric_region().await;
1108        let metadata_region = env.metadata_region();
1109        let physical_region_id = env.default_physical_region_id();
1110        let metadata_region_id = to_metadata_region_id(physical_region_id);
1111        let registry_len = || metadata_region.cache_access_locks.lock().unwrap().len();
1112
1113        // Metadata reads and writes leave no lock entries behind.
1114        let logical_region_id = RegionId::new(1024, 1);
1115        add_test_logical_region(&metadata_region, physical_region_id, logical_region_id)
1116            .await
1117            .unwrap();
1118        metadata_region
1119            .logical_columns(physical_region_id, logical_region_id)
1120            .await
1121            .unwrap();
1122        assert_eq!(registry_len(), 0);
1123
1124        // Concurrent leases share one lock; the entry lives until the last lease drops.
1125        let first = metadata_region.cache_access_lock(metadata_region_id);
1126        let second = metadata_region.cache_access_lock(metadata_region_id);
1127        assert!(Arc::ptr_eq(
1128            first.lock.as_ref().unwrap(),
1129            second.lock.as_ref().unwrap()
1130        ));
1131        drop(first);
1132        assert_eq!(registry_len(), 1);
1133        drop(second);
1134        assert_eq!(registry_len(), 0);
1135
1136        // A new acquisition after cleanup mints a fresh lock and entry.
1137        let third = metadata_region.cache_access_lock(metadata_region_id);
1138        assert_eq!(registry_len(), 1);
1139        drop(third);
1140        assert_eq!(registry_len(), 0);
1141    }
1142}