Skip to main content

metric_engine/engine/
create.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
15mod 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    /// Initialize a physical metric region at given region id.
123    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        // create metadata region
133        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        // create data region
145        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        // remember this table
191        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    /// Create multiple logical regions on the same physical region.
203    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        // Checks the time index unit of each request.
220        for (_, request) in &requests {
221            // Safety: verify_region_create_request() ensures that the request is valid.
222            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        // Filters out the requests that the logical region already exists
245        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            // log skipped regions
259            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        // Finds new columns to add to physical region
269        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        // TODO(weny): we dont need to pass a mutable new_columns here.
292        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                        // Safety: previous steps ensure the physical region exist
313                        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            // Safety: previous steps ensure the physical region exist
324            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        // Writes logical regions metadata to metadata region
334        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    /// Check if
353    /// - internal columns are not occupied
354    /// - required table option is present ([PHYSICAL_TABLE_METADATA_KEY] or
355    ///   [LOGICAL_TABLE_METADATA_KEY])
356    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        // check if internal columns are not occupied or defined in the request
370        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        // check if required table option is present
384        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        // check if field columns are either a normal metric value or a native histogram.
395        let mut field_cols = Vec::new();
396        for col in &request.column_metadatas {
397            // Verified in above steps.
398            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        // make sure the normal field column is float64 type
431        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    /// Build data region id and metadata region id from the given region id.
443    ///
444    /// Return value: (data_region_id, metadata_region_id)
445    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    /// Build [RegionCreateRequest] for metadata region
453    ///
454    /// This method will append [METADATA_REGION_SUBDIR] to the given `region_dir`.
455    pub fn create_request_for_metadata_region(
456        &self,
457        request: &RegionCreateRequest,
458    ) -> RegionCreateRequest {
459        // ts TIME INDEX DEFAULT 0
460        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        // key STRING PRIMARY KEY
474        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        // val STRING
484        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    /// Convert [RegionCreateRequest] for data region.
512    ///
513    /// All tag columns in the original request will be converted to value columns.
514    /// Those columns real semantic type is stored in metadata region.
515    ///
516    /// This will also add internal columns to the request.
517    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        // change nullability for tag columns
531        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        // add internal columns if not defined in the request
545        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
554        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
591/// Returns true if the column is the metric name column.
592pub(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
600/// Returns true if the column is the tsid column.
601pub(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
609/// Groups the create logical region requests by physical region id.
610fn 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
625/// Parses the physical region id from the request.
626fn 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
642/// Creates the region options for metadata region in metric engine.
643pub(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        // internal column is occupied
681        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        // allow reserved internal columns when defined properly
718        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        // valid request
763        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        // original request
946        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        // set up
983        let env = TestEnv::new().await;
984        let engine = MetricEngine::try_new(env.mito(), EngineConfig::default()).unwrap();
985        let engine_inner = engine.inner;
986
987        // check create data region request
988        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        // check create metadata region request
1003        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}