Skip to main content

common_meta/reconciliation/
utils.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::collections::{HashMap, HashSet};
16use std::fmt::{self, Display};
17use std::ops::AddAssign;
18use std::sync::Arc;
19use std::time::Instant;
20
21use api::v1::SemanticType;
22use common_procedure::{Context as ProcedureContext, ProcedureId, watcher};
23use common_telemetry::{error, warn};
24use datatypes::schema::{ColumnSchema, Schema};
25use futures::future::{join_all, try_join_all};
26use snafu::{OptionExt, ResultExt, ensure};
27use store_api::metadata::{ColumnMetadata, RegionMetadata};
28use store_api::storage::consts::ReservedColumnId;
29use store_api::storage::{RegionId, TableId};
30use table::metadata::{TableInfo, TableMeta};
31use table::table_name::TableName;
32use table::table_reference::TableReference;
33
34use crate::cache_invalidator::CacheInvalidatorRef;
35use crate::error::{
36    ColumnIdMismatchSnafu, ColumnNotFoundSnafu, MismatchColumnIdSnafu,
37    MissingColumnInColumnMetadataSnafu, ProcedureStateReceiverNotFoundSnafu,
38    ProcedureStateReceiverSnafu, Result, TimestampMismatchSnafu, UnexpectedSnafu,
39    WaitProcedureSnafu,
40};
41use crate::key::TableMetadataManagerRef;
42use crate::metrics;
43use crate::node_manager::NodeManagerRef;
44use crate::reconciliation::reconcile_logical_tables::ReconcileLogicalTablesProcedure;
45use crate::reconciliation::reconcile_table::ReconcileTableProcedure;
46use crate::reconciliation::reconcile_table::resolve_column_metadata::ResolveStrategy;
47
48#[derive(Debug, PartialEq, Eq)]
49pub(crate) struct PartialRegionMetadata<'a> {
50    pub(crate) column_metadatas: &'a [ColumnMetadata],
51    pub(crate) primary_key: &'a [u32],
52    pub(crate) table_id: TableId,
53}
54
55impl<'a> From<&'a RegionMetadata> for PartialRegionMetadata<'a> {
56    fn from(region_metadata: &'a RegionMetadata) -> Self {
57        Self {
58            column_metadatas: &region_metadata.column_metadatas,
59            primary_key: &region_metadata.primary_key,
60            table_id: region_metadata.region_id.table_id(),
61        }
62    }
63}
64
65/// Checks if the column metadatas are consistent.
66///
67/// The column metadatas are consistent if:
68/// - The column metadatas are the same.
69/// - The primary key are the same.
70/// - The table id of the region metadatas are the same.
71///
72/// ## Panic
73/// Panic if region_metadatas is empty.
74pub(crate) fn check_column_metadatas_consistent(
75    region_metadatas: &[RegionMetadata],
76) -> Option<Vec<ColumnMetadata>> {
77    let is_column_metadata_consistent = region_metadatas
78        .windows(2)
79        .all(|w| PartialRegionMetadata::from(&w[0]) == PartialRegionMetadata::from(&w[1]));
80
81    if !is_column_metadata_consistent {
82        return None;
83    }
84
85    Some(region_metadatas[0].column_metadatas.clone())
86}
87
88/// Returns columns with tag values in primary-key order while preserving the
89/// positions of non-tag columns.
90pub(crate) fn reorder_tag_columns(
91    column_metadatas: &[ColumnMetadata],
92    primary_key: &[u32],
93) -> Result<Vec<ColumnMetadata>> {
94    let tag_count = column_metadatas
95        .iter()
96        .filter(|column| column.semantic_type == SemanticType::Tag)
97        .count();
98    ensure!(
99        primary_key.len() == tag_count,
100        UnexpectedSnafu {
101            err_msg: format!(
102                "Number of primary key columns {} does not match tag columns {}",
103                primary_key.len(),
104                tag_count,
105            ),
106        }
107    );
108
109    let mut primary_key_ids = HashSet::with_capacity(primary_key.len());
110    let mut tags = Vec::with_capacity(primary_key.len());
111    for column_id in primary_key {
112        let column = column_metadatas
113            .iter()
114            .find(|column| column.column_id == *column_id)
115            .with_context(|| UnexpectedSnafu {
116                err_msg: format!(
117                    "Primary key column {} not found in column metadata",
118                    column_id
119                ),
120            })?;
121        ensure!(
122            column.semantic_type == SemanticType::Tag,
123            UnexpectedSnafu {
124                err_msg: format!("Primary key column {} is not a tag", column_id),
125            }
126        );
127        ensure!(
128            primary_key_ids.insert(column_id),
129            UnexpectedSnafu {
130                err_msg: format!("Primary key column {} is duplicated", column_id),
131            }
132        );
133        tags.push(column);
134    }
135
136    let mut tags = tags.into_iter();
137
138    column_metadatas
139        .iter()
140        .map(|column| {
141            if column.semantic_type == SemanticType::Tag {
142                tags.next().cloned().with_context(|| UnexpectedSnafu {
143                    err_msg: "Primary key has fewer columns than tags".to_string(),
144                })
145            } else {
146                Ok(column.clone())
147            }
148        })
149        .collect()
150}
151
152/// Resolves column metadata inconsistencies among the given region metadatas
153/// by using the column metadata from the metasrv as the source of truth.
154///
155/// All region metadatas whose column metadata differs from the given `column_metadatas`
156/// will be marked for reconciliation.
157///
158/// Returns the region ids that need to be reconciled.
159pub(crate) fn resolve_column_metadatas_with_metasrv(
160    column_metadatas: &[ColumnMetadata],
161    region_metadatas: &[RegionMetadata],
162) -> Result<Vec<RegionId>> {
163    let is_same_table = region_metadatas
164        .windows(2)
165        .all(|w| w[0].region_id.table_id() == w[1].region_id.table_id());
166
167    ensure!(
168        is_same_table,
169        UnexpectedSnafu {
170            err_msg: "Region metadatas are not from the same table"
171        }
172    );
173
174    let mut regions_ids = vec![];
175    for region_metadata in region_metadatas {
176        if region_metadata.column_metadatas != column_metadatas {
177            check_column_metadata_invariants(column_metadatas, &region_metadata.column_metadatas)?;
178            regions_ids.push(region_metadata.region_id);
179        }
180    }
181    Ok(regions_ids)
182}
183
184/// Resolves column metadata inconsistencies among the given region metadatas
185/// by selecting the column metadata with the highest schema version.
186///
187/// This strategy assumes that at most two versions of column metadata may exist,
188/// due to the poison mechanism, making the highest schema version a safe choice.
189///
190/// Returns the resolved column metadata and the region ids that need to be reconciled.
191pub(crate) fn resolve_column_metadatas_with_latest(
192    region_metadatas: &[RegionMetadata],
193) -> Result<(Vec<ColumnMetadata>, Vec<RegionId>)> {
194    let is_same_table = region_metadatas
195        .windows(2)
196        .all(|w| w[0].region_id.table_id() == w[1].region_id.table_id());
197
198    ensure!(
199        is_same_table,
200        UnexpectedSnafu {
201            err_msg: "Region metadatas are not from the same table"
202        }
203    );
204
205    let latest_region_metadata = region_metadatas
206        .iter()
207        .max_by_key(|c| c.schema_version)
208        .context(UnexpectedSnafu {
209            err_msg: "All Region metadatas have the same schema version",
210        })?;
211    let latest_column_metadatas = PartialRegionMetadata::from(latest_region_metadata);
212
213    let mut region_ids = vec![];
214    for region_metadata in region_metadatas {
215        if PartialRegionMetadata::from(region_metadata) != latest_column_metadatas {
216            check_column_metadata_invariants(
217                &latest_region_metadata.column_metadatas,
218                &region_metadata.column_metadatas,
219            )?;
220            region_ids.push(region_metadata.region_id);
221        }
222    }
223
224    Ok((
225        reorder_tag_columns(
226            &latest_region_metadata.column_metadatas,
227            &latest_region_metadata.primary_key,
228        )?,
229        region_ids,
230    ))
231}
232
233/// Constructs a vector of [`ColumnMetadata`] from the provided table information.
234///
235/// This function maps each [`ColumnSchema`] to its corresponding [`ColumnMetadata`] by
236/// determining the semantic type (Tag, Timestamp, or Field) and retrieving the column ID
237/// from the `name_to_ids` mapping.
238///
239/// Returns an error if any column name is missing in the mapping.
240pub(crate) fn build_column_metadata_from_table_info(
241    column_schemas: &[ColumnSchema],
242    primary_key_indexes: &[usize],
243    name_to_ids: &HashMap<String, u32>,
244) -> Result<Vec<ColumnMetadata>> {
245    let primary_names = primary_key_indexes
246        .iter()
247        .map(|i| column_schemas[*i].name.as_str())
248        .collect::<HashSet<_>>();
249
250    column_schemas
251        .iter()
252        .map(|column_schema| {
253            let column_id = *name_to_ids
254                .get(column_schema.name.as_str())
255                .with_context(|| UnexpectedSnafu {
256                    err_msg: format!(
257                        "Column name {} not found in name_to_ids",
258                        column_schema.name
259                    ),
260                })?;
261
262            let semantic_type = if primary_names.contains(&column_schema.name.as_str()) {
263                SemanticType::Tag
264            } else if column_schema.is_time_index() {
265                SemanticType::Timestamp
266            } else {
267                SemanticType::Field
268            };
269            Ok(ColumnMetadata {
270                column_schema: column_schema.clone(),
271                semantic_type,
272                column_id,
273            })
274        })
275        .collect::<Result<Vec<_>>>()
276}
277
278/// Builds SyncColumns metadata with tags in the TableInfo primary-key order.
279/// Non-tag columns retain their schema positions.
280pub(crate) fn build_reconciliation_column_metadata(
281    column_schemas: &[ColumnSchema],
282    primary_key_indexes: &[usize],
283    name_to_ids: &HashMap<String, u32>,
284) -> Result<Vec<ColumnMetadata>> {
285    let column_metadatas =
286        build_column_metadata_from_table_info(column_schemas, primary_key_indexes, name_to_ids)?;
287    let primary_key = primary_key_indexes
288        .iter()
289        .map(|index| column_metadatas[*index].column_id)
290        .collect::<Vec<_>>();
291    reorder_tag_columns(&column_metadatas, &primary_key)
292}
293
294/// Checks whether the schema invariants hold between the existing and new column metadata.
295///
296/// Invariants:
297/// - Primary key (Tag) columns must exist in the new metadata, with identical name and ID.
298/// - Timestamp column must remain exactly the same in name and ID.
299pub(crate) fn check_column_metadata_invariants(
300    new_column_metadatas: &[ColumnMetadata],
301    column_metadatas: &[ColumnMetadata],
302) -> Result<()> {
303    let new_primary_keys = new_column_metadatas
304        .iter()
305        .filter(|c| c.semantic_type == SemanticType::Tag)
306        .map(|c| (c.column_schema.name.as_str(), c.column_id))
307        .collect::<HashMap<_, _>>();
308
309    let old_primary_keys = column_metadatas
310        .iter()
311        .filter(|c| c.semantic_type == SemanticType::Tag)
312        .map(|c| (c.column_schema.name.as_str(), c.column_id));
313
314    for (name, id) in old_primary_keys {
315        let column_id = new_primary_keys
316            .get(name)
317            .cloned()
318            .context(ColumnNotFoundSnafu {
319                column_name: name,
320                column_id: id,
321            })?;
322
323        ensure!(
324            column_id == id,
325            ColumnIdMismatchSnafu {
326                column_name: name,
327                expected_column_id: id,
328                actual_column_id: column_id,
329            }
330        );
331    }
332
333    let new_ts_column = new_column_metadatas
334        .iter()
335        .find(|c| c.semantic_type == SemanticType::Timestamp)
336        .map(|c| (c.column_schema.name.as_str(), c.column_id))
337        .context(UnexpectedSnafu {
338            err_msg: "Timestamp column not found in new column metadata",
339        })?;
340
341    let old_ts_column = column_metadatas
342        .iter()
343        .find(|c| c.semantic_type == SemanticType::Timestamp)
344        .map(|c| (c.column_schema.name.as_str(), c.column_id))
345        .context(UnexpectedSnafu {
346            err_msg: "Timestamp column not found in column metadata",
347        })?;
348    ensure!(
349        new_ts_column == old_ts_column,
350        TimestampMismatchSnafu {
351            expected_column_name: old_ts_column.0,
352            expected_column_id: old_ts_column.1,
353            actual_column_name: new_ts_column.0,
354            actual_column_id: new_ts_column.1,
355        }
356    );
357
358    Ok(())
359}
360
361/// Builds a [`TableMeta`] from the provided [`ColumnMetadata`]s.
362///
363/// Returns an error if:
364/// - Any column is missing in the `name_to_ids`(if `name_to_ids` is provided).
365/// - The column id in table metadata is not the same as the column id in the column metadata.(if `name_to_ids` is provided)
366/// - The table index is missing in the column metadata.
367/// - The primary key or partition key columns are missing in the column metadata.
368///
369/// TODO(weny): add tests
370pub(crate) fn build_table_meta_from_column_metadatas(
371    table_id: TableId,
372    table_ref: TableReference,
373    table_meta: &TableMeta,
374    name_to_ids: Option<HashMap<String, u32>>,
375    column_metadata: &[ColumnMetadata],
376) -> Result<TableMeta> {
377    let column_in_column_metadata = column_metadata
378        .iter()
379        .map(|c| (c.column_schema.name.as_str(), c))
380        .collect::<HashMap<_, _>>();
381    let column_schemas = table_meta.schema.column_schemas();
382    let primary_key_names = table_meta
383        .primary_key_indices
384        .iter()
385        .map(|i| column_schemas[*i].name.as_str())
386        .collect::<HashSet<_>>();
387    let partition_key_names = table_meta
388        .partition_key_indices
389        .iter()
390        .map(|i| column_schemas[*i].name.as_str())
391        .collect::<HashSet<_>>();
392    ensure!(
393        column_metadata
394            .iter()
395            .any(|c| c.semantic_type == SemanticType::Timestamp),
396        UnexpectedSnafu {
397            err_msg: format!(
398                "Missing table index in column metadata, table: {}, table_id: {}",
399                table_ref, table_id
400            ),
401        }
402    );
403
404    if let Some(name_to_ids) = &name_to_ids {
405        // Ensures all primary key and partition key exists in the column metadata.
406        for column_name in primary_key_names.iter().chain(partition_key_names.iter()) {
407            let column_in_column_metadata = column_in_column_metadata
408                .get(column_name)
409                .with_context(|| MissingColumnInColumnMetadataSnafu {
410                    column_name: column_name.to_string(),
411                    table_name: table_ref.to_string(),
412                    table_id,
413                })?;
414
415            let column_id = *name_to_ids
416                .get(*column_name)
417                .with_context(|| UnexpectedSnafu {
418                    err_msg: format!("column id not found in name_to_ids: {}", column_name),
419                })?;
420            ensure!(
421                column_id == column_in_column_metadata.column_id,
422                MismatchColumnIdSnafu {
423                    column_name: column_name.to_string(),
424                    column_id,
425                    table_name: table_ref.to_string(),
426                    table_id,
427                }
428            );
429        }
430    } else {
431        warn!(
432            "`name_to_ids` is not provided, table: {}, table_id: {}",
433            table_ref, table_id
434        );
435    }
436
437    let mut new_raw_table_meta = table_meta.clone();
438    let primary_key_indices = &mut new_raw_table_meta.primary_key_indices;
439    let partition_key_indices = &mut new_raw_table_meta.partition_key_indices;
440    let value_indices = &mut new_raw_table_meta.value_indices;
441    let mut columns = Vec::with_capacity(column_metadata.len());
442    let column_ids = &mut new_raw_table_meta.column_ids;
443    let next_column_id = &mut new_raw_table_meta.next_column_id;
444
445    column_ids.clear();
446    value_indices.clear();
447    primary_key_indices.clear();
448    partition_key_indices.clear();
449
450    for (idx, col) in column_metadata.iter().enumerate() {
451        if partition_key_names.contains(&col.column_schema.name.as_str()) {
452            partition_key_indices.push(idx);
453        }
454        match col.semantic_type {
455            SemanticType::Tag => {
456                primary_key_indices.push(idx);
457            }
458            SemanticType::Field => {
459                value_indices.push(idx);
460            }
461            SemanticType::Timestamp => {
462                value_indices.push(idx);
463            }
464        }
465
466        columns.push(col.column_schema.clone());
467        column_ids.push(col.column_id);
468    }
469
470    *next_column_id = column_ids
471        .iter()
472        .filter(|id| !ReservedColumnId::is_reserved(**id))
473        .max()
474        .map(|max| max + 1)
475        .unwrap_or(*next_column_id)
476        .max(*next_column_id);
477
478    new_raw_table_meta.schema = Arc::new(Schema::new_with_version(
479        columns,
480        table_meta.schema.version(),
481    ));
482    Ok(new_raw_table_meta)
483}
484
485/// Returns true if the logical table info needs to be updated.
486///
487/// The logical table only support to add columns, so we can check the length of column metadatas
488/// to determine whether the logical table info needs to be updated.
489pub(crate) fn need_update_logical_table_info(
490    table_info: &TableInfo,
491    column_metadatas: &[ColumnMetadata],
492) -> bool {
493    table_info.meta.schema.column_schemas().len() != column_metadatas.len()
494}
495
496/// The result of waiting for inflight subprocedures.
497pub struct PartialSuccessResult<'a> {
498    pub failed_procedures: Vec<&'a SubprocedureMeta>,
499    pub success_procedures: Vec<&'a SubprocedureMeta>,
500}
501
502/// The result of waiting for inflight subprocedures.
503pub enum WaitForInflightSubproceduresResult<'a> {
504    Success(Vec<&'a SubprocedureMeta>),
505    PartialSuccess(PartialSuccessResult<'a>),
506}
507
508/// Wait for inflight subprocedures.
509///
510/// If `fail_fast` is true, the function will return an error if any subprocedure fails.
511/// Otherwise, the function will continue waiting for all subprocedures to complete.
512pub(crate) async fn wait_for_inflight_subprocedures<'a>(
513    procedure_ctx: &ProcedureContext,
514    subprocedures: &'a [SubprocedureMeta],
515    fail_fast: bool,
516) -> Result<WaitForInflightSubproceduresResult<'a>> {
517    let mut receivers = Vec::with_capacity(subprocedures.len());
518    for subprocedure in subprocedures {
519        let procedure_id = subprocedure.procedure_id();
520        let receiver = procedure_ctx
521            .provider
522            .procedure_state_receiver(procedure_id)
523            .await
524            .context(ProcedureStateReceiverSnafu { procedure_id })?
525            .context(ProcedureStateReceiverNotFoundSnafu { procedure_id })?;
526        receivers.push((receiver, subprocedure));
527    }
528
529    let mut tasks = Vec::with_capacity(receivers.len());
530    for (receiver, subprocedure) in receivers.iter_mut() {
531        tasks.push(async move {
532            watcher::wait(receiver).await.inspect_err(|e| {
533                error!(e; "inflight subprocedure failed, parent procedure_id: {}, procedure: {}", procedure_ctx.procedure_id, subprocedure);
534            })
535        });
536    }
537
538    if fail_fast {
539        try_join_all(tasks).await.context(WaitProcedureSnafu)?;
540        return Ok(WaitForInflightSubproceduresResult::Success(
541            subprocedures.iter().collect(),
542        ));
543    }
544
545    // If fail_fast is false, we need to wait for all subprocedures to complete.
546    let results = join_all(tasks).await;
547    let failed_procedures_num = results.iter().filter(|r| r.is_err()).count();
548    if failed_procedures_num == 0 {
549        return Ok(WaitForInflightSubproceduresResult::Success(
550            subprocedures.iter().collect(),
551        ));
552    }
553    warn!(
554        "{} inflight subprocedures failed, total: {}, parent procedure_id: {}",
555        failed_procedures_num,
556        subprocedures.len(),
557        procedure_ctx.procedure_id
558    );
559
560    let mut failed_procedures = Vec::with_capacity(failed_procedures_num);
561    let mut success_procedures = Vec::with_capacity(subprocedures.len() - failed_procedures_num);
562    for (result, subprocedure) in results.into_iter().zip(subprocedures) {
563        if result.is_err() {
564            failed_procedures.push(subprocedure);
565        } else {
566            success_procedures.push(subprocedure);
567        }
568    }
569
570    Ok(WaitForInflightSubproceduresResult::PartialSuccess(
571        PartialSuccessResult {
572            failed_procedures,
573            success_procedures,
574        },
575    ))
576}
577
578#[derive(Clone)]
579pub struct Context {
580    pub node_manager: NodeManagerRef,
581    pub table_metadata_manager: TableMetadataManagerRef,
582    pub cache_invalidator: CacheInvalidatorRef,
583}
584
585/// Metadata for an inflight physical table subprocedure.
586pub struct PhysicalTableMeta {
587    pub procedure_id: ProcedureId,
588    pub table_id: TableId,
589    pub table_name: TableName,
590}
591
592/// Metadata for an inflight logical table subprocedure.
593pub struct LogicalTableMeta {
594    pub procedure_id: ProcedureId,
595    pub physical_table_id: TableId,
596    pub physical_table_name: TableName,
597    pub logical_tables: Vec<(TableId, TableName)>,
598}
599
600/// Metadata for an inflight database subprocedure.
601pub struct ReconcileDatabaseMeta {
602    pub procedure_id: ProcedureId,
603    pub catalog: String,
604    pub schema: String,
605}
606
607/// The inflight subprocedure metadata.
608pub enum SubprocedureMeta {
609    PhysicalTable(PhysicalTableMeta),
610    LogicalTable(LogicalTableMeta),
611    Database(ReconcileDatabaseMeta),
612}
613
614impl Display for SubprocedureMeta {
615    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
616        match self {
617            SubprocedureMeta::PhysicalTable(meta) => {
618                write!(
619                    f,
620                    "ReconcilePhysicalTable(procedure_id: {}, table_id: {}, table_name: {})",
621                    meta.procedure_id, meta.table_id, meta.table_name
622                )
623            }
624            SubprocedureMeta::LogicalTable(meta) => {
625                write!(
626                    f,
627                    "ReconcileLogicalTable(procedure_id: {}, physical_table_id: {}, physical_table_name: {}, logical_tables: {:?})",
628                    meta.procedure_id,
629                    meta.physical_table_id,
630                    meta.physical_table_name,
631                    meta.logical_tables
632                )
633            }
634            SubprocedureMeta::Database(meta) => {
635                write!(
636                    f,
637                    "ReconcileDatabase(procedure_id: {}, catalog: {}, schema: {})",
638                    meta.procedure_id, meta.catalog, meta.schema
639                )
640            }
641        }
642    }
643}
644
645impl SubprocedureMeta {
646    /// Creates a new logical table subprocedure metadata.
647    pub fn new_logical_table(
648        procedure_id: ProcedureId,
649        physical_table_id: TableId,
650        physical_table_name: TableName,
651        logical_tables: Vec<(TableId, TableName)>,
652    ) -> Self {
653        Self::LogicalTable(LogicalTableMeta {
654            procedure_id,
655            physical_table_id,
656            physical_table_name,
657            logical_tables,
658        })
659    }
660
661    /// Creates a new physical table subprocedure metadata.
662    pub fn new_physical_table(
663        procedure_id: ProcedureId,
664        table_id: TableId,
665        table_name: TableName,
666    ) -> Self {
667        Self::PhysicalTable(PhysicalTableMeta {
668            procedure_id,
669            table_id,
670            table_name,
671        })
672    }
673
674    /// Creates a new reconcile database subprocedure metadata.
675    pub fn new_reconcile_database(
676        procedure_id: ProcedureId,
677        catalog: String,
678        schema: String,
679    ) -> Self {
680        Self::Database(ReconcileDatabaseMeta {
681            procedure_id,
682            catalog,
683            schema,
684        })
685    }
686
687    /// Returns the procedure id of the subprocedure.
688    pub fn procedure_id(&self) -> ProcedureId {
689        match self {
690            SubprocedureMeta::PhysicalTable(meta) => meta.procedure_id,
691            SubprocedureMeta::LogicalTable(meta) => meta.procedure_id,
692            SubprocedureMeta::Database(meta) => meta.procedure_id,
693        }
694    }
695
696    /// Returns the number of tables will be reconciled.
697    pub fn table_num(&self) -> usize {
698        match self {
699            SubprocedureMeta::PhysicalTable(_) => 1,
700            SubprocedureMeta::LogicalTable(meta) => meta.logical_tables.len(),
701            SubprocedureMeta::Database(_) => 0,
702        }
703    }
704
705    /// Returns the number of databases will be reconciled.
706    pub fn database_num(&self) -> usize {
707        match self {
708            SubprocedureMeta::Database(_) => 1,
709            _ => 0,
710        }
711    }
712}
713
714/// The metrics of reconciling catalog.
715#[derive(Clone, Default)]
716pub struct ReconcileCatalogMetrics {
717    pub succeeded_databases: usize,
718    pub failed_databases: usize,
719}
720
721impl AddAssign for ReconcileCatalogMetrics {
722    fn add_assign(&mut self, other: Self) {
723        self.succeeded_databases += other.succeeded_databases;
724        self.failed_databases += other.failed_databases;
725    }
726}
727
728impl Display for ReconcileCatalogMetrics {
729    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
730        write!(
731            f,
732            "succeeded_databases: {}, failed_databases: {}",
733            self.succeeded_databases, self.failed_databases
734        )
735    }
736}
737
738impl From<WaitForInflightSubproceduresResult<'_>> for ReconcileCatalogMetrics {
739    fn from(result: WaitForInflightSubproceduresResult<'_>) -> Self {
740        match result {
741            WaitForInflightSubproceduresResult::Success(subprocedures) => ReconcileCatalogMetrics {
742                succeeded_databases: subprocedures.len(),
743                failed_databases: 0,
744            },
745            WaitForInflightSubproceduresResult::PartialSuccess(PartialSuccessResult {
746                failed_procedures,
747                success_procedures,
748            }) => {
749                let succeeded_databases = success_procedures
750                    .iter()
751                    .map(|subprocedure| subprocedure.database_num())
752                    .sum();
753                let failed_databases = failed_procedures
754                    .iter()
755                    .map(|subprocedure| subprocedure.database_num())
756                    .sum();
757                ReconcileCatalogMetrics {
758                    succeeded_databases,
759                    failed_databases,
760                }
761            }
762        }
763    }
764}
765
766/// The metrics of reconciling database.
767#[derive(Clone, Default)]
768pub struct ReconcileDatabaseMetrics {
769    pub succeeded_tables: usize,
770    pub failed_tables: usize,
771    pub succeeded_procedures: usize,
772    pub failed_procedures: usize,
773}
774
775impl Display for ReconcileDatabaseMetrics {
776    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
777        write!(
778            f,
779            "succeeded_tables: {}, failed_tables: {}, succeeded_procedures: {}, failed_procedures: {}",
780            self.succeeded_tables,
781            self.failed_tables,
782            self.succeeded_procedures,
783            self.failed_procedures
784        )
785    }
786}
787
788impl AddAssign for ReconcileDatabaseMetrics {
789    fn add_assign(&mut self, other: Self) {
790        self.succeeded_tables += other.succeeded_tables;
791        self.failed_tables += other.failed_tables;
792        self.succeeded_procedures += other.succeeded_procedures;
793        self.failed_procedures += other.failed_procedures;
794    }
795}
796
797impl From<WaitForInflightSubproceduresResult<'_>> for ReconcileDatabaseMetrics {
798    fn from(result: WaitForInflightSubproceduresResult<'_>) -> Self {
799        match result {
800            WaitForInflightSubproceduresResult::Success(subprocedures) => {
801                let table_num = subprocedures
802                    .iter()
803                    .map(|subprocedure| subprocedure.table_num())
804                    .sum();
805                ReconcileDatabaseMetrics {
806                    succeeded_procedures: subprocedures.len(),
807                    failed_procedures: 0,
808                    succeeded_tables: table_num,
809                    failed_tables: 0,
810                }
811            }
812            WaitForInflightSubproceduresResult::PartialSuccess(PartialSuccessResult {
813                failed_procedures,
814                success_procedures,
815            }) => {
816                let succeeded_tables = success_procedures
817                    .iter()
818                    .map(|subprocedure| subprocedure.table_num())
819                    .sum();
820                let failed_tables = failed_procedures
821                    .iter()
822                    .map(|subprocedure| subprocedure.table_num())
823                    .sum();
824                ReconcileDatabaseMetrics {
825                    succeeded_procedures: success_procedures.len(),
826                    failed_procedures: failed_procedures.len(),
827                    succeeded_tables,
828                    failed_tables,
829                }
830            }
831        }
832    }
833}
834
835/// The metrics of reconciling logical tables.
836#[derive(Clone)]
837pub struct ReconcileLogicalTableMetrics {
838    pub start_time: Instant,
839    pub update_table_info_count: usize,
840    pub create_tables_count: usize,
841    pub column_metadata_consistent_count: usize,
842    pub column_metadata_inconsistent_count: usize,
843}
844
845impl Default for ReconcileLogicalTableMetrics {
846    fn default() -> Self {
847        Self {
848            start_time: Instant::now(),
849            update_table_info_count: 0,
850            create_tables_count: 0,
851            column_metadata_consistent_count: 0,
852            column_metadata_inconsistent_count: 0,
853        }
854    }
855}
856
857const CREATE_TABLES: &str = "create_tables";
858const UPDATE_TABLE_INFO: &str = "update_table_info";
859const COLUMN_METADATA_CONSISTENT: &str = "column_metadata_consistent";
860const COLUMN_METADATA_INCONSISTENT: &str = "column_metadata_inconsistent";
861
862impl ReconcileLogicalTableMetrics {
863    /// The total number of tables that have been reconciled.
864    pub fn total_table_count(&self) -> usize {
865        self.create_tables_count
866            + self.column_metadata_consistent_count
867            + self.column_metadata_inconsistent_count
868    }
869}
870
871impl Drop for ReconcileLogicalTableMetrics {
872    fn drop(&mut self) {
873        let procedure_name = ReconcileLogicalTablesProcedure::TYPE_NAME;
874        metrics::METRIC_META_RECONCILIATION_STATS
875            .with_label_values(&[procedure_name, metrics::TABLE_TYPE_LOGICAL, CREATE_TABLES])
876            .inc_by(self.create_tables_count as u64);
877        metrics::METRIC_META_RECONCILIATION_STATS
878            .with_label_values(&[
879                procedure_name,
880                metrics::TABLE_TYPE_LOGICAL,
881                UPDATE_TABLE_INFO,
882            ])
883            .inc_by(self.update_table_info_count as u64);
884        metrics::METRIC_META_RECONCILIATION_STATS
885            .with_label_values(&[
886                procedure_name,
887                metrics::TABLE_TYPE_LOGICAL,
888                COLUMN_METADATA_CONSISTENT,
889            ])
890            .inc_by(self.column_metadata_consistent_count as u64);
891        metrics::METRIC_META_RECONCILIATION_STATS
892            .with_label_values(&[
893                procedure_name,
894                metrics::TABLE_TYPE_LOGICAL,
895                COLUMN_METADATA_INCONSISTENT,
896            ])
897            .inc_by(self.column_metadata_inconsistent_count as u64);
898    }
899}
900
901impl Display for ReconcileLogicalTableMetrics {
902    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
903        let elapsed = self.start_time.elapsed();
904        if self.create_tables_count > 0 {
905            write!(f, "create_tables_count: {}, ", self.create_tables_count)?;
906        }
907        if self.update_table_info_count > 0 {
908            write!(
909                f,
910                "update_table_info_count: {}, ",
911                self.update_table_info_count
912            )?;
913        }
914        if self.column_metadata_consistent_count > 0 {
915            write!(
916                f,
917                "column_metadata_consistent_count: {}, ",
918                self.column_metadata_consistent_count
919            )?;
920        }
921        if self.column_metadata_inconsistent_count > 0 {
922            write!(
923                f,
924                "column_metadata_inconsistent_count: {}, ",
925                self.column_metadata_inconsistent_count
926            )?;
927        }
928
929        write!(
930            f,
931            "total_table_count: {}, elapsed: {:?}",
932            self.total_table_count(),
933            elapsed
934        )
935    }
936}
937
938/// The result of resolving column metadata.
939#[derive(Clone, Copy)]
940pub enum ResolveColumnMetadataResult {
941    Consistent,
942    Inconsistent(ResolveStrategy),
943}
944
945impl Display for ResolveColumnMetadataResult {
946    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
947        match self {
948            ResolveColumnMetadataResult::Consistent => write!(f, "Consistent"),
949            ResolveColumnMetadataResult::Inconsistent(strategy) => {
950                let strategy_str = strategy.as_ref();
951                write!(f, "Inconsistent({})", strategy_str)
952            }
953        }
954    }
955}
956
957/// The metrics of reconciling physical tables.
958#[derive(Clone)]
959pub struct ReconcileTableMetrics {
960    /// The start time of the reconciliation.
961    pub start_time: Instant,
962    /// The result of resolving column metadata.
963    pub resolve_column_metadata_result: Option<ResolveColumnMetadataResult>,
964    /// Whether the table info has been updated.
965    pub update_table_info: bool,
966}
967
968impl Drop for ReconcileTableMetrics {
969    fn drop(&mut self) {
970        if let Some(resolve_column_metadata_result) = self.resolve_column_metadata_result {
971            match resolve_column_metadata_result {
972                ResolveColumnMetadataResult::Consistent => {
973                    metrics::METRIC_META_RECONCILIATION_STATS
974                        .with_label_values(&[
975                            ReconcileTableProcedure::TYPE_NAME,
976                            metrics::TABLE_TYPE_PHYSICAL,
977                            COLUMN_METADATA_CONSISTENT,
978                        ])
979                        .inc();
980                }
981                ResolveColumnMetadataResult::Inconsistent(strategy) => {
982                    metrics::METRIC_META_RECONCILIATION_STATS
983                        .with_label_values(&[
984                            ReconcileTableProcedure::TYPE_NAME,
985                            metrics::TABLE_TYPE_PHYSICAL,
986                            COLUMN_METADATA_INCONSISTENT,
987                        ])
988                        .inc();
989                    metrics::METRIC_META_RECONCILIATION_RESOLVED_COLUMN_METADATA
990                        .with_label_values(&[strategy.as_ref()])
991                        .inc();
992                }
993            }
994        }
995        if self.update_table_info {
996            metrics::METRIC_META_RECONCILIATION_STATS
997                .with_label_values(&[
998                    ReconcileTableProcedure::TYPE_NAME,
999                    metrics::TABLE_TYPE_PHYSICAL,
1000                    UPDATE_TABLE_INFO,
1001                ])
1002                .inc();
1003        }
1004    }
1005}
1006
1007impl Default for ReconcileTableMetrics {
1008    fn default() -> Self {
1009        Self {
1010            start_time: Instant::now(),
1011            resolve_column_metadata_result: None,
1012            update_table_info: false,
1013        }
1014    }
1015}
1016
1017impl Display for ReconcileTableMetrics {
1018    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1019        let elapsed = self.start_time.elapsed();
1020        if let Some(resolve_column_metadata_result) = self.resolve_column_metadata_result {
1021            write!(
1022                f,
1023                "resolve_column_metadata_result: {}, ",
1024                resolve_column_metadata_result
1025            )?;
1026        }
1027        write!(
1028            f,
1029            "update_table_info: {}, elapsed: {:?}",
1030            self.update_table_info, elapsed
1031        )
1032    }
1033}
1034
1035#[cfg(test)]
1036mod tests {
1037    use std::assert_matches;
1038    use std::collections::HashMap;
1039    use std::sync::Arc;
1040
1041    use api::v1::SemanticType;
1042    use datatypes::prelude::ConcreteDataType;
1043    use datatypes::schema::{ColumnSchema, Schema, SchemaBuilder};
1044    use store_api::metadata::ColumnMetadata;
1045    use store_api::storage::RegionId;
1046    use table::metadata::TableMetaBuilder;
1047    use table::table_reference::TableReference;
1048
1049    use super::*;
1050    use crate::ddl::test_util::region_metadata::build_region_metadata;
1051    use crate::error::Error;
1052    use crate::reconciliation::utils::check_column_metadatas_consistent;
1053
1054    fn new_test_schema() -> Schema {
1055        let column_schemas = vec![
1056            ColumnSchema::new("col1", ConcreteDataType::int32_datatype(), true),
1057            ColumnSchema::new(
1058                "ts",
1059                ConcreteDataType::timestamp_millisecond_datatype(),
1060                false,
1061            )
1062            .with_time_index(true),
1063            ColumnSchema::new("col2", ConcreteDataType::int32_datatype(), true),
1064        ];
1065        SchemaBuilder::try_from(column_schemas)
1066            .unwrap()
1067            .version(123)
1068            .build()
1069            .unwrap()
1070    }
1071
1072    fn new_test_column_metadatas() -> Vec<ColumnMetadata> {
1073        vec![
1074            ColumnMetadata {
1075                column_schema: ColumnSchema::new("col1", ConcreteDataType::int32_datatype(), true),
1076                semantic_type: SemanticType::Tag,
1077                column_id: 0,
1078            },
1079            ColumnMetadata {
1080                column_schema: ColumnSchema::new(
1081                    "ts",
1082                    ConcreteDataType::timestamp_millisecond_datatype(),
1083                    false,
1084                )
1085                .with_time_index(true),
1086                semantic_type: SemanticType::Timestamp,
1087                column_id: 1,
1088            },
1089            ColumnMetadata {
1090                column_schema: ColumnSchema::new("col2", ConcreteDataType::int32_datatype(), true),
1091                semantic_type: SemanticType::Field,
1092                column_id: 2,
1093            },
1094        ]
1095    }
1096
1097    fn new_test_raw_table_info() -> TableMeta {
1098        let mut table_meta_builder = TableMetaBuilder::empty();
1099        table_meta_builder
1100            .schema(Arc::new(new_test_schema()))
1101            .primary_key_indices(vec![0])
1102            .partition_key_indices(vec![2])
1103            .next_column_id(4)
1104            .build()
1105            .unwrap()
1106    }
1107
1108    #[test]
1109    fn test_build_table_info_from_column_metadatas_identical() {
1110        let column_metadatas = new_test_column_metadatas();
1111        let table_id = 1;
1112        let table_ref = TableReference::full("test_catalog", "test_schema", "test_table");
1113        let mut table_meta = new_test_raw_table_info();
1114        table_meta.column_ids = vec![0, 1, 2];
1115        let name_to_ids = HashMap::from([
1116            ("col1".to_string(), 0),
1117            ("ts".to_string(), 1),
1118            ("col2".to_string(), 2),
1119        ]);
1120
1121        let new_table_meta = build_table_meta_from_column_metadatas(
1122            table_id,
1123            table_ref,
1124            &table_meta,
1125            Some(name_to_ids),
1126            &column_metadatas,
1127        )
1128        .unwrap();
1129        assert_eq!(new_table_meta, table_meta);
1130    }
1131
1132    #[test]
1133    fn test_build_table_info_from_column_metadatas() {
1134        let mut column_metadatas = new_test_column_metadatas();
1135        column_metadatas.push(ColumnMetadata {
1136            column_schema: ColumnSchema::new(
1137                "__table_id",
1138                ConcreteDataType::string_datatype(),
1139                true,
1140            ),
1141            semantic_type: SemanticType::Tag,
1142            column_id: ReservedColumnId::table_id(),
1143        });
1144
1145        let table_id = 1;
1146        let table_ref = TableReference::full("test_catalog", "test_schema", "test_table");
1147        let table_meta = new_test_raw_table_info();
1148        let name_to_ids = HashMap::from([
1149            ("col1".to_string(), 0),
1150            ("ts".to_string(), 1),
1151            ("col2".to_string(), 2),
1152        ]);
1153
1154        let new_table_meta = build_table_meta_from_column_metadatas(
1155            table_id,
1156            table_ref,
1157            &table_meta,
1158            Some(name_to_ids),
1159            &column_metadatas,
1160        )
1161        .unwrap();
1162
1163        assert_eq!(new_table_meta.primary_key_indices, vec![0, 3]);
1164        assert_eq!(new_table_meta.partition_key_indices, vec![2]);
1165        assert_eq!(new_table_meta.value_indices, vec![1, 2]);
1166        assert_eq!(new_table_meta.schema.timestamp_index(), Some(1));
1167        assert_eq!(
1168            new_table_meta.column_ids,
1169            vec![0, 1, 2, ReservedColumnId::table_id()]
1170        );
1171        assert_eq!(new_table_meta.next_column_id, table_meta.next_column_id);
1172    }
1173
1174    #[test]
1175    fn test_build_table_info_from_column_metadatas_with_incorrect_name_to_ids() {
1176        let column_metadatas = new_test_column_metadatas();
1177        let table_id = 1;
1178        let table_ref = TableReference::full("test_catalog", "test_schema", "test_table");
1179        let table_meta = new_test_raw_table_info();
1180        let name_to_ids = HashMap::from([
1181            ("col1".to_string(), 0),
1182            ("ts".to_string(), 1),
1183            // Change column id of col2 to 3.
1184            ("col2".to_string(), 3),
1185        ]);
1186
1187        let err = build_table_meta_from_column_metadatas(
1188            table_id,
1189            table_ref,
1190            &table_meta,
1191            Some(name_to_ids),
1192            &column_metadatas,
1193        )
1194        .unwrap_err();
1195
1196        assert_matches!(err, Error::MismatchColumnId { .. });
1197    }
1198
1199    #[test]
1200    fn test_build_table_info_from_column_metadatas_with_missing_time_index() {
1201        let mut column_metadatas = new_test_column_metadatas();
1202        column_metadatas.retain(|c| c.semantic_type != SemanticType::Timestamp);
1203        let table_id = 1;
1204        let table_ref = TableReference::full("test_catalog", "test_schema", "test_table");
1205        let table_meta = new_test_raw_table_info();
1206        let name_to_ids = HashMap::from([
1207            ("col1".to_string(), 0),
1208            ("ts".to_string(), 1),
1209            ("col2".to_string(), 2),
1210        ]);
1211
1212        let err = build_table_meta_from_column_metadatas(
1213            table_id,
1214            table_ref,
1215            &table_meta,
1216            Some(name_to_ids),
1217            &column_metadatas,
1218        )
1219        .unwrap_err();
1220
1221        assert!(
1222            err.to_string()
1223                .contains("Missing table index in column metadata"),
1224            "err: {}",
1225            err
1226        );
1227    }
1228
1229    #[test]
1230    fn test_build_table_info_from_column_metadatas_with_missing_column() {
1231        let mut column_metadatas = new_test_column_metadatas();
1232        // Remove primary key column.
1233        column_metadatas.retain(|c| c.column_id != 0);
1234        let table_id = 1;
1235        let table_ref = TableReference::full("test_catalog", "test_schema", "test_table");
1236        let table_meta = new_test_raw_table_info();
1237        let name_to_ids = HashMap::from([
1238            ("col1".to_string(), 0),
1239            ("ts".to_string(), 1),
1240            ("col2".to_string(), 2),
1241        ]);
1242
1243        let err = build_table_meta_from_column_metadatas(
1244            table_id,
1245            table_ref,
1246            &table_meta,
1247            Some(name_to_ids.clone()),
1248            &column_metadatas,
1249        )
1250        .unwrap_err();
1251        assert_matches!(err, Error::MissingColumnInColumnMetadata { .. });
1252
1253        let mut column_metadatas = new_test_column_metadatas();
1254        // Remove partition key column.
1255        column_metadatas.retain(|c| c.column_id != 2);
1256
1257        let err = build_table_meta_from_column_metadatas(
1258            table_id,
1259            table_ref,
1260            &table_meta,
1261            Some(name_to_ids),
1262            &column_metadatas,
1263        )
1264        .unwrap_err();
1265        assert_matches!(err, Error::MissingColumnInColumnMetadata { .. });
1266    }
1267
1268    #[test]
1269    fn test_reconciliation_columns_follow_primary_key_order() {
1270        let columns = vec![
1271            ColumnMetadata {
1272                column_schema: ColumnSchema::new(
1273                    "tag_b",
1274                    ConcreteDataType::string_datatype(),
1275                    true,
1276                ),
1277                semantic_type: SemanticType::Tag,
1278                column_id: 1,
1279            },
1280            ColumnMetadata {
1281                column_schema: ColumnSchema::new("field", ConcreteDataType::int32_datatype(), true),
1282                semantic_type: SemanticType::Field,
1283                column_id: 2,
1284            },
1285            ColumnMetadata {
1286                column_schema: ColumnSchema::new(
1287                    "tag_a",
1288                    ConcreteDataType::string_datatype(),
1289                    true,
1290                ),
1291                semantic_type: SemanticType::Tag,
1292                column_id: 3,
1293            },
1294            ColumnMetadata {
1295                column_schema: ColumnSchema::new(
1296                    "ts",
1297                    ConcreteDataType::timestamp_millisecond_datatype(),
1298                    false,
1299                )
1300                .with_time_index(true),
1301                semantic_type: SemanticType::Timestamp,
1302                column_id: 4,
1303            },
1304        ];
1305        let schemas = columns
1306            .iter()
1307            .map(|column| column.column_schema.clone())
1308            .collect::<Vec<_>>();
1309        let ids = columns
1310            .iter()
1311            .map(|column| (column.column_schema.name.clone(), column.column_id))
1312            .collect();
1313        let expected = vec![3, 2, 1, 4];
1314        assert_eq!(
1315            build_reconciliation_column_metadata(&schemas, &[2, 0], &ids)
1316                .unwrap()
1317                .iter()
1318                .map(|column| column.column_id)
1319                .collect::<Vec<_>>(),
1320            expected
1321        );
1322
1323        let mut metadata = build_region_metadata(RegionId::new(1024, 0), &columns);
1324        metadata.primary_key = vec![3, 1];
1325        metadata.schema_version = 2;
1326
1327        assert_eq!(
1328            check_column_metadatas_consistent(&[metadata.clone()]).unwrap(),
1329            columns
1330        );
1331
1332        assert_eq!(
1333            resolve_column_metadatas_with_latest(&[metadata])
1334                .unwrap()
1335                .0
1336                .iter()
1337                .map(|column| column.column_id)
1338                .collect::<Vec<_>>(),
1339            expected
1340        );
1341    }
1342
1343    #[test]
1344    fn test_reorder_tag_columns_rejects_invalid_primary_key() {
1345        let columns = new_test_column_metadatas();
1346
1347        let err = reorder_tag_columns(&columns, &[999]).unwrap_err();
1348        assert_matches!(err, Error::Unexpected { .. });
1349        assert!(err.to_string().contains("not found in column metadata"));
1350
1351        let err = reorder_tag_columns(&columns, &[]).unwrap_err();
1352        assert_matches!(err, Error::Unexpected { .. });
1353        assert!(err.to_string().contains("does not match tag columns"));
1354
1355        let err = reorder_tag_columns(&columns, &[2]).unwrap_err();
1356        assert_matches!(err, Error::Unexpected { .. });
1357        assert!(err.to_string().contains("is not a tag"));
1358
1359        let mut duplicate_tags = columns.clone();
1360        duplicate_tags[2].semantic_type = SemanticType::Tag;
1361        let err = reorder_tag_columns(&duplicate_tags, &[0, 0]).unwrap_err();
1362        assert_matches!(err, Error::Unexpected { .. });
1363        assert!(err.to_string().contains("is duplicated"));
1364    }
1365
1366    #[test]
1367    fn test_check_column_metadatas_consistent() {
1368        let column_metadatas = new_test_column_metadatas();
1369        let region_metadata1 = build_region_metadata(RegionId::new(1024, 0), &column_metadatas);
1370        let region_metadata2 = build_region_metadata(RegionId::new(1024, 1), &column_metadatas);
1371        let result =
1372            check_column_metadatas_consistent(&[region_metadata1, region_metadata2]).unwrap();
1373        assert_eq!(result, column_metadatas);
1374
1375        let region_metadata1 = build_region_metadata(RegionId::new(1025, 0), &column_metadatas);
1376        let region_metadata2 = build_region_metadata(RegionId::new(1024, 1), &column_metadatas);
1377        let result = check_column_metadatas_consistent(&[region_metadata1, region_metadata2]);
1378        assert!(result.is_none());
1379    }
1380
1381    #[test]
1382    fn test_check_column_metadata_invariants() {
1383        let column_metadatas = new_test_column_metadatas();
1384        let mut new_column_metadatas = column_metadatas.clone();
1385        new_column_metadatas.push(ColumnMetadata {
1386            column_schema: ColumnSchema::new("col3", ConcreteDataType::int32_datatype(), true),
1387            semantic_type: SemanticType::Field,
1388            column_id: 3,
1389        });
1390        check_column_metadata_invariants(&new_column_metadatas, &column_metadatas).unwrap();
1391    }
1392
1393    #[test]
1394    fn test_check_column_metadata_invariants_missing_primary_key_column_or_ts_column() {
1395        let column_metadatas = new_test_column_metadatas();
1396        let mut new_column_metadatas = column_metadatas.clone();
1397        new_column_metadatas.retain(|c| c.semantic_type != SemanticType::Timestamp);
1398        check_column_metadata_invariants(&new_column_metadatas, &column_metadatas).unwrap_err();
1399
1400        let column_metadatas = new_test_column_metadatas();
1401        let mut new_column_metadatas = column_metadatas.clone();
1402        new_column_metadatas.retain(|c| c.semantic_type != SemanticType::Tag);
1403        check_column_metadata_invariants(&new_column_metadatas, &column_metadatas).unwrap_err();
1404    }
1405
1406    #[test]
1407    fn test_check_column_metadata_invariants_mismatch_column_id() {
1408        let column_metadatas = new_test_column_metadatas();
1409        let mut new_column_metadatas = column_metadatas.clone();
1410        if let Some(col) = new_column_metadatas
1411            .iter_mut()
1412            .find(|c| c.semantic_type == SemanticType::Timestamp)
1413        {
1414            col.column_id = 100;
1415        }
1416        check_column_metadata_invariants(&new_column_metadatas, &column_metadatas).unwrap_err();
1417
1418        let column_metadatas = new_test_column_metadatas();
1419        let mut new_column_metadatas = column_metadatas.clone();
1420        if let Some(col) = new_column_metadatas
1421            .iter_mut()
1422            .find(|c| c.semantic_type == SemanticType::Tag)
1423        {
1424            col.column_id = 100;
1425        }
1426        check_column_metadata_invariants(&new_column_metadatas, &column_metadatas).unwrap_err();
1427    }
1428
1429    #[test]
1430    fn test_resolve_column_metadatas_with_use_metasrv_strategy() {
1431        let column_metadatas = new_test_column_metadatas();
1432        let region_metadata1 = build_region_metadata(RegionId::new(1024, 0), &column_metadatas);
1433        let mut metasrv_column_metadatas = region_metadata1.column_metadatas.clone();
1434        metasrv_column_metadatas.push(ColumnMetadata {
1435            column_schema: ColumnSchema::new("col3", ConcreteDataType::int32_datatype(), true),
1436            semantic_type: SemanticType::Field,
1437            column_id: 3,
1438        });
1439        let result =
1440            resolve_column_metadatas_with_metasrv(&metasrv_column_metadatas, &[region_metadata1])
1441                .unwrap();
1442
1443        assert_eq!(result, vec![RegionId::new(1024, 0)]);
1444    }
1445
1446    #[test]
1447    fn test_resolve_column_metadatas_with_use_latest_strategy() {
1448        let column_metadatas = new_test_column_metadatas();
1449        let region_metadata1 = build_region_metadata(RegionId::new(1024, 0), &column_metadatas);
1450        let mut new_column_metadatas = column_metadatas.clone();
1451        new_column_metadatas.push(ColumnMetadata {
1452            column_schema: ColumnSchema::new("col3", ConcreteDataType::int32_datatype(), true),
1453            semantic_type: SemanticType::Field,
1454            column_id: 3,
1455        });
1456
1457        let mut region_metadata2 =
1458            build_region_metadata(RegionId::new(1024, 1), &new_column_metadatas);
1459        region_metadata2.schema_version = 2;
1460
1461        let (resolved_column_metadatas, region_ids) =
1462            resolve_column_metadatas_with_latest(&[region_metadata1, region_metadata2]).unwrap();
1463        assert_eq!(region_ids, vec![RegionId::new(1024, 0)]);
1464        assert_eq!(resolved_column_metadatas, new_column_metadatas);
1465    }
1466}