1pub(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
60pub 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
82pub 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
100pub 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 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
181pub 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
196pub 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
215pub 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
226pub enum MultipleResults<T> {
234 Ok(Vec<T>),
235 PartialRetryable(Error),
236 PartialNonRetryable(Error),
237 AllRetryable(Error),
238 AllNonRetryable(Error),
239}
240
241pub 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 MultipleResults::PartialNonRetryable(non_retryable_results.into_iter().next().unwrap())
303}
304
305pub 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
320pub 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
332pub 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(®ion_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(®ion_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 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
436pub 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 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 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}