Skip to main content

common_meta/ddl/
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
15pub(crate) mod raw_table_info;
16#[allow(dead_code)]
17pub(crate) mod region_metadata_lister;
18pub(crate) mod table_id;
19pub(crate) mod table_info;
20
21use std::collections::HashMap;
22use std::fmt::Debug;
23
24use api::region::RegionResponse;
25use api::v1::region::sync_request::ManifestInfo;
26use api::v1::region::{
27    MetricManifestInfo, MitoManifestInfo, RegionRequest, RegionRequestHeader, SyncRequest,
28    region_request,
29};
30use common_catalog::consts::{METRIC_ENGINE, MITO_ENGINE};
31use common_error::ext::BoxedError;
32use common_procedure::error::Error as ProcedureError;
33use common_telemetry::tracing_context::TracingContext;
34use common_telemetry::{error, info, warn};
35use futures::future::join_all;
36use snafu::{OptionExt, ResultExt, ensure};
37use store_api::metadata::ColumnMetadata;
38use store_api::metric_engine_consts::{LOGICAL_TABLE_METADATA_KEY, MANIFEST_INFO_EXTENSION_KEY};
39use store_api::region_engine::RegionManifestInfo;
40use store_api::storage::RegionId;
41use table::metadata::TableId;
42#[cfg(feature = "enterprise")]
43use table::metadata::TableInfo;
44use table::table_reference::TableReference;
45
46use crate::ddl::{DdlContext, DetectingRegion};
47use crate::error::{
48    self, DecodeJsonSnafu, Error, MetadataCorruptionSnafu, OperateDatanodeSnafu, Result,
49    TableNotFoundSnafu, UnsupportedSnafu,
50};
51use crate::key::datanode_table::DatanodeTableValue;
52use crate::key::table_name::TableNameKey;
53use crate::key::table_route::TableRouteValue;
54use crate::key::{TableMetadataManager, TableMetadataManagerRef};
55use crate::peer::Peer;
56use crate::rpc::ddl::CreateTableTask;
57use crate::rpc::router::{RegionRoute, find_follower_regions, find_followers};
58use crate::wal_provider::RegionWalOptions;
59
60/// Adds [Peer] context if the error is unretryable.
61pub fn add_peer_context_if_needed(datanode: Peer) -> impl FnOnce(Error) -> Error {
62    move |err| {
63        error!(err; "Failed to operate datanode, peer: {}", datanode);
64        if !err.is_retry_later() {
65            return Err::<(), BoxedError>(BoxedError::new(err))
66                .context(OperateDatanodeSnafu { peer: datanode })
67                .unwrap_err();
68        }
69        err
70    }
71}
72
73#[cfg(feature = "enterprise")]
74pub(crate) fn is_metric_engine_logical_table(
75    table_info: &TableInfo,
76    table_route_value: &TableRouteValue,
77) -> bool {
78    table_info.meta.engine == METRIC_ENGINE
79        && matches!(table_route_value, TableRouteValue::Logical(_))
80}
81
82/// Maps the error to the corresponding procedure error.
83///
84/// This function determines whether the error should be retried and if poison cleanup is needed,
85/// then maps it to the appropriate procedure error variant.
86pub fn map_to_procedure_error(e: Error) -> ProcedureError {
87    match (e.is_retry_later(), e.need_clean_poisons()) {
88        (true, true) => ProcedureError::retry_later_and_clean_poisons(e),
89        (true, false) => ProcedureError::retry_later(e),
90        (false, true) => ProcedureError::external_and_clean_poisons(e),
91        (false, false) => ProcedureError::external(e),
92    }
93}
94
95#[inline]
96pub fn region_storage_path(catalog: &str, schema: &str) -> String {
97    format!("{}/{}", catalog, schema)
98}
99
100/// Extracts catalog and schema from the path that created by [region_storage_path].
101pub fn get_catalog_and_schema(path: &str) -> Option<(String, String)> {
102    let mut split = path.split('/');
103    Some((split.next()?.to_string(), split.next()?.to_string()))
104}
105
106pub async fn check_and_get_physical_table_id(
107    table_metadata_manager: &TableMetadataManagerRef,
108    tasks: &[CreateTableTask],
109) -> Result<TableId> {
110    let mut physical_table_name = None;
111    for task in tasks {
112        ensure!(
113            task.create_table.engine == METRIC_ENGINE,
114            UnsupportedSnafu {
115                operation: format!("create table with engine {}", task.create_table.engine)
116            }
117        );
118        let current_physical_table_name = task
119            .create_table
120            .table_options
121            .get(LOGICAL_TABLE_METADATA_KEY)
122            .context(UnsupportedSnafu {
123                operation: format!(
124                    "create table without table options {}",
125                    LOGICAL_TABLE_METADATA_KEY,
126                ),
127            })?;
128        let current_physical_table_name = TableNameKey::new(
129            &task.create_table.catalog_name,
130            &task.create_table.schema_name,
131            current_physical_table_name,
132        );
133
134        physical_table_name = match physical_table_name {
135            Some(name) => {
136                ensure!(
137                    name == current_physical_table_name,
138                    UnsupportedSnafu {
139                        operation: format!(
140                            "create table with different physical table name {} and {}",
141                            name, current_physical_table_name
142                        )
143                    }
144                );
145                Some(name)
146            }
147            None => Some(current_physical_table_name),
148        };
149    }
150    // Safety: `physical_table_name` is `Some` here
151    let physical_table_name = physical_table_name.unwrap();
152    table_metadata_manager
153        .table_name_manager()
154        .get(physical_table_name)
155        .await?
156        .with_context(|| TableNotFoundSnafu {
157            table_name: TableReference::from(physical_table_name).to_string(),
158        })
159        .map(|table| table.table_id())
160}
161
162pub async fn get_physical_table_id(
163    table_metadata_manager: &TableMetadataManagerRef,
164    logical_table_name: TableNameKey<'_>,
165) -> Result<TableId> {
166    let logical_table_id = table_metadata_manager
167        .table_name_manager()
168        .get(logical_table_name)
169        .await?
170        .with_context(|| TableNotFoundSnafu {
171            table_name: TableReference::from(logical_table_name).to_string(),
172        })
173        .map(|table| table.table_id())?;
174
175    table_metadata_manager
176        .table_route_manager()
177        .get_physical_table_id(logical_table_id)
178        .await
179}
180
181/// Converts a list of [`RegionRoute`] to a list of [`DetectingRegion`].
182pub fn convert_region_routes_to_detecting_regions(
183    region_routes: &[RegionRoute],
184) -> Vec<DetectingRegion> {
185    region_routes
186        .iter()
187        .flat_map(|route| {
188            route
189                .leader_peer
190                .as_ref()
191                .map(|peer| (peer.id, route.region.id))
192        })
193        .collect::<Vec<_>>()
194}
195
196/// Gets the wal options for a table.
197pub async fn get_region_wal_options(
198    table_metadata_manager: &TableMetadataManager,
199    table_route_value: &TableRouteValue,
200    physical_table_id: TableId,
201) -> Result<RegionWalOptions> {
202    let region_wal_options =
203        if let TableRouteValue::Physical(table_route_value) = &table_route_value {
204            let datanode_table_values = table_metadata_manager
205                .datanode_table_manager()
206                .regions(physical_table_id, table_route_value)
207                .await?;
208            extract_region_wal_options(&datanode_table_values)?
209        } else {
210            HashMap::new()
211        };
212    Ok(region_wal_options)
213}
214
215/// Extracts region wal options from [DatanodeTableValue]s.
216pub fn extract_region_wal_options(
217    datanode_table_values: &Vec<DatanodeTableValue>,
218) -> Result<RegionWalOptions> {
219    let mut region_wal_options = RegionWalOptions::new();
220    for value in datanode_table_values {
221        region_wal_options.extend(value.region_info.region_wal_options.clone());
222    }
223    Ok(region_wal_options)
224}
225
226/// The result of multiple operations.
227///
228/// - Ok: all operations are successful.
229/// - PartialRetryable: if any operation is retryable and without non retryable error, the result is retryable.
230/// - PartialNonRetryable: if any operation is non retryable, the result is non retryable.
231/// - AllRetryable: all operations are retryable.
232/// - AllNonRetryable: all operations are not retryable.
233pub enum MultipleResults<T> {
234    Ok(Vec<T>),
235    PartialRetryable(Error),
236    PartialNonRetryable(Error),
237    AllRetryable(Error),
238    AllNonRetryable(Error),
239}
240
241/// Handles the results of alter region requests.
242///
243/// For partial success, we need to check if the errors are retryable.
244/// If all the errors are retryable, we return a retryable error.
245/// Otherwise, we return the first error.
246pub fn handle_multiple_results<T: Debug>(results: Vec<Result<T>>) -> MultipleResults<T> {
247    if results.is_empty() {
248        return MultipleResults::Ok(Vec::new());
249    }
250    let num_results = results.len();
251    let mut retryable_results = Vec::new();
252    let mut non_retryable_results = Vec::new();
253    let mut ok_results = Vec::new();
254
255    for result in results {
256        match result {
257            Ok(value) => ok_results.push(value),
258            Err(err) => {
259                if err.is_retry_later() {
260                    retryable_results.push(err);
261                } else {
262                    non_retryable_results.push(err);
263                }
264            }
265        }
266    }
267
268    common_telemetry::debug!(
269        "retryable_results: {}, non_retryable_results: {}, ok_results: {}",
270        retryable_results.len(),
271        non_retryable_results.len(),
272        ok_results.len()
273    );
274
275    if retryable_results.len() == num_results {
276        return MultipleResults::AllRetryable(retryable_results.into_iter().next().unwrap());
277    } else if non_retryable_results.len() == num_results {
278        warn!("all non retryable results: {}", non_retryable_results.len());
279        for err in &non_retryable_results {
280            error!(err; "non retryable error");
281        }
282        return MultipleResults::AllNonRetryable(non_retryable_results.into_iter().next().unwrap());
283    } else if ok_results.len() == num_results {
284        return MultipleResults::Ok(ok_results);
285    } else if !retryable_results.is_empty()
286        && !ok_results.is_empty()
287        && non_retryable_results.is_empty()
288    {
289        return MultipleResults::PartialRetryable(retryable_results.into_iter().next().unwrap());
290    }
291
292    warn!(
293        "partial non retryable results: {}, retryable results: {}, ok results: {}",
294        non_retryable_results.len(),
295        retryable_results.len(),
296        ok_results.len()
297    );
298    for err in &non_retryable_results {
299        error!(err; "non retryable error");
300    }
301    // non_retryable_results.len() > 0
302    MultipleResults::PartialNonRetryable(non_retryable_results.into_iter().next().unwrap())
303}
304
305/// Parses manifest infos from extensions.
306pub fn parse_manifest_infos_from_extensions(
307    extensions: &HashMap<String, Vec<u8>>,
308) -> Result<Vec<(RegionId, RegionManifestInfo)>> {
309    let data_manifest_version =
310        extensions
311            .get(MANIFEST_INFO_EXTENSION_KEY)
312            .context(error::UnexpectedSnafu {
313                err_msg: "manifest info extension not found",
314            })?;
315    let data_manifest_version =
316        RegionManifestInfo::decode_list(data_manifest_version).context(error::SerdeJsonSnafu {})?;
317    Ok(data_manifest_version)
318}
319
320/// Parses column metadatas from extensions.
321pub fn parse_column_metadatas(
322    extensions: &HashMap<String, Vec<u8>>,
323    key: &str,
324) -> Result<Vec<ColumnMetadata>> {
325    let value = extensions.get(key).context(error::UnexpectedSnafu {
326        err_msg: format!("column metadata extension not found: {}", key),
327    })?;
328    let column_metadatas = ColumnMetadata::decode_list(value).context(error::SerdeJsonSnafu {})?;
329    Ok(column_metadatas)
330}
331
332/// Sync follower regions on datanodes.
333pub async fn sync_follower_regions(
334    context: &DdlContext,
335    table_id: TableId,
336    results: &[RegionResponse],
337    region_routes: &[RegionRoute],
338    engine: &str,
339) -> Result<()> {
340    if engine != MITO_ENGINE && engine != METRIC_ENGINE {
341        info!(
342            "Skip submitting sync region requests for table_id: {}, engine: {}",
343            table_id, engine
344        );
345        return Ok(());
346    }
347
348    let results = results
349        .iter()
350        .map(|response| parse_manifest_infos_from_extensions(&response.extensions))
351        .collect::<Result<Vec<_>>>()?
352        .into_iter()
353        .flatten()
354        .collect::<HashMap<_, _>>();
355
356    let is_mito_engine = engine == MITO_ENGINE;
357
358    let followers = find_followers(region_routes);
359    if followers.is_empty() {
360        return Ok(());
361    }
362    let mut sync_region_tasks = Vec::with_capacity(followers.len());
363    for datanode in followers {
364        let requester = context.node_manager.datanode(&datanode).await;
365        let regions = find_follower_regions(region_routes, &datanode);
366        for region in regions {
367            let region_id = RegionId::new(table_id, region);
368            let manifest_info = if is_mito_engine {
369                let region_manifest_info =
370                    results.get(&region_id).context(error::UnexpectedSnafu {
371                        err_msg: format!("No manifest info found for region {}", region_id),
372                    })?;
373                ensure!(
374                    region_manifest_info.is_mito(),
375                    error::UnexpectedSnafu {
376                        err_msg: format!("Region {} is not a mito region", region_id)
377                    }
378                );
379                ManifestInfo::MitoManifestInfo(MitoManifestInfo {
380                    data_manifest_version: region_manifest_info.data_manifest_version(),
381                })
382            } else {
383                let region_manifest_info =
384                    results.get(&region_id).context(error::UnexpectedSnafu {
385                        err_msg: format!("No manifest info found for region {}", region_id),
386                    })?;
387                ensure!(
388                    region_manifest_info.is_metric(),
389                    error::UnexpectedSnafu {
390                        err_msg: format!("Region {} is not a metric region", region_id)
391                    }
392                );
393                ManifestInfo::MetricManifestInfo(MetricManifestInfo {
394                    data_manifest_version: region_manifest_info.data_manifest_version(),
395                    metadata_manifest_version: region_manifest_info
396                        .metadata_manifest_version()
397                        .unwrap_or_default(),
398                })
399            };
400            let request = RegionRequest {
401                header: Some(RegionRequestHeader {
402                    tracing_context: TracingContext::from_current_span().to_w3c(),
403                    ..Default::default()
404                }),
405                body: Some(region_request::Body::Sync(SyncRequest {
406                    region_id: region_id.as_u64(),
407                    manifest_info: Some(manifest_info),
408                })),
409            };
410
411            let datanode = datanode.clone();
412            let requester = requester.clone();
413            sync_region_tasks.push(async move {
414                requester
415                    .handle(request)
416                    .await
417                    .map_err(add_peer_context_if_needed(datanode))
418            });
419        }
420    }
421
422    // Failure to sync region is not critical.
423    // We try our best to sync the regions.
424    if let Err(err) = join_all(sync_region_tasks)
425        .await
426        .into_iter()
427        .collect::<Result<Vec<_>>>()
428    {
429        error!(err; "Failed to sync follower regions on datanodes, table_id: {}", table_id);
430    }
431    info!("Sync follower regions on datanodes, table_id: {}", table_id);
432
433    Ok(())
434}
435
436/// Extracts column metadatas from extensions.
437pub fn extract_column_metadatas(
438    results: &mut [RegionResponse],
439    key: &str,
440) -> Result<Option<Vec<ColumnMetadata>>> {
441    let mut schemas = results
442        .iter_mut()
443        .map(|r| r.extensions.remove(key))
444        .collect::<Vec<_>>();
445
446    if schemas.is_empty() {
447        warn!("extract_column_metadatas: no extension key `{key}` found in results");
448        return Ok(None);
449    }
450
451    // Verify all the physical schemas are the same
452    // Safety: previous check ensures this vec is not empty
453    let first_column_metadatas = schemas
454        .swap_remove(0)
455        .map(|first_bytes| ColumnMetadata::decode_list(&first_bytes).context(DecodeJsonSnafu))
456        .transpose()?;
457
458    for s in schemas {
459        // check decoded column metadata instead of bytes because it contains extension map.
460        let column_metadata = s
461            .map(|bytes| ColumnMetadata::decode_list(&bytes).context(DecodeJsonSnafu))
462            .transpose()?;
463        ensure!(
464            column_metadata == first_column_metadatas,
465            MetadataCorruptionSnafu {
466                err_msg: format!(
467                    "The table column metadata schemas from datanodes are not the same. First: {:?}, Current: {:?}",
468                    first_column_metadatas, column_metadata,
469                ),
470            }
471        );
472    }
473    Ok(first_column_metadatas)
474}
475
476#[cfg(test)]
477mod tests {
478    use super::*;
479
480    #[test]
481    fn test_get_catalog_and_schema() {
482        let test_catalog = "my_catalog";
483        let test_schema = "my_schema";
484        let path = region_storage_path(test_catalog, test_schema);
485        let (catalog, schema) = get_catalog_and_schema(&path).unwrap();
486        assert_eq!(catalog, test_catalog);
487        assert_eq!(schema, test_schema);
488    }
489}