1mod extract_new_columns;
16
17use std::collections::{HashMap, HashSet};
18
19use api::v1::SemanticType;
20use common_query::native_histogram::is_native_histogram_value_type;
21use common_telemetry::info;
22use common_time::{FOREVER, Timestamp};
23use datatypes::data_type::ConcreteDataType;
24use datatypes::schema::{ColumnSchema, SkippingIndexOptions};
25use datatypes::value::Value;
26use mito2::engine::MITO_ENGINE_NAME;
27use snafu::{OptionExt, ResultExt, ensure};
28use store_api::metadata::ColumnMetadata;
29use store_api::metric_engine_consts::{
30 ALTER_PHYSICAL_EXTENSION_KEY, DATA_REGION_SUBDIR, DATA_SCHEMA_TABLE_ID_COLUMN_NAME,
31 DATA_SCHEMA_TSID_COLUMN_NAME, LOGICAL_TABLE_METADATA_KEY, METADATA_REGION_SUBDIR,
32 METADATA_SCHEMA_KEY_COLUMN_INDEX, METADATA_SCHEMA_KEY_COLUMN_NAME,
33 METADATA_SCHEMA_TIMESTAMP_COLUMN_INDEX, METADATA_SCHEMA_TIMESTAMP_COLUMN_NAME,
34 METADATA_SCHEMA_VALUE_COLUMN_INDEX, METADATA_SCHEMA_VALUE_COLUMN_NAME,
35 is_metric_engine_internal_column,
36};
37use store_api::mito_engine_options::{TTL_KEY, WAL_OPTIONS_KEY};
38use store_api::region_engine::RegionEngine;
39use store_api::region_request::{AffectedRows, PathType, RegionCreateRequest, RegionRequest};
40use store_api::storage::RegionId;
41use store_api::storage::consts::ReservedColumnId;
42
43use crate::engine::MetricEngineInner;
44use crate::engine::create::extract_new_columns::extract_new_columns;
45use crate::engine::options::{PhysicalRegionOptions, set_data_region_options};
46use crate::error::{
47 ColumnTypeMismatchSnafu, ConflictRegionOptionSnafu, CreateMitoRegionSnafu,
48 InternalColumnOccupiedSnafu, InvalidMetadataSnafu, MissingRegionOptionSnafu,
49 MultipleFieldColumnSnafu, NoFieldColumnSnafu, ParseRegionIdSnafu, PhysicalRegionNotFoundSnafu,
50 Result, SerializeColumnMetadataSnafu, UnexpectedRequestSnafu,
51};
52use crate::metrics::PHYSICAL_REGION_COUNT;
53use crate::utils::{
54 self, append_manifest_info, encode_manifest_info_to_extensions, to_data_region_id,
55 to_metadata_region_id,
56};
57
58const DEFAULT_TABLE_ID_SKIPPING_INDEX_GRANULARITY: u32 = 1024;
59const DEFAULT_TABLE_ID_SKIPPING_INDEX_FALSE_POSITIVE_RATE: f64 = 0.01;
60
61impl MetricEngineInner {
62 pub async fn create_regions(
63 &self,
64 mut requests: Vec<(RegionId, RegionCreateRequest)>,
65 extension_return_value: &mut HashMap<String, Vec<u8>>,
66 ) -> Result<AffectedRows> {
67 if requests.is_empty() {
68 return Ok(0);
69 }
70
71 for (_, request) in requests.iter() {
72 Self::verify_region_create_request(request)?;
73 }
74
75 let first_request = &requests.first().unwrap().1;
76 if first_request.is_physical_table() {
77 ensure!(
78 requests.len() == 1,
79 UnexpectedRequestSnafu {
80 reason: "Physical table must be created with single request".to_string(),
81 }
82 );
83 let (region_id, request) = requests.pop().unwrap();
84 self.create_physical_region(region_id, request, extension_return_value)
85 .await?;
86
87 return Ok(0);
88 } else if first_request
89 .options
90 .contains_key(LOGICAL_TABLE_METADATA_KEY)
91 {
92 if requests.len() == 1 {
93 let request = &requests.first().unwrap().1;
94 let physical_region_id = parse_physical_region_id(request)?;
95 let mut manifest_infos = Vec::with_capacity(1);
96 self.create_logical_regions(physical_region_id, requests, extension_return_value)
97 .await?;
98 append_manifest_info(&self.mito, physical_region_id, &mut manifest_infos);
99 encode_manifest_info_to_extensions(&manifest_infos, extension_return_value)?;
100 } else {
101 let grouped_requests =
102 group_create_logical_region_requests_by_physical_region_id(requests)?;
103 let mut manifest_infos = Vec::with_capacity(grouped_requests.len());
104 for (physical_region_id, requests) in grouped_requests {
105 self.create_logical_regions(
106 physical_region_id,
107 requests,
108 extension_return_value,
109 )
110 .await?;
111 append_manifest_info(&self.mito, physical_region_id, &mut manifest_infos);
112 }
113 encode_manifest_info_to_extensions(&manifest_infos, extension_return_value)?;
114 }
115 } else {
116 return MissingRegionOptionSnafu {}.fail();
117 }
118
119 Ok(0)
120 }
121
122 async fn create_physical_region(
124 &self,
125 region_id: RegionId,
126 request: RegionCreateRequest,
127 extension_return_value: &mut HashMap<String, Vec<u8>>,
128 ) -> Result<()> {
129 let physical_region_options = PhysicalRegionOptions::try_from(&request.options)?;
130 let (data_region_id, metadata_region_id) = Self::transform_region_id(region_id);
131
132 let create_metadata_region_request = self.create_request_for_metadata_region(&request);
134 self.mito
135 .handle_request(
136 metadata_region_id,
137 RegionRequest::Create(create_metadata_region_request),
138 )
139 .await
140 .with_context(|_| CreateMitoRegionSnafu {
141 region_type: METADATA_REGION_SUBDIR,
142 })?;
143
144 let create_data_region_request = self.create_request_for_data_region(&request);
146 let physical_columns = create_data_region_request
147 .column_metadatas
148 .iter()
149 .map(|metadata| (metadata.column_schema.name.clone(), metadata.clone()))
150 .collect::<HashMap<_, _>>();
151 let time_index_unit = create_data_region_request
152 .column_metadatas
153 .iter()
154 .find_map(|metadata| {
155 if metadata.semantic_type == SemanticType::Timestamp {
156 metadata
157 .column_schema
158 .data_type
159 .as_timestamp()
160 .map(|data_type| data_type.unit())
161 } else {
162 None
163 }
164 })
165 .context(UnexpectedRequestSnafu {
166 reason: "No time index column found",
167 })?;
168 let response = self
169 .mito
170 .handle_request(
171 data_region_id,
172 RegionRequest::Create(create_data_region_request),
173 )
174 .await
175 .with_context(|_| CreateMitoRegionSnafu {
176 region_type: DATA_REGION_SUBDIR,
177 })?;
178 let primary_key_encoding = self.mito.get_primary_key_encoding(data_region_id).context(
179 PhysicalRegionNotFoundSnafu {
180 region_id: data_region_id,
181 },
182 )?;
183 extension_return_value.extend(response.extensions);
184
185 info!(
186 "Created physical metric region {region_id}, primary key encoding={primary_key_encoding}, physical_region_options={physical_region_options:?}"
187 );
188 PHYSICAL_REGION_COUNT.inc();
189
190 self.state.write().unwrap().add_physical_region(
192 data_region_id,
193 physical_columns,
194 primary_key_encoding,
195 physical_region_options,
196 time_index_unit,
197 );
198
199 Ok(())
200 }
201
202 async fn create_logical_regions(
204 &self,
205 physical_region_id: RegionId,
206 requests: Vec<(RegionId, RegionCreateRequest)>,
207 extension_return_value: &mut HashMap<String, Vec<u8>>,
208 ) -> Result<()> {
209 let data_region_id = utils::to_data_region_id(physical_region_id);
210
211 let unit = self
212 .state
213 .read()
214 .unwrap()
215 .physical_region_time_index_unit(physical_region_id)
216 .context(PhysicalRegionNotFoundSnafu {
217 region_id: data_region_id,
218 })?;
219 for (_, request) in &requests {
221 let time_index_column = request
223 .column_metadatas
224 .iter()
225 .find(|col| col.semantic_type == SemanticType::Timestamp)
226 .unwrap();
227 let request_unit = time_index_column
228 .column_schema
229 .data_type
230 .as_timestamp()
231 .unwrap()
232 .unit();
233 ensure!(
234 request_unit == unit,
235 UnexpectedRequestSnafu {
236 reason: format!(
237 "Metric has differenttime unit ({:?}) than the physical region ({:?})",
238 request_unit, unit
239 ),
240 }
241 );
242 }
243
244 let requests = {
246 let state = self.state.read().unwrap();
247 let mut skipped = Vec::with_capacity(requests.len());
248 let mut kept_requests = Vec::with_capacity(requests.len());
249
250 for (region_id, request) in requests {
251 if state.is_logical_region_exist(region_id) {
252 skipped.push(region_id);
253 } else {
254 kept_requests.push((region_id, request));
255 }
256 }
257
258 if !skipped.is_empty() {
260 info!(
261 "Skipped creating logical regions {skipped:?} because they already exist",
262 skipped = skipped
263 );
264 }
265 kept_requests
266 };
267
268 let mut new_column_names = HashSet::new();
270 let mut new_columns = Vec::new();
271
272 let index_option = {
273 let state = &self.state.read().unwrap();
274 let region_state = state
275 .physical_region_states()
276 .get(&data_region_id)
277 .with_context(|| PhysicalRegionNotFoundSnafu {
278 region_id: data_region_id,
279 })?;
280 let physical_columns = region_state.physical_columns();
281
282 extract_new_columns(
283 &requests,
284 physical_columns,
285 &mut new_column_names,
286 &mut new_columns,
287 )?;
288 region_state.options().index
289 };
290
291 self.data_region
293 .add_columns(data_region_id, new_columns, index_option)
294 .await?;
295
296 let physical_columns = self.data_region.physical_columns(data_region_id).await?;
297 let physical_schema_map = physical_columns
298 .iter()
299 .map(|metadata| (metadata.column_schema.name.as_str(), metadata))
300 .collect::<HashMap<_, _>>();
301 let logical_regions = requests
302 .iter()
303 .map(|(region_id, _)| *region_id)
304 .collect::<Vec<_>>();
305 let logical_region_columns = requests.iter().map(|(region_id, request)| {
306 (
307 *region_id,
308 request
309 .column_metadatas
310 .iter()
311 .map(|metadata| {
312 let column_metadata = *physical_schema_map
314 .get(metadata.column_schema.name.as_str())
315 .unwrap();
316 (metadata.column_schema.name.as_str(), column_metadata)
317 })
318 .collect::<HashMap<_, _>>(),
319 )
320 });
321
322 let new_add_columns = new_column_names.iter().map(|name| {
323 let column_metadata = *physical_schema_map.get(name).unwrap();
325 (name.to_string(), column_metadata.clone())
326 });
327
328 extension_return_value.insert(
329 ALTER_PHYSICAL_EXTENSION_KEY.to_string(),
330 ColumnMetadata::encode_list(&physical_columns).context(SerializeColumnMetadataSnafu)?,
331 );
332
333 self.metadata_region
335 .add_logical_regions(physical_region_id, true, logical_region_columns)
336 .await?;
337
338 {
339 let mut state = self.state.write().unwrap();
340 state.add_physical_columns(data_region_id, new_add_columns);
341 state.add_logical_regions(physical_region_id, logical_regions.clone());
342 }
343 for logical_region_id in logical_regions {
344 self.metadata_region
345 .open_logical_region(logical_region_id)
346 .await;
347 }
348
349 Ok(())
350 }
351
352 fn verify_region_create_request(request: &RegionCreateRequest) -> Result<()> {
357 request.validate().context(InvalidMetadataSnafu)?;
358
359 let name_to_index = request
360 .column_metadatas
361 .iter()
362 .enumerate()
363 .map(|(idx, metadata)| (metadata.column_schema.name.clone(), idx))
364 .collect::<HashMap<String, usize>>();
365
366 let table_id_col_def = request.column_metadatas.iter().any(is_metric_name_col);
367 let tsid_col_def = request.column_metadatas.iter().any(is_tsid_col);
368
369 ensure!(
371 !name_to_index.contains_key(DATA_SCHEMA_TABLE_ID_COLUMN_NAME) || table_id_col_def,
372 InternalColumnOccupiedSnafu {
373 column: DATA_SCHEMA_TABLE_ID_COLUMN_NAME,
374 }
375 );
376 ensure!(
377 !name_to_index.contains_key(DATA_SCHEMA_TSID_COLUMN_NAME) || tsid_col_def,
378 InternalColumnOccupiedSnafu {
379 column: DATA_SCHEMA_TSID_COLUMN_NAME,
380 }
381 );
382
383 ensure!(
385 request.is_physical_table() || request.options.contains_key(LOGICAL_TABLE_METADATA_KEY),
386 MissingRegionOptionSnafu {}
387 );
388 ensure!(
389 !(request.is_physical_table()
390 && request.options.contains_key(LOGICAL_TABLE_METADATA_KEY)),
391 ConflictRegionOptionSnafu {}
392 );
393
394 let mut field_cols = Vec::new();
396 for col in &request.column_metadatas {
397 if is_metric_engine_internal_column(&col.column_schema.name) {
399 continue;
400 }
401 match col.semantic_type {
402 SemanticType::Tag => ensure!(
403 col.column_schema.data_type == ConcreteDataType::string_datatype(),
404 ColumnTypeMismatchSnafu {
405 expect: ConcreteDataType::string_datatype(),
406 actual: col.column_schema.data_type.clone(),
407 }
408 ),
409 SemanticType::Field => {
410 field_cols.push(col);
411 }
412 SemanticType::Timestamp => {}
413 }
414 }
415 let [field_col] = field_cols.as_slice() else {
416 if field_cols.is_empty() {
417 NoFieldColumnSnafu.fail()?;
418 }
419 return MultipleFieldColumnSnafu {
420 previous: field_cols[0].column_schema.name.clone(),
421 current: field_cols[1].column_schema.name.clone(),
422 }
423 .fail();
424 };
425
426 if is_native_histogram_value_type(&field_col.column_schema.data_type) {
427 return Ok(());
428 }
429
430 ensure!(
432 field_col.column_schema.data_type == ConcreteDataType::float64_datatype(),
433 ColumnTypeMismatchSnafu {
434 expect: ConcreteDataType::float64_datatype(),
435 actual: field_col.column_schema.data_type.clone(),
436 }
437 );
438
439 Ok(())
440 }
441
442 fn transform_region_id(region_id: RegionId) -> (RegionId, RegionId) {
446 (
447 to_data_region_id(region_id),
448 to_metadata_region_id(region_id),
449 )
450 }
451
452 pub fn create_request_for_metadata_region(
456 &self,
457 request: &RegionCreateRequest,
458 ) -> RegionCreateRequest {
459 let timestamp_column_metadata = ColumnMetadata {
461 column_id: METADATA_SCHEMA_TIMESTAMP_COLUMN_INDEX as _,
462 semantic_type: SemanticType::Timestamp,
463 column_schema: ColumnSchema::new(
464 METADATA_SCHEMA_TIMESTAMP_COLUMN_NAME,
465 ConcreteDataType::timestamp_millisecond_datatype(),
466 false,
467 )
468 .with_default_constraint(Some(datatypes::schema::ColumnDefaultConstraint::Value(
469 Value::Timestamp(Timestamp::new_millisecond(0)),
470 )))
471 .unwrap(),
472 };
473 let key_column_metadata = ColumnMetadata {
475 column_id: METADATA_SCHEMA_KEY_COLUMN_INDEX as _,
476 semantic_type: SemanticType::Tag,
477 column_schema: ColumnSchema::new(
478 METADATA_SCHEMA_KEY_COLUMN_NAME,
479 ConcreteDataType::string_datatype(),
480 false,
481 ),
482 };
483 let value_column_metadata = ColumnMetadata {
485 column_id: METADATA_SCHEMA_VALUE_COLUMN_INDEX as _,
486 semantic_type: SemanticType::Field,
487 column_schema: ColumnSchema::new(
488 METADATA_SCHEMA_VALUE_COLUMN_NAME,
489 ConcreteDataType::string_datatype(),
490 true,
491 ),
492 };
493
494 let options = region_options_for_metadata_region(&request.options);
495 RegionCreateRequest {
496 engine: MITO_ENGINE_NAME.to_string(),
497 column_metadatas: vec![
498 timestamp_column_metadata,
499 key_column_metadata,
500 value_column_metadata,
501 ],
502 primary_key: vec![METADATA_SCHEMA_KEY_COLUMN_INDEX as _],
503 options,
504 table_dir: request.table_dir.clone(),
505 path_type: PathType::Metadata,
506 partition_expr_json: Some("".to_string()),
507 requirements: request.requirements,
508 }
509 }
510
511 pub fn create_request_for_data_region(
518 &self,
519 request: &RegionCreateRequest,
520 ) -> RegionCreateRequest {
521 let mut data_region_request = request.clone();
522 let mut primary_key = vec![ReservedColumnId::table_id(), ReservedColumnId::tsid()];
523
524 data_region_request.table_dir = request.table_dir.clone();
525 data_region_request.path_type = PathType::Data;
526
527 let table_id_col_def = request.column_metadatas.iter().any(is_metric_name_col);
528 let tsid_col_def = request.column_metadatas.iter().any(is_tsid_col);
529
530 data_region_request
532 .column_metadatas
533 .iter_mut()
534 .for_each(|metadata| {
535 if metadata.semantic_type == SemanticType::Tag
536 && !is_metric_name_col(metadata)
537 && !is_tsid_col(metadata)
538 {
539 metadata.column_schema.set_nullable();
540 primary_key.push(metadata.column_id);
541 }
542 });
543
544 if !table_id_col_def {
546 data_region_request.column_metadatas.push(table_id_col());
547 }
548 if !tsid_col_def {
549 data_region_request.column_metadatas.push(tsid_col());
550 }
551 data_region_request.primary_key = primary_key;
552
553 set_data_region_options(&mut data_region_request.options);
555
556 data_region_request
557 }
558}
559
560fn table_id_col() -> ColumnMetadata {
561 ColumnMetadata {
562 column_id: ReservedColumnId::table_id(),
563 semantic_type: SemanticType::Tag,
564 column_schema: ColumnSchema::new(
565 DATA_SCHEMA_TABLE_ID_COLUMN_NAME,
566 ConcreteDataType::uint32_datatype(),
567 false,
568 )
569 .with_skipping_options(SkippingIndexOptions::new_unchecked(
570 DEFAULT_TABLE_ID_SKIPPING_INDEX_GRANULARITY,
571 DEFAULT_TABLE_ID_SKIPPING_INDEX_FALSE_POSITIVE_RATE,
572 datatypes::schema::SkippingIndexType::BloomFilter,
573 ))
574 .unwrap(),
575 }
576}
577
578fn tsid_col() -> ColumnMetadata {
579 ColumnMetadata {
580 column_id: ReservedColumnId::tsid(),
581 semantic_type: SemanticType::Tag,
582 column_schema: ColumnSchema::new(
583 DATA_SCHEMA_TSID_COLUMN_NAME,
584 ConcreteDataType::uint64_datatype(),
585 false,
586 )
587 .with_inverted_index(false),
588 }
589}
590
591pub(crate) fn is_metric_name_col(column: &ColumnMetadata) -> bool {
593 column.column_id == ReservedColumnId::table_id()
594 && column.semantic_type == SemanticType::Tag
595 && column.column_schema.data_type == ConcreteDataType::uint32_datatype()
596 && column.column_schema.name == DATA_SCHEMA_TABLE_ID_COLUMN_NAME
597 && !column.column_schema.is_nullable()
598}
599
600pub(crate) fn is_tsid_col(column: &ColumnMetadata) -> bool {
602 column.column_id == ReservedColumnId::tsid()
603 && column.semantic_type == SemanticType::Tag
604 && column.column_schema.data_type == ConcreteDataType::uint64_datatype()
605 && column.column_schema.name == DATA_SCHEMA_TSID_COLUMN_NAME
606 && !column.column_schema.is_nullable()
607}
608
609fn group_create_logical_region_requests_by_physical_region_id(
611 requests: Vec<(RegionId, RegionCreateRequest)>,
612) -> Result<HashMap<RegionId, Vec<(RegionId, RegionCreateRequest)>>> {
613 let mut result = HashMap::with_capacity(requests.len());
614 for (region_id, request) in requests {
615 let physical_region_id = parse_physical_region_id(&request)?;
616 result
617 .entry(physical_region_id)
618 .or_insert_with(Vec::new)
619 .push((region_id, request));
620 }
621
622 Ok(result)
623}
624
625fn parse_physical_region_id(request: &RegionCreateRequest) -> Result<RegionId> {
627 let physical_region_id_raw = request
628 .options
629 .get(LOGICAL_TABLE_METADATA_KEY)
630 .ok_or(MissingRegionOptionSnafu {}.build())?;
631
632 let physical_region_id: RegionId = physical_region_id_raw
633 .parse::<u64>()
634 .with_context(|_| ParseRegionIdSnafu {
635 raw: physical_region_id_raw,
636 })?
637 .into();
638
639 Ok(physical_region_id)
640}
641
642pub(crate) fn region_options_for_metadata_region(
644 original: &HashMap<String, String>,
645) -> HashMap<String, String> {
646 let mut metadata_region_options = HashMap::new();
647 metadata_region_options.insert(TTL_KEY.to_string(), FOREVER.to_string());
648
649 if let Some(wal_options) = original.get(WAL_OPTIONS_KEY) {
650 metadata_region_options.insert(WAL_OPTIONS_KEY.to_string(), wal_options.clone());
651 }
652
653 metadata_region_options
654}
655
656#[cfg(test)]
657mod test {
658 use common_meta::ddl::test_util::assert_column_name_and_id;
659 use common_meta::ddl::utils::{parse_column_metadatas, parse_manifest_infos_from_extensions};
660 use common_query::native_histogram::native_histogram_value_type;
661 use common_query::prelude::{greptime_native_histogram, greptime_timestamp, greptime_value};
662 use store_api::metric_engine_consts::{METRIC_ENGINE_NAME, PHYSICAL_TABLE_METADATA_KEY};
663 use store_api::region_request::{BatchRegionDdlRequest, RegionRequirements};
664
665 use super::*;
666 use crate::config::EngineConfig;
667 use crate::engine::MetricEngine;
668 use crate::test_util::{TestEnv, create_logical_region_request};
669
670 #[test]
671 fn test_internal_column_metadata() {
672 let table_id_col = table_id_col();
673 let tsid_col = tsid_col();
674 assert!(is_metric_name_col(&table_id_col));
675 assert!(is_tsid_col(&tsid_col));
676 }
677
678 #[test]
679 fn test_verify_region_create_request() {
680 let request = RegionCreateRequest {
682 column_metadatas: vec![
683 ColumnMetadata {
684 column_id: 0,
685 semantic_type: SemanticType::Timestamp,
686 column_schema: ColumnSchema::new(
687 METADATA_SCHEMA_TIMESTAMP_COLUMN_NAME,
688 ConcreteDataType::timestamp_millisecond_datatype(),
689 false,
690 ),
691 },
692 ColumnMetadata {
693 column_id: 1,
694 semantic_type: SemanticType::Tag,
695 column_schema: ColumnSchema::new(
696 DATA_SCHEMA_TABLE_ID_COLUMN_NAME,
697 ConcreteDataType::uint32_datatype(),
698 false,
699 ),
700 },
701 ],
702 table_dir: "test_dir".to_string(),
703 path_type: PathType::Bare,
704 engine: METRIC_ENGINE_NAME.to_string(),
705 primary_key: vec![],
706 options: HashMap::new(),
707 partition_expr_json: Some("".to_string()),
708 requirements: RegionRequirements::object_storage(),
709 };
710 let result = MetricEngineInner::verify_region_create_request(&request);
711 assert!(result.is_err());
712 assert_eq!(
713 result.unwrap_err().to_string(),
714 "Internal column __table_id is reserved".to_string()
715 );
716
717 let request = RegionCreateRequest {
719 column_metadatas: vec![
720 ColumnMetadata {
721 column_id: 0,
722 semantic_type: SemanticType::Timestamp,
723 column_schema: ColumnSchema::new(
724 METADATA_SCHEMA_TIMESTAMP_COLUMN_NAME,
725 ConcreteDataType::timestamp_millisecond_datatype(),
726 false,
727 ),
728 },
729 ColumnMetadata {
730 column_id: 1,
731 semantic_type: SemanticType::Tag,
732 column_schema: ColumnSchema::new(
733 "column1".to_string(),
734 ConcreteDataType::string_datatype(),
735 false,
736 ),
737 },
738 ColumnMetadata {
739 column_id: 2,
740 semantic_type: SemanticType::Field,
741 column_schema: ColumnSchema::new(
742 "column2".to_string(),
743 ConcreteDataType::float64_datatype(),
744 false,
745 ),
746 },
747 table_id_col(),
748 tsid_col(),
749 ],
750 table_dir: "test_dir".to_string(),
751 path_type: PathType::Bare,
752 engine: METRIC_ENGINE_NAME.to_string(),
753 primary_key: vec![],
754 options: [(PHYSICAL_TABLE_METADATA_KEY.to_string(), String::new())]
755 .into_iter()
756 .collect(),
757 partition_expr_json: Some("".to_string()),
758 requirements: Default::default(),
759 };
760 MetricEngineInner::verify_region_create_request(&request).unwrap();
761
762 let request = RegionCreateRequest {
764 column_metadatas: vec![
765 ColumnMetadata {
766 column_id: 0,
767 semantic_type: SemanticType::Timestamp,
768 column_schema: ColumnSchema::new(
769 METADATA_SCHEMA_TIMESTAMP_COLUMN_NAME,
770 ConcreteDataType::timestamp_millisecond_datatype(),
771 false,
772 ),
773 },
774 ColumnMetadata {
775 column_id: 1,
776 semantic_type: SemanticType::Tag,
777 column_schema: ColumnSchema::new(
778 "column1".to_string(),
779 ConcreteDataType::string_datatype(),
780 false,
781 ),
782 },
783 ColumnMetadata {
784 column_id: 2,
785 semantic_type: SemanticType::Field,
786 column_schema: ColumnSchema::new(
787 "column2".to_string(),
788 ConcreteDataType::float64_datatype(),
789 false,
790 ),
791 },
792 ],
793 table_dir: "test_dir".to_string(),
794 path_type: PathType::Bare,
795 engine: METRIC_ENGINE_NAME.to_string(),
796 primary_key: vec![],
797 options: [(PHYSICAL_TABLE_METADATA_KEY.to_string(), String::new())]
798 .into_iter()
799 .collect(),
800 partition_expr_json: Some("".to_string()),
801 requirements: Default::default(),
802 };
803 MetricEngineInner::verify_region_create_request(&request).unwrap();
804 }
805
806 #[test]
807 fn test_verify_region_create_request_options() {
808 let mut request = RegionCreateRequest {
809 column_metadatas: vec![
810 ColumnMetadata {
811 column_id: 0,
812 semantic_type: SemanticType::Timestamp,
813 column_schema: ColumnSchema::new(
814 METADATA_SCHEMA_TIMESTAMP_COLUMN_NAME,
815 ConcreteDataType::timestamp_millisecond_datatype(),
816 false,
817 ),
818 },
819 ColumnMetadata {
820 column_id: 1,
821 semantic_type: SemanticType::Field,
822 column_schema: ColumnSchema::new(
823 "val".to_string(),
824 ConcreteDataType::float64_datatype(),
825 false,
826 ),
827 },
828 ],
829 table_dir: "test_dir".to_string(),
830 path_type: PathType::Bare,
831 engine: METRIC_ENGINE_NAME.to_string(),
832 primary_key: vec![],
833 options: HashMap::new(),
834 partition_expr_json: Some("".to_string()),
835 requirements: Default::default(),
836 };
837 MetricEngineInner::verify_region_create_request(&request).unwrap_err();
838
839 let mut options = HashMap::new();
840 options.insert(PHYSICAL_TABLE_METADATA_KEY.to_string(), "value".to_string());
841 request.options.clone_from(&options);
842 MetricEngineInner::verify_region_create_request(&request).unwrap();
843
844 options.insert(LOGICAL_TABLE_METADATA_KEY.to_string(), "value".to_string());
845 request.options.clone_from(&options);
846 MetricEngineInner::verify_region_create_request(&request).unwrap_err();
847
848 options.remove(PHYSICAL_TABLE_METADATA_KEY).unwrap();
849 request.options = options;
850 MetricEngineInner::verify_region_create_request(&request).unwrap();
851 }
852
853 #[test]
854 fn test_verify_region_create_request_native_histogram_fields() {
855 let native_histogram_columns = vec![
856 ColumnMetadata {
857 column_id: 0,
858 semantic_type: SemanticType::Timestamp,
859 column_schema: ColumnSchema::new(
860 greptime_timestamp(),
861 ConcreteDataType::timestamp_millisecond_datatype(),
862 false,
863 ),
864 },
865 ColumnMetadata {
866 column_id: 1,
867 semantic_type: SemanticType::Tag,
868 column_schema: ColumnSchema::new("job", ConcreteDataType::string_datatype(), true),
869 },
870 ColumnMetadata {
871 column_id: 2,
872 semantic_type: SemanticType::Field,
873 column_schema: ColumnSchema::new(
874 greptime_native_histogram(),
875 native_histogram_value_type().clone(),
876 true,
877 ),
878 },
879 ];
880 let request = RegionCreateRequest {
881 column_metadatas: native_histogram_columns,
882 table_dir: "test_dir".to_string(),
883 path_type: PathType::Bare,
884 engine: METRIC_ENGINE_NAME.to_string(),
885 primary_key: vec![],
886 options: [(
887 LOGICAL_TABLE_METADATA_KEY.to_string(),
888 "physical".to_string(),
889 )]
890 .into_iter()
891 .collect(),
892 partition_expr_json: Some("".to_string()),
893 requirements: Default::default(),
894 };
895 MetricEngineInner::verify_region_create_request(&request).unwrap();
896
897 let request = RegionCreateRequest {
898 column_metadatas: vec![
899 ColumnMetadata {
900 column_id: 0,
901 semantic_type: SemanticType::Timestamp,
902 column_schema: ColumnSchema::new(
903 greptime_timestamp(),
904 ConcreteDataType::timestamp_millisecond_datatype(),
905 false,
906 ),
907 },
908 ColumnMetadata {
909 column_id: 1,
910 semantic_type: SemanticType::Field,
911 column_schema: ColumnSchema::new(
912 "value_a",
913 ConcreteDataType::float64_datatype(),
914 true,
915 ),
916 },
917 ColumnMetadata {
918 column_id: 2,
919 semantic_type: SemanticType::Field,
920 column_schema: ColumnSchema::new(
921 "value_b",
922 ConcreteDataType::float64_datatype(),
923 true,
924 ),
925 },
926 ],
927 table_dir: "test_dir".to_string(),
928 path_type: PathType::Bare,
929 engine: METRIC_ENGINE_NAME.to_string(),
930 primary_key: vec![],
931 options: [(
932 LOGICAL_TABLE_METADATA_KEY.to_string(),
933 "physical".to_string(),
934 )]
935 .into_iter()
936 .collect(),
937 partition_expr_json: Some("".to_string()),
938 requirements: Default::default(),
939 };
940 assert!(MetricEngineInner::verify_region_create_request(&request).is_err());
941 }
942
943 #[tokio::test]
944 async fn test_create_request_for_physical_regions() {
945 let options: HashMap<_, _> = [
947 ("ttl".to_string(), "60m".to_string()),
948 ("skip_wal".to_string(), "true".to_string()),
949 ]
950 .into_iter()
951 .collect();
952 let request = RegionCreateRequest {
953 engine: METRIC_ENGINE_NAME.to_string(),
954 column_metadatas: vec![
955 ColumnMetadata {
956 column_id: 0,
957 semantic_type: SemanticType::Timestamp,
958 column_schema: ColumnSchema::new(
959 "timestamp",
960 ConcreteDataType::timestamp_millisecond_datatype(),
961 false,
962 ),
963 },
964 ColumnMetadata {
965 column_id: 1,
966 semantic_type: SemanticType::Tag,
967 column_schema: ColumnSchema::new(
968 "tag",
969 ConcreteDataType::string_datatype(),
970 false,
971 ),
972 },
973 ],
974 primary_key: vec![0],
975 options,
976 table_dir: "/test_dir".to_string(),
977 path_type: PathType::Bare,
978 partition_expr_json: Some("".to_string()),
979 requirements: RegionRequirements::object_storage(),
980 };
981
982 let env = TestEnv::new().await;
984 let engine = MetricEngine::try_new(env.mito(), EngineConfig::default()).unwrap();
985 let engine_inner = engine.inner;
986
987 let data_region_request = engine_inner.create_request_for_data_region(&request);
989 assert_eq!(data_region_request.table_dir, "/test_dir".to_string());
990 assert_eq!(data_region_request.path_type, PathType::Data);
991 assert_eq!(data_region_request.column_metadatas.len(), 4);
992 assert_eq!(
993 data_region_request.primary_key,
994 vec![ReservedColumnId::table_id(), ReservedColumnId::tsid(), 1]
995 );
996 assert!(data_region_request.options.contains_key("ttl"));
997 assert_eq!(
998 data_region_request.requirements,
999 RegionRequirements::object_storage()
1000 );
1001
1002 let metadata_region_request = engine_inner.create_request_for_metadata_region(&request);
1004 assert_eq!(metadata_region_request.table_dir, "/test_dir".to_string());
1005 assert_eq!(metadata_region_request.path_type, PathType::Metadata);
1006 assert_eq!(
1007 metadata_region_request.options.get("ttl").unwrap(),
1008 "forever"
1009 );
1010 assert!(!metadata_region_request.options.contains_key("skip_wal"));
1011 assert_eq!(
1012 metadata_region_request.requirements,
1013 RegionRequirements::object_storage()
1014 );
1015 }
1016
1017 #[tokio::test]
1018 async fn test_create_request_for_physical_regions_with_internal_columns() {
1019 let options: HashMap<_, _> = [
1020 ("ttl".to_string(), "60m".to_string()),
1021 ("skip_wal".to_string(), "true".to_string()),
1022 ]
1023 .into_iter()
1024 .collect();
1025 let request = RegionCreateRequest {
1026 engine: METRIC_ENGINE_NAME.to_string(),
1027 column_metadatas: vec![
1028 ColumnMetadata {
1029 column_id: 0,
1030 semantic_type: SemanticType::Timestamp,
1031 column_schema: ColumnSchema::new(
1032 "timestamp",
1033 ConcreteDataType::timestamp_millisecond_datatype(),
1034 false,
1035 ),
1036 },
1037 ColumnMetadata {
1038 column_id: 1,
1039 semantic_type: SemanticType::Tag,
1040 column_schema: ColumnSchema::new(
1041 "tag",
1042 ConcreteDataType::string_datatype(),
1043 false,
1044 ),
1045 },
1046 ColumnMetadata {
1047 column_id: 2,
1048 semantic_type: SemanticType::Field,
1049 column_schema: ColumnSchema::new(
1050 "value",
1051 ConcreteDataType::float64_datatype(),
1052 false,
1053 ),
1054 },
1055 table_id_col(),
1056 tsid_col(),
1057 ],
1058 primary_key: vec![0],
1059 options,
1060 table_dir: "/test_dir".to_string(),
1061 path_type: PathType::Bare,
1062 partition_expr_json: Some("".to_string()),
1063 requirements: Default::default(),
1064 };
1065
1066 let env = TestEnv::new().await;
1067 let engine = MetricEngine::try_new(env.mito(), EngineConfig::default()).unwrap();
1068 let engine_inner = engine.inner;
1069
1070 let data_region_request = engine_inner.create_request_for_data_region(&request);
1071 assert_eq!(data_region_request.column_metadatas.len(), 5);
1072 assert_eq!(
1073 data_region_request.primary_key,
1074 vec![ReservedColumnId::table_id(), ReservedColumnId::tsid(), 1]
1075 );
1076
1077 let table_id_count = data_region_request
1078 .column_metadatas
1079 .iter()
1080 .filter(|metadata| metadata.column_schema.name == DATA_SCHEMA_TABLE_ID_COLUMN_NAME)
1081 .count();
1082 let tsid_count = data_region_request
1083 .column_metadatas
1084 .iter()
1085 .filter(|metadata| metadata.column_schema.name == DATA_SCHEMA_TSID_COLUMN_NAME)
1086 .count();
1087 assert_eq!(table_id_count, 1);
1088 assert_eq!(tsid_count, 1);
1089
1090 let tag_metadata = data_region_request
1091 .column_metadatas
1092 .iter()
1093 .find(|metadata| metadata.column_schema.name == "tag")
1094 .unwrap();
1095 assert!(tag_metadata.column_schema.is_nullable());
1096
1097 let table_id_metadata = data_region_request
1098 .column_metadatas
1099 .iter()
1100 .find(|metadata| metadata.column_schema.name == DATA_SCHEMA_TABLE_ID_COLUMN_NAME)
1101 .unwrap();
1102 assert!(is_metric_name_col(table_id_metadata));
1103
1104 let tsid_metadata = data_region_request
1105 .column_metadatas
1106 .iter()
1107 .find(|metadata| metadata.column_schema.name == DATA_SCHEMA_TSID_COLUMN_NAME)
1108 .unwrap();
1109 assert!(is_tsid_col(tsid_metadata));
1110 }
1111
1112 #[tokio::test]
1113 async fn test_create_logical_regions() {
1114 let env = TestEnv::new().await;
1115 let engine = env.metric();
1116 let physical_region_id1 = RegionId::new(1024, 0);
1117 let physical_region_id2 = RegionId::new(1024, 1);
1118 let logical_region_id1 = RegionId::new(1025, 0);
1119 let logical_region_id2 = RegionId::new(1025, 1);
1120 env.create_physical_region(physical_region_id1, "/test_dir1", vec![])
1121 .await;
1122 env.create_physical_region(physical_region_id2, "/test_dir2", vec![])
1123 .await;
1124
1125 let region_create_request1 =
1126 create_logical_region_request(&["job"], physical_region_id1, "logical1");
1127 let region_create_request2 =
1128 create_logical_region_request(&["job"], physical_region_id2, "logical2");
1129
1130 let response = engine
1131 .handle_batch_ddl_requests(BatchRegionDdlRequest::Create(vec![
1132 (logical_region_id1, region_create_request1),
1133 (logical_region_id2, region_create_request2),
1134 ]))
1135 .await
1136 .unwrap();
1137
1138 let manifest_infos = parse_manifest_infos_from_extensions(&response.extensions).unwrap();
1139 assert_eq!(manifest_infos.len(), 2);
1140 let region_ids = manifest_infos.into_iter().map(|i| i.0).collect::<Vec<_>>();
1141 assert!(region_ids.contains(&physical_region_id1));
1142 assert!(region_ids.contains(&physical_region_id2));
1143
1144 let column_metadatas =
1145 parse_column_metadatas(&response.extensions, ALTER_PHYSICAL_EXTENSION_KEY).unwrap();
1146 assert_column_name_and_id(
1147 &column_metadatas,
1148 &[
1149 (greptime_timestamp(), 0),
1150 (greptime_value(), 1),
1151 ("__table_id", ReservedColumnId::table_id()),
1152 ("__tsid", ReservedColumnId::tsid()),
1153 ("job", 2),
1154 ],
1155 );
1156 }
1157}