1use 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
59pub struct MetadataRegion {
72 pub(crate) mito: MitoEngine,
73 cache: Cache<RegionId, RegionMetadataCacheEntry>,
78 cache_access_locks: CacheAccessLockRegistry,
85 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
98type CacheAccessLockRegistry = Arc<Mutex<HashMap<RegionId, Weak<RwLock<()>>>>>;
102
103struct CacheAccessLockLease {
107 registry: CacheAccessLockRegistry,
108 region_id: RegionId,
109 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 .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 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 if current.ptr_eq(&weak) && current.upgrade().is_none() {
153 registry.remove(&self.region_id);
154 }
155 }
156}
157
158const MAX_CACHE_SIZE: u64 = ReadableSize::mb(128).as_bytes();
160const 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 .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 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 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 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 pub async fn remove_logical_region(
237 &self,
238 physical_region_id: RegionId,
239 logical_region_id: RegionId,
240 ) -> Result<()> {
241 let region_id = utils::to_metadata_region_id(physical_region_id);
243 let region_key = Self::concat_region_key(logical_region_id);
244
245 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 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 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, ®ion_column_prefix)
280 .await?
281 {
282 if !k.starts_with(®ion_column_prefix) {
283 continue;
284 }
285 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 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 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
316impl MetadataRegion {
318 pub fn concat_region_key(region_id: RegionId) -> String {
319 format!("{REGION_PREFIX}{}", region_id.as_u64())
320 }
321
322 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 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 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
374pub 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
389fn 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
402fn decode_record_batch_to_key(batch: RecordBatch) -> Vec<String> {
404 batch.iter_column_as_string(0).flatten().collect::<Vec<_>>()
405}
406
407impl 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 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(®ion_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 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 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 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 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 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 let (src_logical_region_id, column_name) = Self::parse_column_key(k)?.unwrap();
701 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 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 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 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 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 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 let logical_regions = metadata_region
1010 .logical_regions(physical_region_id)
1011 .await
1012 .unwrap();
1013 assert_eq!(logical_regions.len(), 2);
1014
1015 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 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 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 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}