1mod alter;
16mod bulk_insert;
17mod catchup;
18mod close;
19mod create;
20mod drop;
21mod flush;
22mod open;
23mod options;
24mod put;
25mod read;
26mod region_metadata;
27mod staging;
28mod state;
29mod sync;
30
31use std::any::Any;
32use std::collections::HashMap;
33use std::sync::{Arc, RwLock};
34
35use api::region::RegionResponse;
36use async_trait::async_trait;
37use common_error::ext::{BoxedError, ErrorExt};
38use common_error::status_code::StatusCode;
39use common_runtime::RepeatedTask;
40use mito2::engine::MitoEngine;
41pub(crate) use options::IndexOptions;
42use snafu::{OptionExt, ResultExt};
43pub(crate) use state::MetricEngineState;
44use store_api::metadata::RegionMetadataRef;
45use store_api::metric_engine_consts::METRIC_ENGINE_NAME;
46use store_api::region_engine::{
47 BatchResponses, RegionEngine, RegionRole, RegionScannerRef, RegionStatistic,
48 RemapManifestsRequest, RemapManifestsResponse, SetRegionRoleStateResponse,
49 SetRegionRoleStateSuccess, SettableRegionRoleState, SyncRegionFromRequest,
50 SyncRegionFromResponse,
51};
52use store_api::region_request::{
53 AffectedRows, BatchRegionDdlRequest, RegionCatchupRequest, RegionOpenRequest, RegionPutRequest,
54 RegionRequest, RegionTruncateRequest,
55};
56use store_api::storage::{RegionId, ScanRequest, SequenceNumber};
57
58use crate::config::EngineConfig;
59use crate::data_region::DataRegion;
60use crate::error::{
61 self, Error, Result, StartRepeatedTaskSnafu, UnsupportedRegionRequestSnafu,
62 UnsupportedRemapManifestsRequestSnafu,
63};
64use crate::metadata_region::MetadataRegion;
65use crate::repeated_task::FlushMetadataRegionTask;
66use crate::row_modifier::RowModifier;
67use crate::utils::{self, get_region_statistic};
68
69#[cfg_attr(doc, aquamarine::aquamarine)]
70#[derive(Clone)]
132pub struct MetricEngine {
133 inner: Arc<MetricEngineInner>,
134}
135
136#[async_trait]
137impl RegionEngine for MetricEngine {
138 fn name(&self) -> &str {
140 METRIC_ENGINE_NAME
141 }
142
143 async fn handle_batch_open_requests(
144 &self,
145 parallelism: usize,
146 requests: Vec<(RegionId, RegionOpenRequest)>,
147 ) -> Result<BatchResponses, BoxedError> {
148 self.inner
149 .handle_batch_open_requests(parallelism, requests)
150 .await
151 .map_err(BoxedError::new)
152 }
153
154 async fn handle_batch_catchup_requests(
155 &self,
156 parallelism: usize,
157 requests: Vec<(RegionId, RegionCatchupRequest)>,
158 ) -> Result<BatchResponses, BoxedError> {
159 self.inner
160 .handle_batch_catchup_requests(parallelism, requests)
161 .await
162 .map_err(BoxedError::new)
163 }
164
165 async fn handle_batch_ddl_requests(
166 &self,
167 batch_request: BatchRegionDdlRequest,
168 ) -> Result<RegionResponse, BoxedError> {
169 match batch_request {
170 BatchRegionDdlRequest::Create(requests) => {
171 let mut extension_return_value = HashMap::new();
172 let rows = self
173 .inner
174 .create_regions(requests, &mut extension_return_value)
175 .await
176 .map_err(BoxedError::new)?;
177
178 Ok(RegionResponse {
179 affected_rows: rows,
180 extensions: extension_return_value,
181 metadata: Vec::new(),
182 })
183 }
184 BatchRegionDdlRequest::Alter(requests) => {
185 let mut extension_return_value = HashMap::new();
186 let rows = self
187 .inner
188 .alter_regions(requests, &mut extension_return_value)
189 .await
190 .map_err(BoxedError::new)?;
191
192 Ok(RegionResponse {
193 affected_rows: rows,
194 extensions: extension_return_value,
195 metadata: Vec::new(),
196 })
197 }
198 BatchRegionDdlRequest::Drop(requests) => {
199 self.handle_requests(
200 requests
201 .into_iter()
202 .map(|(region_id, req)| (region_id, RegionRequest::Drop(req))),
203 )
204 .await
205 }
206 }
207 }
208
209 async fn handle_request(
211 &self,
212 region_id: RegionId,
213 request: RegionRequest,
214 ) -> Result<RegionResponse, BoxedError> {
215 let mut extension_return_value = HashMap::new();
216
217 let result = match request {
218 RegionRequest::EnterStaging(_) => {
219 if self.inner.is_physical_region(region_id) {
220 self.handle_enter_staging_request(region_id, request).await
221 } else {
222 UnsupportedRegionRequestSnafu { request }.fail()
223 }
224 }
225 RegionRequest::ApplyStagingManifest(_) => {
226 if self.inner.is_physical_region(region_id) {
227 return self.inner.mito.handle_request(region_id, request).await;
228 } else {
229 UnsupportedRegionRequestSnafu { request }.fail()
230 }
231 }
232 RegionRequest::Put(put) => self.inner.put_region(region_id, put).await,
233 RegionRequest::Create(create) => {
234 self.inner
235 .create_regions(vec![(region_id, create)], &mut extension_return_value)
236 .await
237 }
238 RegionRequest::Drop(drop) => self.inner.drop_region(region_id, drop).await,
239 RegionRequest::Open(open) => self.inner.open_region(region_id, open).await,
240 RegionRequest::CleanUp(clean_up) => {
241 self.inner.clean_up_region(region_id, clean_up).await
242 }
243 RegionRequest::Close(close) => self.inner.close_region(region_id, close).await,
244 RegionRequest::Alter(alter) => {
245 self.inner
246 .alter_regions(vec![(region_id, alter)], &mut extension_return_value)
247 .await
248 }
249 RegionRequest::Compact(_) => {
250 if self.inner.is_physical_region(region_id) {
251 self.inner
252 .mito
253 .handle_request(region_id, request)
254 .await
255 .context(error::MitoFlushOperationSnafu)
256 .map(|response| response.affected_rows)
257 } else {
258 UnsupportedRegionRequestSnafu { request }.fail()
259 }
260 }
261 RegionRequest::Flush(req) => self.inner.flush_region(region_id, req).await,
262 RegionRequest::BuildIndex(_) => {
263 if self.inner.is_physical_region(region_id) {
264 self.inner
265 .mito
266 .handle_request(region_id, request)
267 .await
268 .context(error::MitoFlushOperationSnafu)
269 .map(|response| response.affected_rows)
270 } else {
271 UnsupportedRegionRequestSnafu { request }.fail()
272 }
273 }
274 RegionRequest::Truncate(RegionTruncateRequest::Unflushed) => {
275 if self.inner.is_physical_region(region_id) {
276 self.inner
277 .mito
278 .handle_request(
279 utils::to_data_region_id(region_id),
280 RegionRequest::Truncate(RegionTruncateRequest::Unflushed),
281 )
282 .await
283 .context(error::MitoTruncateOperationSnafu)
284 .map(|response| response.affected_rows)
285 } else {
286 UnsupportedRegionRequestSnafu {
287 request: RegionRequest::Truncate(RegionTruncateRequest::Unflushed),
288 }
289 .fail()
290 }
291 }
292 RegionRequest::Truncate(request) => UnsupportedRegionRequestSnafu {
293 request: RegionRequest::Truncate(request),
294 }
295 .fail(),
296 RegionRequest::Delete(delete) => self.inner.delete_region(region_id, delete).await,
297 RegionRequest::Catchup(_) => {
298 let mut response = self
299 .inner
300 .handle_batch_catchup_requests(
301 1,
302 vec![(region_id, RegionCatchupRequest::default())],
303 )
304 .await
305 .map_err(BoxedError::new)?;
306 debug_assert_eq!(response.len(), 1);
307 let (resp_region_id, response) = response
308 .pop()
309 .context(error::UnexpectedRequestSnafu {
310 reason: "expected 1 response, but got zero responses",
311 })
312 .map_err(BoxedError::new)?;
313 debug_assert_eq!(region_id, resp_region_id);
314 return response;
315 }
316 RegionRequest::BulkInserts(bulk) => {
317 self.inner.bulk_insert_region(region_id, bulk).await
318 }
319 };
320
321 result.map_err(BoxedError::new).map(|rows| RegionResponse {
322 affected_rows: rows,
323 extensions: extension_return_value,
324 metadata: Vec::new(),
325 })
326 }
327
328 async fn handle_query(
329 &self,
330 region_id: RegionId,
331 request: ScanRequest,
332 ) -> Result<RegionScannerRef, BoxedError> {
333 self.handle_query(region_id, request).await
334 }
335
336 async fn get_committed_sequence(
337 &self,
338 region_id: RegionId,
339 ) -> Result<SequenceNumber, BoxedError> {
340 self.inner
341 .get_last_seq_num(region_id)
342 .await
343 .map_err(BoxedError::new)
344 }
345
346 async fn get_metadata(&self, region_id: RegionId) -> Result<RegionMetadataRef, BoxedError> {
348 self.inner
349 .load_region_metadata(region_id)
350 .await
351 .map_err(BoxedError::new)
352 }
353
354 fn region_statistic(&self, region_id: RegionId) -> Option<RegionStatistic> {
358 if self.inner.is_physical_region(region_id) {
359 get_region_statistic(&self.inner.mito, region_id)
360 } else {
361 None
362 }
363 }
364
365 async fn stop(&self) -> Result<(), BoxedError> {
367 Ok(())
369 }
370
371 fn set_region_role(&self, region_id: RegionId, role: RegionRole) -> Result<(), BoxedError> {
372 for x in [
374 utils::to_metadata_region_id(region_id),
375 utils::to_data_region_id(region_id),
376 ] {
377 if let Err(e) = self.inner.mito.set_region_role(x, role)
378 && e.status_code() != StatusCode::RegionNotFound
379 {
380 return Err(e);
381 }
382 }
383 Ok(())
384 }
385
386 async fn sync_region(
387 &self,
388 region_id: RegionId,
389 request: SyncRegionFromRequest,
390 ) -> Result<SyncRegionFromResponse, BoxedError> {
391 match request {
392 SyncRegionFromRequest::FromManifest(manifest_info) => self
393 .inner
394 .sync_region_from_manifest(region_id, manifest_info)
395 .await
396 .map_err(BoxedError::new),
397 SyncRegionFromRequest::FromRegion {
398 source_region_id,
399 parallelism,
400 } => {
401 if self.inner.is_physical_region(region_id) {
402 self.inner
403 .sync_region_from_region(region_id, source_region_id, parallelism)
404 .await
405 .map_err(BoxedError::new)
406 } else {
407 Err(BoxedError::new(
408 error::UnsupportedSyncRegionFromRequestSnafu { region_id }.build(),
409 ))
410 }
411 }
412 }
413 }
414
415 async fn remap_manifests(
416 &self,
417 request: RemapManifestsRequest,
418 ) -> Result<RemapManifestsResponse, BoxedError> {
419 let region_id = request.region_id;
420 if self.inner.is_physical_region(region_id) {
421 self.inner.mito.remap_manifests(request).await
422 } else {
423 Err(BoxedError::new(
424 UnsupportedRemapManifestsRequestSnafu { region_id }.build(),
425 ))
426 }
427 }
428
429 async fn set_region_role_state_gracefully(
430 &self,
431 region_id: RegionId,
432 region_role_state: SettableRegionRoleState,
433 ) -> std::result::Result<SetRegionRoleStateResponse, BoxedError> {
434 let metadata_result = match self
435 .inner
436 .mito
437 .set_region_role_state_gracefully(
438 utils::to_metadata_region_id(region_id),
439 region_role_state,
440 )
441 .await?
442 {
443 SetRegionRoleStateResponse::Success(success) => success,
444 SetRegionRoleStateResponse::NotFound => {
445 return Ok(SetRegionRoleStateResponse::NotFound);
446 }
447 SetRegionRoleStateResponse::InvalidTransition(error) => {
448 return Ok(SetRegionRoleStateResponse::InvalidTransition(error));
449 }
450 };
451
452 let data_result = match self
453 .inner
454 .mito
455 .set_region_role_state_gracefully(region_id, region_role_state)
456 .await?
457 {
458 SetRegionRoleStateResponse::Success(success) => success,
459 SetRegionRoleStateResponse::NotFound => {
460 return Ok(SetRegionRoleStateResponse::NotFound);
461 }
462 SetRegionRoleStateResponse::InvalidTransition(error) => {
463 return Ok(SetRegionRoleStateResponse::InvalidTransition(error));
464 }
465 };
466
467 Ok(SetRegionRoleStateResponse::success(
468 SetRegionRoleStateSuccess::metric(
469 data_result.last_entry_id().unwrap_or_default(),
470 metadata_result.last_entry_id().unwrap_or_default(),
471 ),
472 ))
473 }
474
475 fn role(&self, region_id: RegionId) -> Option<RegionRole> {
479 if self.inner.is_physical_region(region_id) {
480 self.inner.mito.role(region_id)
481 } else {
482 None
483 }
484 }
485
486 fn as_any(&self) -> &dyn Any {
487 self
488 }
489}
490
491impl MetricEngine {
492 pub fn try_new(mito: MitoEngine, mut config: EngineConfig) -> Result<Self> {
493 let metadata_region = MetadataRegion::new(mito.clone());
494 let data_region = DataRegion::new(mito.clone());
495 let state = Arc::new(RwLock::default());
496 config.sanitize();
497 let flush_interval = config.flush_metadata_region_interval;
498 let inner = Arc::new(MetricEngineInner {
499 mito: mito.clone(),
500 metadata_region,
501 data_region,
502 state: state.clone(),
503 row_modifier: RowModifier::default(),
504 flush_task: RepeatedTask::new(
505 flush_interval,
506 Box::new(FlushMetadataRegionTask {
507 state: state.clone(),
508 mito: mito.clone(),
509 }),
510 ),
511 });
512 inner
513 .flush_task
514 .start(common_runtime::global_runtime())
515 .context(StartRepeatedTaskSnafu { name: "flush_task" })?;
516 Ok(Self { inner })
517 }
518
519 pub fn mito(&self) -> MitoEngine {
520 self.inner.mito.clone()
521 }
522
523 pub async fn put_regions_batch(
527 &self,
528 requests: impl ExactSizeIterator<Item = (RegionId, RegionPutRequest)>,
529 ) -> Result<AffectedRows> {
530 self.inner.put_regions_batch(requests).await
531 }
532
533 pub async fn logical_regions(&self, physical_region_id: RegionId) -> Result<Vec<RegionId>> {
535 self.inner
536 .metadata_region
537 .logical_regions(physical_region_id)
538 .await
539 }
540
541 async fn handle_query(
543 &self,
544 region_id: RegionId,
545 request: ScanRequest,
546 ) -> Result<RegionScannerRef, BoxedError> {
547 self.inner
548 .read_region(region_id, request)
549 .await
550 .map_err(BoxedError::new)
551 }
552
553 async fn handle_requests(
554 &self,
555 requests: impl IntoIterator<Item = (RegionId, RegionRequest)>,
556 ) -> Result<RegionResponse, BoxedError> {
557 let mut affected_rows = 0;
558 let mut extensions = HashMap::new();
559 for (region_id, request) in requests {
560 let response = self.handle_request(region_id, request).await?;
561 affected_rows += response.affected_rows;
562 extensions.extend(response.extensions);
563 }
564
565 Ok(RegionResponse {
566 affected_rows,
567 extensions,
568 metadata: Vec::new(),
569 })
570 }
571}
572
573#[cfg(test)]
574impl MetricEngine {
575 pub async fn scan_to_stream(
576 &self,
577 region_id: RegionId,
578 request: ScanRequest,
579 ) -> Result<common_recordbatch::SendableRecordBatchStream, BoxedError> {
580 self.inner.scan_to_stream(region_id, request).await
581 }
582}
583
584struct MetricEngineInner {
585 mito: MitoEngine,
586 metadata_region: MetadataRegion,
587 data_region: DataRegion,
588 state: Arc<RwLock<MetricEngineState>>,
589 row_modifier: RowModifier,
590 flush_task: RepeatedTask<Error>,
591}
592
593#[cfg(test)]
594mod test {
595 use std::assert_matches;
596 use std::collections::HashMap;
597
598 use api::v1::Rows;
599 use common_recordbatch::RecordBatches;
600 use common_telemetry::info;
601 use common_wal::options::{KafkaWalOptions, WalOptions};
602 use mito2::sst::location::region_dir_from_table_dir;
603 use mito2::test_util::{kafka_log_store_factory, prepare_test_for_kafka_log_store};
604 use store_api::metric_engine_consts::PHYSICAL_TABLE_METADATA_KEY;
605 use store_api::mito_engine_options::WAL_OPTIONS_KEY;
606 use store_api::region_request::{
607 PathType, RegionCleanUpRequest, RegionCloseRequest, RegionDropRequest, RegionFlushRequest,
608 RegionOpenRequest, RegionPutRequest, RegionRequest, RegionTruncateRequest,
609 };
610
611 use super::*;
612 use crate::maybe_skip_kafka_log_store_integration_test;
613 use crate::test_util::{
614 TestEnv, build_rows, create_logical_region_request, row_schema_with_tags,
615 };
616
617 #[tokio::test]
618 async fn test_build_series_index_forwarding() {
619 use api::v1::region::build_index_request;
620 use store_api::region_request::RegionBuildIndexRequest;
621
622 let env = TestEnv::new().await;
623 env.init_metric_region().await;
624 let request = || {
625 RegionRequest::BuildIndex(RegionBuildIndexRequest {
626 options: Some(build_index_request::Options::SeriesIndex(Default::default())),
627 })
628 };
629 let error = env
631 .metric()
632 .handle_request(env.default_physical_region_id(), request())
633 .await
634 .unwrap_err();
635 let error = format!("{error:?}");
636 assert!(error.contains("series index is disabled"), "{error}");
637 let error = env
638 .metric()
639 .handle_request(env.default_logical_region_id(), request())
640 .await
641 .unwrap_err();
642 assert_eq!(
643 common_error::status_code::StatusCode::Unsupported,
644 error.status_code()
645 );
646 }
647
648 #[tokio::test]
649 async fn test_discard_unflushed_data_only() {
650 let env = TestEnv::new().await;
651 env.init_metric_region().await;
652
653 let engine = env.metric();
654 let mito = env.mito();
655 let physical_region_id = env.default_physical_region_id();
656 let logical_region_id = env.default_logical_region_id();
657 let data_region_id = utils::to_data_region_id(physical_region_id);
658 let metadata_region_id = utils::to_metadata_region_id(physical_region_id);
659
660 let metadata_memtable_size = mito
661 .region_statistic(metadata_region_id)
662 .unwrap()
663 .memtable_size;
664 assert!(metadata_memtable_size > 0);
665
666 engine
667 .handle_request(
668 logical_region_id,
669 RegionRequest::Put(RegionPutRequest {
670 skip_wal: false,
671 rows: Rows {
672 schema: row_schema_with_tags(&["job"]),
673 rows: build_rows(1, 5),
674 },
675 hint: None,
676 partition_expr_version: None,
677 }),
678 )
679 .await
680 .unwrap();
681 assert!(mito.region_statistic(data_region_id).unwrap().memtable_size > 0);
682
683 engine
684 .handle_request(
685 physical_region_id,
686 RegionRequest::Truncate(RegionTruncateRequest::Unflushed),
687 )
688 .await
689 .unwrap();
690
691 assert_eq!(
692 0,
693 mito.region_statistic(data_region_id).unwrap().memtable_size
694 );
695 assert_eq!(
696 metadata_memtable_size,
697 mito.region_statistic(metadata_region_id)
698 .unwrap()
699 .memtable_size
700 );
701
702 let stream = engine
703 .scan_to_stream(logical_region_id, ScanRequest::default())
704 .await
705 .unwrap();
706 let batches = RecordBatches::try_collect(stream).await.unwrap();
707 assert_eq!(
708 0,
709 batches.iter().map(|batch| batch.num_rows()).sum::<usize>()
710 );
711
712 engine
714 .handle_request(
715 logical_region_id,
716 RegionRequest::Truncate(RegionTruncateRequest::Unflushed),
717 )
718 .await
719 .unwrap_err();
720 }
721
722 #[tokio::test]
723 async fn close_open_regions() {
724 let env = TestEnv::new().await;
725 env.init_metric_region().await;
726 let engine = env.metric();
727
728 let physical_region_id = env.default_physical_region_id();
730 engine
731 .handle_request(
732 physical_region_id,
733 RegionRequest::Close(RegionCloseRequest::default()),
734 )
735 .await
736 .unwrap();
737
738 let physical_region_option = [(PHYSICAL_TABLE_METADATA_KEY.to_string(), String::new())]
740 .into_iter()
741 .collect();
742 let open_request = RegionOpenRequest {
743 engine: METRIC_ENGINE_NAME.to_string(),
744 table_dir: TestEnv::default_table_dir(),
745 path_type: PathType::Bare, options: physical_region_option,
747 skip_wal_replay: false,
748 checkpoint: None,
749 requirements: Default::default(),
750 };
751 engine
752 .handle_request(physical_region_id, RegionRequest::Open(open_request))
753 .await
754 .unwrap();
755
756 let nonexistent_region_id = RegionId::new(12313, 12);
758 engine
759 .handle_request(
760 nonexistent_region_id,
761 RegionRequest::Close(RegionCloseRequest::default()),
762 )
763 .await
764 .unwrap();
765
766 let invalid_open_request = RegionOpenRequest {
768 engine: METRIC_ENGINE_NAME.to_string(),
769 table_dir: TestEnv::default_table_dir(),
770 path_type: PathType::Bare, options: HashMap::new(),
772 skip_wal_replay: false,
773 checkpoint: None,
774 requirements: Default::default(),
775 };
776 engine
777 .handle_request(
778 nonexistent_region_id,
779 RegionRequest::Open(invalid_open_request),
780 )
781 .await
782 .unwrap();
783 }
784
785 #[tokio::test]
786 async fn test_offline_cleanup_physical_region() {
787 let env = TestEnv::new().await;
788 env.init_metric_region().await;
789 let engine = env.metric();
790 let mito = env.mito();
791 let physical_region_id = env.default_physical_region_id();
792 let metadata_region_id = crate::utils::to_metadata_region_id(physical_region_id);
793 let data_region_id = crate::utils::to_data_region_id(physical_region_id);
794
795 engine
796 .handle_request(
797 physical_region_id,
798 RegionRequest::Close(RegionCloseRequest::default()),
799 )
800 .await
801 .unwrap();
802
803 let object_store = env.get_object_store().unwrap();
804 let metadata_region_dir = region_dir_from_table_dir(
805 &TestEnv::default_table_dir(),
806 metadata_region_id,
807 PathType::Metadata,
808 );
809 let data_region_dir = region_dir_from_table_dir(
810 &TestEnv::default_table_dir(),
811 data_region_id,
812 PathType::Data,
813 );
814 assert!(object_store.exists(&metadata_region_dir).await.unwrap());
815 assert!(object_store.exists(&data_region_dir).await.unwrap());
816
817 let physical_region_option = [(PHYSICAL_TABLE_METADATA_KEY.to_string(), String::new())]
818 .into_iter()
819 .collect();
820 let clean_up_request = RegionCleanUpRequest {
821 engine: METRIC_ENGINE_NAME.to_string(),
822 table_dir: TestEnv::default_table_dir(),
823 path_type: PathType::Bare,
824 options: physical_region_option,
825 };
826 engine
827 .handle_request(physical_region_id, RegionRequest::CleanUp(clean_up_request))
828 .await
829 .unwrap();
830
831 assert!(!mito.is_region_exists(metadata_region_id));
832 assert!(!mito.is_region_exists(data_region_id));
833 assert!(!object_store.exists(&metadata_region_dir).await.unwrap());
834 assert!(!object_store.exists(&data_region_dir).await.unwrap());
835 }
836
837 #[tokio::test]
838 async fn test_role() {
839 let env = TestEnv::new().await;
840 env.init_metric_region().await;
841
842 let logical_region_id = env.default_logical_region_id();
843 let physical_region_id = env.default_physical_region_id();
844
845 assert!(env.metric().role(logical_region_id).is_none());
846 assert!(env.metric().role(physical_region_id).is_some());
847 }
848
849 #[tokio::test]
850 async fn test_region_disk_usage() {
851 let env = TestEnv::new().await;
852 env.init_metric_region().await;
853
854 let logical_region_id = env.default_logical_region_id();
855 let physical_region_id = env.default_physical_region_id();
856
857 assert!(env.metric().region_statistic(logical_region_id).is_none());
858 assert!(env.metric().region_statistic(physical_region_id).is_some());
859 }
860
861 #[tokio::test]
862 async fn test_open_region_failure() {
863 let env = TestEnv::new().await;
864 env.init_metric_region().await;
865 let physical_region_id = env.default_physical_region_id();
866
867 let metric_engine = env.metric();
868 metric_engine
869 .handle_request(
870 physical_region_id,
871 RegionRequest::Flush(RegionFlushRequest::default()),
872 )
873 .await
874 .unwrap();
875
876 let path = region_dir_from_table_dir(
877 &TestEnv::default_table_dir(),
878 physical_region_id,
879 PathType::Metadata,
880 );
881 let object_store = env.get_object_store().unwrap();
882 let list = object_store.list(&path).await.unwrap();
883 for entry in list {
885 if entry.metadata().is_dir() {
886 continue;
887 }
888 if entry.name().ends_with("parquet") {
889 info!("deleting {}", entry.path());
890 object_store.delete(entry.path()).await.unwrap();
891 }
892 }
893
894 let physical_region_option = [(PHYSICAL_TABLE_METADATA_KEY.to_string(), String::new())]
895 .into_iter()
896 .collect();
897 let open_request = RegionOpenRequest {
898 engine: METRIC_ENGINE_NAME.to_string(),
899 table_dir: TestEnv::default_table_dir(),
900 path_type: PathType::Bare,
901 options: physical_region_option,
902 skip_wal_replay: false,
903 checkpoint: None,
904 requirements: Default::default(),
905 };
906 metric_engine
909 .handle_request(physical_region_id, RegionRequest::Open(open_request))
910 .await
911 .unwrap();
912
913 metric_engine
915 .handle_request(
916 physical_region_id,
917 RegionRequest::Close(RegionCloseRequest::default()),
918 )
919 .await
920 .unwrap();
921
922 let physical_region_option = [(PHYSICAL_TABLE_METADATA_KEY.to_string(), String::new())]
924 .into_iter()
925 .collect();
926 let open_request = RegionOpenRequest {
927 engine: METRIC_ENGINE_NAME.to_string(),
928 table_dir: TestEnv::default_table_dir(),
929 path_type: PathType::Bare,
930 options: physical_region_option,
931 skip_wal_replay: false,
932 checkpoint: None,
933 requirements: Default::default(),
934 };
935 let err = metric_engine
936 .handle_request(physical_region_id, RegionRequest::Open(open_request))
937 .await
938 .unwrap_err();
939 assert_eq!(err.status_code(), StatusCode::StorageUnavailable);
941
942 let mito_engine = metric_engine.mito();
943 let data_region_id = utils::to_data_region_id(physical_region_id);
944 let metadata_region_id = utils::to_metadata_region_id(physical_region_id);
945 let err = mito_engine.get_metadata(data_region_id).await.unwrap_err();
947 assert_eq!(err.status_code(), StatusCode::RegionNotFound);
948 let err = mito_engine
949 .get_metadata(metadata_region_id)
950 .await
951 .unwrap_err();
952 assert_eq!(err.status_code(), StatusCode::RegionNotFound);
953 }
954
955 #[tokio::test]
956 async fn test_catchup_regions() {
957 common_telemetry::init_default_ut_logging();
958 maybe_skip_kafka_log_store_integration_test!();
959 let kafka_log_store_factory = kafka_log_store_factory().unwrap();
960 let mito_env = mito2::test_util::TestEnv::new()
961 .await
962 .with_log_store_factory(kafka_log_store_factory.clone());
963 let env = TestEnv::with_mito_env(mito_env).await;
964 let table_dir = |region_id| format!("table/{region_id}");
965 let mut physical_region_ids = vec![];
966 let mut logical_region_ids = vec![];
967
968 let num_topics = 3;
969 let num_physical_regions = 8;
970 let num_logical_regions = 16;
971 let parallelism = 2;
972 let mut topics = Vec::with_capacity(num_topics);
973 for _ in 0..num_topics {
974 let topic = prepare_test_for_kafka_log_store(&kafka_log_store_factory)
975 .await
976 .unwrap();
977 topics.push(topic);
978 }
979
980 let topic_idx = |id| (id as usize) % num_topics;
981 for i in 0..num_physical_regions {
983 let physical_region_id = RegionId::new(1, i);
984 physical_region_ids.push(physical_region_id);
985
986 let wal_options = WalOptions::Kafka(KafkaWalOptions::new(topics[topic_idx(i)].clone()));
987 env.create_physical_region(
988 physical_region_id,
989 &table_dir(physical_region_id),
990 vec![(
991 WAL_OPTIONS_KEY.to_string(),
992 serde_json::to_string(&wal_options).unwrap(),
993 )],
994 )
995 .await;
996 for j in 0..num_logical_regions {
998 let logical_region_id = RegionId::new(1024 + i, j);
999 logical_region_ids.push(logical_region_id);
1000 env.create_logical_region(physical_region_id, logical_region_id)
1001 .await;
1002 }
1003 }
1004
1005 let metric_engine = env.metric();
1006 for region_id in logical_region_ids.iter().chain(physical_region_ids.iter()) {
1008 metric_engine
1009 .handle_request(
1010 *region_id,
1011 RegionRequest::Close(RegionCloseRequest::default()),
1012 )
1013 .await
1014 .unwrap();
1015 }
1016
1017 let requests = physical_region_ids
1019 .iter()
1020 .enumerate()
1021 .map(|(idx, region_id)| {
1022 let mut options = HashMap::new();
1023 let wal_options =
1024 WalOptions::Kafka(KafkaWalOptions::new(topics[topic_idx(idx as u32)].clone()));
1025 options.insert(PHYSICAL_TABLE_METADATA_KEY.to_string(), String::new());
1026 options.insert(
1027 WAL_OPTIONS_KEY.to_string(),
1028 serde_json::to_string(&wal_options).unwrap(),
1029 );
1030 (
1031 *region_id,
1032 RegionOpenRequest {
1033 engine: METRIC_ENGINE_NAME.to_string(),
1034 table_dir: table_dir(*region_id),
1035 path_type: PathType::Bare,
1036 options: options.clone(),
1037 skip_wal_replay: true,
1038 checkpoint: None,
1039 requirements: Default::default(),
1040 },
1041 )
1042 })
1043 .collect::<Vec<_>>();
1044 info!("Open batch regions with parallelism: {parallelism}");
1045 metric_engine
1046 .handle_batch_open_requests(parallelism, requests)
1047 .await
1048 .unwrap();
1049 {
1050 let state = metric_engine.inner.state.read().unwrap();
1051 for logical_region in &logical_region_ids {
1052 assert!(!state.logical_regions().contains_key(logical_region));
1053 }
1054 }
1055
1056 let catch_requests = physical_region_ids
1057 .iter()
1058 .map(|region_id| {
1059 (
1060 *region_id,
1061 RegionCatchupRequest {
1062 set_writable: true,
1063 ..Default::default()
1064 },
1065 )
1066 })
1067 .collect::<Vec<_>>();
1068 metric_engine
1069 .handle_batch_catchup_requests(parallelism, catch_requests)
1070 .await
1071 .unwrap();
1072 {
1073 let state = metric_engine.inner.state.read().unwrap();
1074 for logical_region in &logical_region_ids {
1075 assert!(state.logical_regions().contains_key(logical_region));
1076 }
1077 }
1078 }
1079
1080 #[tokio::test]
1081 async fn test_drop_region() {
1082 let env = TestEnv::new().await;
1083 let engine = env.metric();
1084 let physical_region_id1 = RegionId::new(1024, 0);
1085 let logical_region_id1 = RegionId::new(1025, 0);
1086 env.create_physical_region(physical_region_id1, "/test_dir1", vec![])
1087 .await;
1088 let region_create_request1 =
1089 create_logical_region_request(&["job"], physical_region_id1, "logical1");
1090 engine
1091 .handle_batch_ddl_requests(BatchRegionDdlRequest::Create(vec![(
1092 logical_region_id1,
1093 region_create_request1,
1094 )]))
1095 .await
1096 .unwrap();
1097 let err = engine
1098 .handle_request(
1099 physical_region_id1,
1100 RegionRequest::Drop(RegionDropRequest {
1101 fast_path: false,
1102 force: false,
1103 partial_drop: false,
1104 }),
1105 )
1106 .await
1107 .unwrap_err();
1108 assert_matches!(
1109 err.as_any().downcast_ref::<Error>().unwrap(),
1110 &Error::PhysicalRegionBusy { .. }
1111 );
1112
1113 engine
1114 .handle_request(
1115 physical_region_id1,
1116 RegionRequest::Drop(RegionDropRequest {
1117 fast_path: false,
1118 force: true,
1119 partial_drop: false,
1120 }),
1121 )
1122 .await
1123 .unwrap();
1124 assert!(
1125 engine
1126 .inner
1127 .state
1128 .read()
1129 .unwrap()
1130 .physical_region_states()
1131 .get(&physical_region_id1)
1132 .is_none()
1133 );
1134 assert!(
1135 engine
1136 .inner
1137 .state
1138 .read()
1139 .unwrap()
1140 .logical_regions()
1141 .get(&logical_region_id1)
1142 .is_none()
1143 );
1144 }
1145}