1use api::region::RegionResponse;
18use api::v1::SemanticType;
19use common_error::ext::BoxedError;
20use common_telemetry::{error, info, warn};
21use datafusion::common::HashMap;
22use mito2::engine::MITO_ENGINE_NAME;
23use snafu::{OptionExt, ResultExt};
24use store_api::region_engine::{BatchResponses, RegionEngine};
25use store_api::region_request::{
26 AffectedRows, PathType, RegionCleanUpRequest, RegionOpenRequest, RegionRequest,
27 ReplayCheckpoint,
28};
29use store_api::storage::RegionId;
30
31use crate::engine::MetricEngineInner;
32use crate::engine::create::region_options_for_metadata_region;
33use crate::engine::options::{PhysicalRegionOptions, set_data_region_options};
34use crate::error::{
35 BatchOpenMitoRegionSnafu, CleanUpMitoRegionSnafu, NoOpenRegionResultSnafu, OpenMitoRegionSnafu,
36 PhysicalRegionNotFoundSnafu, Result,
37};
38use crate::metrics::{LOGICAL_REGION_COUNT, PHYSICAL_REGION_COUNT};
39use crate::utils;
40
41impl MetricEngineInner {
42 pub async fn handle_batch_open_requests(
43 &self,
44 parallelism: usize,
45 requests: Vec<(RegionId, RegionOpenRequest)>,
46 ) -> Result<BatchResponses> {
47 let mut all_requests = Vec::with_capacity(requests.len() * 2);
49 let mut physical_region_ids = HashMap::with_capacity(requests.len());
50
51 for (region_id, request) in requests {
52 if !request.is_physical_table() {
53 warn!("Skipping non-physical table open request: {region_id}");
54 continue;
55 }
56 let physical_region_options = PhysicalRegionOptions::try_from(&request.options)?;
57 let metadata_region_id = utils::to_metadata_region_id(region_id);
58 let data_region_id = utils::to_data_region_id(region_id);
59 let (open_metadata_region_request, open_data_region_request) =
60 self.transform_open_physical_region_request(request);
61 all_requests.push((metadata_region_id, open_metadata_region_request));
62 all_requests.push((data_region_id, open_data_region_request));
63 physical_region_ids.insert(region_id, physical_region_options);
64 }
65
66 let mut results = self
67 .mito
68 .handle_batch_open_requests(parallelism, all_requests)
69 .await
70 .context(BatchOpenMitoRegionSnafu {})?
71 .into_iter()
72 .collect::<HashMap<_, _>>();
73
74 let mut responses = Vec::with_capacity(physical_region_ids.len());
75 for (physical_region_id, physical_region_options) in physical_region_ids {
76 let metadata_region_id = utils::to_metadata_region_id(physical_region_id);
77 let data_region_id = utils::to_data_region_id(physical_region_id);
78 let metadata_region_result = results.remove(&metadata_region_id);
79 let data_region_result: Option<std::result::Result<RegionResponse, BoxedError>> =
80 results.remove(&data_region_id);
81 let response = self
86 .recover_physical_region_with_results(
87 metadata_region_result,
88 data_region_result,
89 physical_region_id,
90 physical_region_options,
91 true,
92 )
93 .await
94 .map_err(BoxedError::new);
95 responses.push((physical_region_id, response));
96 }
97
98 Ok(responses)
99 }
100
101 async fn close_physical_region_on_recovery_failure(&self, physical_region_id: RegionId) {
106 info!(
107 "Closing metadata region {} and data region {} on metadata recovery failure",
108 utils::to_metadata_region_id(physical_region_id),
109 utils::to_data_region_id(physical_region_id)
110 );
111 if let Err(err) = self.close_physical_region(physical_region_id, false).await {
112 error!(err; "Failed to close physical region {}", physical_region_id);
113 }
114 }
115
116 pub(crate) async fn recover_physical_region_with_results(
117 &self,
118 metadata_region_result: Option<std::result::Result<RegionResponse, BoxedError>>,
119 data_region_result: Option<std::result::Result<RegionResponse, BoxedError>>,
120 physical_region_id: RegionId,
121 physical_region_options: PhysicalRegionOptions,
122 close_region_on_failure: bool,
123 ) -> Result<RegionResponse> {
124 let metadata_region_id = utils::to_metadata_region_id(physical_region_id);
125 let data_region_id = utils::to_data_region_id(physical_region_id);
126 let _ = metadata_region_result
127 .context(NoOpenRegionResultSnafu {
128 region_id: metadata_region_id,
129 })?
130 .context(OpenMitoRegionSnafu {
131 region_type: "metadata",
132 })?;
133
134 let data_region_response = data_region_result
135 .context(NoOpenRegionResultSnafu {
136 region_id: data_region_id,
137 })?
138 .context(OpenMitoRegionSnafu {
139 region_type: "data",
140 })?;
141
142 if let Err(err) = self
143 .recover_states(physical_region_id, physical_region_options)
144 .await
145 {
146 if close_region_on_failure {
147 self.close_physical_region_on_recovery_failure(physical_region_id)
148 .await;
149 }
150 return Err(err);
151 }
152 Ok(data_region_response)
153 }
154
155 pub async fn open_region(
165 &self,
166 region_id: RegionId,
167 request: RegionOpenRequest,
168 ) -> Result<AffectedRows> {
169 if request.is_physical_table() {
170 if self
171 .state
172 .read()
173 .unwrap()
174 .physical_region_states()
175 .get(®ion_id)
176 .is_some()
177 {
178 warn!(
179 "The physical region {} is already open, ignore the open request",
180 region_id
181 );
182 return Ok(0);
183 }
184 let physical_region_options = PhysicalRegionOptions::try_from(&request.options)?;
186 self.open_physical_region(region_id, request).await?;
187 if let Err(err) = self
188 .recover_states(region_id, physical_region_options)
189 .await
190 {
191 self.close_physical_region_on_recovery_failure(region_id)
192 .await;
193 return Err(err);
194 }
195
196 Ok(0)
197 } else {
198 Ok(0)
203 }
204 }
205
206 pub async fn clean_up_region(
207 &self,
208 region_id: RegionId,
209 request: RegionCleanUpRequest,
210 ) -> Result<AffectedRows> {
211 if request.is_physical_table() {
212 self.cleanup_physical_region_offline(region_id, request)
213 .await
214 } else {
215 Ok(0)
216 }
217 }
218
219 async fn cleanup_physical_region_offline(
220 &self,
221 region_id: RegionId,
222 request: RegionCleanUpRequest,
223 ) -> Result<AffectedRows> {
224 let metadata_region_id = utils::to_metadata_region_id(region_id);
225 let data_region_id = utils::to_data_region_id(region_id);
226 let (clean_up_metadata_region_request, clean_up_data_region_request) =
227 self.transform_clean_up_physical_region_request(request);
228 let _ = self
229 .mito
230 .handle_request(
231 metadata_region_id,
232 RegionRequest::CleanUp(clean_up_metadata_region_request),
233 )
234 .await
235 .context(CleanUpMitoRegionSnafu {
236 region_type: "metadata",
237 })?;
238 let data_region_response = self
239 .mito
240 .handle_request(
241 data_region_id,
242 RegionRequest::CleanUp(clean_up_data_region_request),
243 )
244 .await
245 .context(CleanUpMitoRegionSnafu {
246 region_type: "data",
247 })?;
248
249 if self.state.read().unwrap().exist_physical_region(region_id) {
250 self.state
251 .write()
252 .unwrap()
253 .remove_physical_region(region_id)?;
254 PHYSICAL_REGION_COUNT.dec();
255 }
256
257 Ok(data_region_response.affected_rows)
258 }
259
260 fn transform_clean_up_physical_region_request(
266 &self,
267 request: RegionCleanUpRequest,
268 ) -> (RegionCleanUpRequest, RegionCleanUpRequest) {
269 let clean_up_metadata_region_request = RegionCleanUpRequest {
270 table_dir: request.table_dir.clone(),
271 path_type: PathType::Metadata,
272 options: region_options_for_metadata_region(&request.options),
273 engine: MITO_ENGINE_NAME.to_string(),
274 };
275
276 let mut data_region_options = request.options;
277 set_data_region_options(&mut data_region_options);
278 let clean_up_data_region_request = RegionCleanUpRequest {
279 table_dir: request.table_dir,
280 path_type: PathType::Data,
281 options: data_region_options,
282 engine: MITO_ENGINE_NAME.to_string(),
283 };
284
285 (
286 clean_up_metadata_region_request,
287 clean_up_data_region_request,
288 )
289 }
290
291 fn transform_open_physical_region_request(
297 &self,
298 request: RegionOpenRequest,
299 ) -> (RegionOpenRequest, RegionOpenRequest) {
300 let metadata_region_options = region_options_for_metadata_region(&request.options);
301 let checkpoint = request.checkpoint;
302
303 let open_metadata_region_request = RegionOpenRequest {
304 table_dir: request.table_dir.clone(),
305 path_type: PathType::Metadata,
306 options: metadata_region_options,
307 engine: MITO_ENGINE_NAME.to_string(),
308 skip_wal_replay: request.skip_wal_replay,
309 checkpoint: checkpoint.map(|checkpoint| ReplayCheckpoint {
310 entry_id: checkpoint.metadata_entry_id.unwrap_or_default(),
311 metadata_entry_id: None,
312 }),
313 requirements: request.requirements,
314 };
315
316 let mut data_region_options = request.options;
317 set_data_region_options(&mut data_region_options);
318 let open_data_region_request = RegionOpenRequest {
319 table_dir: request.table_dir.clone(),
320 path_type: PathType::Data,
321 options: data_region_options,
322 engine: MITO_ENGINE_NAME.to_string(),
323 skip_wal_replay: request.skip_wal_replay,
324 checkpoint: checkpoint.map(|checkpoint| ReplayCheckpoint {
325 entry_id: checkpoint.entry_id,
326 metadata_entry_id: None,
327 }),
328 requirements: request.requirements,
329 };
330
331 (open_metadata_region_request, open_data_region_request)
332 }
333
334 async fn open_physical_region(
336 &self,
337 region_id: RegionId,
338 request: RegionOpenRequest,
339 ) -> Result<AffectedRows> {
340 let metadata_region_id = utils::to_metadata_region_id(region_id);
341 let data_region_id = utils::to_data_region_id(region_id);
342 let (open_metadata_region_request, open_data_region_request) =
343 self.transform_open_physical_region_request(request);
344 let _ = self
345 .mito
346 .handle_batch_open_requests(
347 2,
348 vec![
349 (metadata_region_id, open_metadata_region_request),
350 (data_region_id, open_data_region_request),
351 ],
352 )
353 .await
354 .context(BatchOpenMitoRegionSnafu {})?;
355
356 info!("Opened physical metric region {region_id}");
357 PHYSICAL_REGION_COUNT.inc();
358
359 Ok(0)
360 }
361
362 pub(crate) async fn recover_states(
371 &self,
372 physical_region_id: RegionId,
373 physical_region_options: PhysicalRegionOptions,
374 ) -> Result<Vec<RegionId>> {
375 let logical_regions = self
377 .metadata_region
378 .logical_regions(physical_region_id)
379 .await?;
380 common_telemetry::debug!(
381 "Recover states for physical region {}, logical regions: {:?}",
382 physical_region_id,
383 logical_regions
384 );
385 let physical_columns = self
386 .data_region
387 .physical_columns(physical_region_id)
388 .await?;
389 let primary_key_encoding = self
390 .mito
391 .get_primary_key_encoding(physical_region_id)
392 .context(PhysicalRegionNotFoundSnafu {
393 region_id: physical_region_id,
394 })?;
395
396 {
397 let mut state = self.state.write().unwrap();
398 let time_index_unit = physical_columns
402 .iter()
403 .find_map(|col| {
404 if col.semantic_type == SemanticType::Timestamp {
405 col.column_schema
406 .data_type
407 .as_timestamp()
408 .map(|data_type| data_type.unit())
409 } else {
410 None
411 }
412 })
413 .unwrap();
414 let physical_columns = physical_columns
415 .into_iter()
416 .map(|col| (col.column_schema.name.clone(), col))
417 .collect();
418 state.add_physical_region(
419 physical_region_id,
420 physical_columns,
421 primary_key_encoding,
422 physical_region_options,
423 time_index_unit,
424 );
425 for logical_region_id in &logical_regions {
427 state.add_logical_region(physical_region_id, *logical_region_id);
428 }
429 }
430
431 let mut opened_logical_region_ids = Vec::new();
432 for logical_region_id in logical_regions {
435 if self
436 .metadata_region
437 .open_logical_region(logical_region_id)
438 .await
439 {
440 opened_logical_region_ids.push(logical_region_id);
441 }
442 }
443
444 LOGICAL_REGION_COUNT.add(opened_logical_region_ids.len() as i64);
445
446 Ok(opened_logical_region_ids)
447 }
448}
449
450