Skip to main content

common_meta/reconciliation/reconcile_table/
resolve_column_metadata.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 async_trait::async_trait;
16use common_procedure::{Context as ProcedureContext, Status};
17use common_telemetry::info;
18use serde::{Deserialize, Serialize};
19use snafu::OptionExt;
20use store_api::metadata::RegionMetadata;
21use strum::AsRefStr;
22
23use crate::error::{self, MissingColumnIdsSnafu, Result};
24use crate::reconciliation::reconcile_table::reconcile_regions::ReconcileRegions;
25use crate::reconciliation::reconcile_table::update_table_info::UpdateTableInfo;
26use crate::reconciliation::reconcile_table::{ReconcileTableContext, State, TableMetadataState};
27use crate::reconciliation::utils::{
28    ResolveColumnMetadataResult, build_reconciliation_column_metadata,
29    check_column_metadatas_consistent, reorder_tag_columns, resolve_column_metadatas_with_latest,
30    resolve_column_metadatas_with_metasrv,
31};
32
33/// Strategy for resolving column metadata inconsistencies.
34#[derive(Debug, Serialize, Deserialize, Clone, Copy, Default, AsRefStr)]
35pub enum ResolveStrategy {
36    #[default]
37    /// Trusts the latest column metadata from datanode.
38    UseLatest,
39
40    /// Always uses the column metadata from metasrv.
41    UseMetasrv,
42
43    /// Aborts the resolution process if inconsistencies are detected.
44    AbortOnConflict,
45}
46
47impl From<api::v1::meta::ResolveStrategy> for ResolveStrategy {
48    fn from(strategy: api::v1::meta::ResolveStrategy) -> Self {
49        match strategy {
50            api::v1::meta::ResolveStrategy::UseMetasrv => Self::UseMetasrv,
51            api::v1::meta::ResolveStrategy::UseLatest => Self::UseLatest,
52            api::v1::meta::ResolveStrategy::AbortOnConflict => Self::AbortOnConflict,
53        }
54    }
55}
56
57/// State responsible for resolving inconsistencies in column metadata across physical regions.
58#[derive(Debug, Serialize, Deserialize)]
59pub struct ResolveColumnMetadata {
60    strategy: ResolveStrategy,
61    region_metadata: Vec<RegionMetadata>,
62}
63
64impl ResolveColumnMetadata {
65    pub fn new(strategy: ResolveStrategy, region_metadata: Vec<RegionMetadata>) -> Self {
66        Self {
67            strategy,
68            region_metadata,
69        }
70    }
71}
72
73#[async_trait]
74#[typetag::serde]
75impl State for ResolveColumnMetadata {
76    async fn next(
77        &mut self,
78        ctx: &mut ReconcileTableContext,
79        _procedure_ctx: &ProcedureContext,
80    ) -> Result<(Box<dyn State>, Status)> {
81        let table_id = ctx.persistent_ctx.table_id;
82        let table_name = &ctx.persistent_ctx.table_name;
83
84        let table_info_value = ctx
85            .table_metadata_manager
86            .table_info_manager()
87            .get(table_id)
88            .await?
89            .with_context(|| error::TableNotFoundSnafu {
90                table_name: table_name.to_string(),
91            })?;
92        ctx.persistent_ctx.table_info_value = Some(table_info_value);
93
94        if let Some(column_metadatas) = check_column_metadatas_consistent(&self.region_metadata) {
95            let column_metadatas =
96                reorder_tag_columns(&column_metadatas, &self.region_metadata[0].primary_key)?;
97            // Safety: fetched in the above.
98            let table_info_value = ctx.persistent_ctx.table_info_value.clone().unwrap();
99            info!(
100                "Column metadatas are consistent for table: {}, table_id: {}.",
101                table_name, table_id
102            );
103
104            ctx.volatile_ctx.result_summary.record_resolved_columns(
105                TableMetadataState::Consistent,
106                None,
107                Some(column_metadatas.len()),
108            );
109
110            // Update metrics.
111            ctx.mut_metrics().resolve_column_metadata_result =
112                Some(ResolveColumnMetadataResult::Consistent);
113            return Ok((
114                Box::new(UpdateTableInfo::new(table_info_value, column_metadatas)),
115                Status::executing(false),
116            ));
117        };
118
119        ctx.volatile_ctx
120            .result_summary
121            .record_metadata_state(TableMetadataState::Inconsistent);
122
123        match self.strategy {
124            ResolveStrategy::UseMetasrv => {
125                let table_info_value = ctx.persistent_ctx.table_info_value.as_ref().unwrap();
126                let name_to_ids = table_info_value
127                    .table_info
128                    .name_to_ids()
129                    .context(MissingColumnIdsSnafu)?;
130                let column_metadata = build_reconciliation_column_metadata(
131                    table_info_value.table_info.meta.schema.column_schemas(),
132                    &table_info_value.table_info.meta.primary_key_indices,
133                    &name_to_ids,
134                )?;
135
136                let region_ids =
137                    resolve_column_metadatas_with_metasrv(&column_metadata, &self.region_metadata)?;
138
139                ctx.volatile_ctx.result_summary.record_resolved_columns(
140                    TableMetadataState::Inconsistent,
141                    Some(self.strategy),
142                    Some(column_metadata.len()),
143                );
144
145                // Update metrics.
146                let metrics = ctx.mut_metrics();
147                metrics.resolve_column_metadata_result =
148                    Some(ResolveColumnMetadataResult::Inconsistent(self.strategy));
149                Ok((
150                    Box::new(ReconcileRegions::new(column_metadata, region_ids)),
151                    Status::executing(true),
152                ))
153            }
154            ResolveStrategy::UseLatest => {
155                let (column_metadatas, region_ids) =
156                    resolve_column_metadatas_with_latest(&self.region_metadata)?;
157
158                ctx.volatile_ctx.result_summary.record_resolved_columns(
159                    TableMetadataState::Inconsistent,
160                    Some(self.strategy),
161                    Some(column_metadatas.len()),
162                );
163
164                // Update metrics.
165                let metrics = ctx.mut_metrics();
166                metrics.resolve_column_metadata_result =
167                    Some(ResolveColumnMetadataResult::Inconsistent(self.strategy));
168                Ok((
169                    Box::new(ReconcileRegions::new(column_metadatas, region_ids)),
170                    Status::executing(true),
171                ))
172            }
173            ResolveStrategy::AbortOnConflict => {
174                let table_name = table_name.to_string();
175
176                ctx.volatile_ctx
177                    .result_summary
178                    .record_resolution_strategy(self.strategy);
179
180                // Update metrics.
181                let metrics = ctx.mut_metrics();
182                metrics.resolve_column_metadata_result =
183                    Some(ResolveColumnMetadataResult::Inconsistent(self.strategy));
184                error::ColumnMetadataConflictsSnafu {
185                    table_name,
186                    table_id,
187                }
188                .fail()
189            }
190        }
191    }
192}