Skip to main content

metric_engine/
engine.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
15mod 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/// # Metric Engine
71///
72/// ## Regions
73///
74/// Regions in this metric engine has several roles. There is `PhysicalRegion`,
75/// which refer to the region that actually stores the data. And `LogicalRegion`
76/// that is "simulated" over physical regions. Each logical region is associated
77/// with one physical region group, which is a group of two physical regions.
78/// Their relationship is illustrated below:
79///
80/// ```mermaid
81/// erDiagram
82///     LogicalRegion ||--o{ PhysicalRegionGroup : corresponds
83///     PhysicalRegionGroup ||--|| DataRegion : contains
84///     PhysicalRegionGroup ||--|| MetadataRegion : contains
85/// ```
86///
87/// Metric engine uses two region groups. One is for data region
88/// ([METRIC_DATA_REGION_GROUP](crate::consts::METRIC_DATA_REGION_GROUP)), and the
89/// other is for metadata region ([METRIC_METADATA_REGION_GROUP](crate::consts::METRIC_METADATA_REGION_GROUP)).
90/// From the definition of [`RegionId`], we can convert between these two physical
91/// region ids easily. Thus in the code base we usually refer to one "physical
92/// region id", and convert it to the other one when necessary.
93///
94/// The logical region, in contrast, is a virtual region. It doesn't has dedicated
95/// storage or region group. Only a region id that is allocated by meta server.
96/// And all other things is shared with other logical region that are associated
97/// with the same physical region group.
98///
99/// For more document about physical regions, please refer to [`MetadataRegion`]
100/// and [`DataRegion`].
101///
102/// ## Operations
103///
104/// Both physical and logical region are accessible to user. But the operation
105/// they support are different. List below:
106///
107/// | Operations | Logical Region | Physical Region |
108/// | ---------- | -------------- | --------------- |
109/// |   Create   |       ✅        |        ✅        |
110/// |    Drop    |       ✅        |        ❓*       |
111/// |   Write    |       ✅        |        ❌        |
112/// |    Read    |       ✅        |        ✅        |
113/// |   Close    |       ✅        |        ✅        |
114/// |    Open    |       ✅        |        ✅        |
115/// |   Alter    |       ✅        |        ❓*       |
116///
117/// *: Physical region can be dropped only when all related logical regions are dropped.
118/// *: Alter: Physical regions only support altering region options.
119///
120/// ## Internal Columns
121///
122/// The physical data region contains two internal columns. Should
123/// mention that "internal" here is for metric engine itself. Mito
124/// engine will add it's internal columns to the region as well.
125///
126/// Their column id is registered in [`ReservedColumnId`]. And column name is
127/// defined in [`DATA_SCHEMA_TSID_COLUMN_NAME`] and [`DATA_SCHEMA_TABLE_ID_COLUMN_NAME`].
128///
129/// Tsid is generated by hashing all tags. And table id is retrieved from logical region
130/// id to distinguish data from different logical tables.
131#[derive(Clone)]
132pub struct MetricEngine {
133    inner: Arc<MetricEngineInner>,
134}
135
136#[async_trait]
137impl RegionEngine for MetricEngine {
138    /// Name of this engine
139    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    /// Handles non-query request to the region. Returns the count of affected rows.
210    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    /// Retrieves region's metadata.
347    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    /// Retrieves region's disk usage.
355    ///
356    /// Note: Returns `None` if it's a logical region.
357    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    /// Stops the engine
366    async fn stop(&self) -> Result<(), BoxedError> {
367        // don't need to stop the underlying mito engine
368        Ok(())
369    }
370
371    fn set_region_role(&self, region_id: RegionId, role: RegionRole) -> Result<(), BoxedError> {
372        // ignore the region not found error
373        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    /// Returns the physical region role.
476    ///
477    /// Note: Returns `None` if it's a logical region.
478    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    /// Batch put operation for multiple logical regions.
524    /// Requests are grouped by physical region; a failure can leave earlier
525    /// physical-region groups committed.
526    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    /// Returns all logical regions associated with the physical region.
534    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    /// Handles substrait query and return a stream of record batches
542    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        // Default Mito config disables series indexes: reaching that error proves forwarding.
630        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        // This maintenance operation applies to a physical region group only.
713        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        // close physical region
729        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        // reopen physical region
739        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, // Use Bare path type for engine regions
746            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        // close nonexistent region won't report error
757        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        // open nonexistent region won't report error
767        let invalid_open_request = RegionOpenRequest {
768            engine: METRIC_ENGINE_NAME.to_string(),
769            table_dir: TestEnv::default_table_dir(),
770            path_type: PathType::Bare, // Use Bare path type for engine regions
771            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        // Delete parquet files in metadata region
884        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        // Opening an already opened region should succeed.
907        // Since the region is already open, no metadata recovery operations will be performed.
908        metric_engine
909            .handle_request(physical_region_id, RegionRequest::Open(open_request))
910            .await
911            .unwrap();
912
913        // Close the region
914        metric_engine
915            .handle_request(
916                physical_region_id,
917                RegionRequest::Close(RegionCloseRequest::default()),
918            )
919            .await
920            .unwrap();
921
922        // Try to reopen region.
923        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        // Failed to open region because of missing parquet files.
940        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        // The metadata/data region should be closed.
946        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        // Creates physical regions
982        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            // Creates logical regions for each physical region
997            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        // Closes all regions
1007        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        // Opens all regions and skip the wal
1018        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}