Skip to main content

mito2/
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
15//! Mito region engine.
16
17#[cfg(test)]
18mod alter_test;
19#[cfg(test)]
20mod append_mode_test;
21#[cfg(test)]
22mod basic_test;
23#[cfg(test)]
24mod batch_catchup_test;
25#[cfg(test)]
26mod batch_open_test;
27#[cfg(test)]
28mod bump_committed_sequence_test;
29#[cfg(test)]
30mod catchup_test;
31#[cfg(test)]
32mod close_test;
33#[cfg(test)]
34pub(crate) mod compaction_test;
35#[cfg(test)]
36mod create_test;
37#[cfg(test)]
38mod drop_test;
39#[cfg(test)]
40mod edit_region_test;
41#[cfg(test)]
42mod file_ref_test;
43#[cfg(test)]
44mod filter_deleted_test;
45#[cfg(test)]
46mod flush_test;
47#[cfg(test)]
48mod index_build_test;
49#[cfg(any(test, feature = "test"))]
50pub mod listener;
51#[cfg(test)]
52mod merge_mode_test;
53#[cfg(test)]
54mod open_test;
55#[cfg(test)]
56mod parallel_test;
57#[cfg(test)]
58mod projection_test;
59#[cfg(test)]
60mod prune_test;
61pub mod region_hook;
62#[cfg(test)]
63mod row_selector_test;
64#[cfg(test)]
65mod scan_corrupt;
66#[cfg(test)]
67mod scan_test;
68#[cfg(test)]
69mod set_role_state_test;
70#[cfg(test)]
71mod skip_wal_test;
72#[cfg(test)]
73mod staging_test;
74#[cfg(test)]
75mod sync_test;
76#[cfg(test)]
77mod truncate_test;
78
79#[cfg(test)]
80mod copy_region_from_test;
81#[cfg(test)]
82mod remap_manifests_test;
83
84#[cfg(test)]
85mod apply_staging_manifest_test;
86#[cfg(test)]
87mod partition_filter_test;
88mod puffin_index;
89
90use std::any::Any;
91use std::collections::{HashMap, HashSet};
92use std::sync::Arc;
93use std::time::Instant;
94
95use api::region::RegionResponse;
96use async_trait::async_trait;
97use common_base::Plugins;
98use common_error::ext::BoxedError;
99use common_meta::error::UnexpectedSnafu;
100use common_meta::key::SchemaMetadataManagerRef;
101use common_recordbatch::{QueryMemoryTracker, SendableRecordBatchStream};
102use common_stat::get_total_memory_bytes;
103use common_telemetry::{info, tracing, warn};
104use common_wal::options::WalOptions;
105use futures::future::{join_all, try_join_all};
106use futures::stream::{self, Stream, StreamExt};
107use object_store::manager::ObjectStoreManagerRef;
108use region_hook::RegionHookRef;
109use snafu::{OptionExt, ResultExt, ensure};
110use store_api::ManifestVersion;
111use store_api::codec::PrimaryKeyEncoding;
112use store_api::logstore::LogStore;
113use store_api::logstore::provider::{KafkaProvider, Provider};
114use store_api::metadata::{ColumnMetadata, RegionMetadataRef};
115use store_api::metric_engine_consts::{
116    MANIFEST_INFO_EXTENSION_KEY, TABLE_COLUMN_METADATA_EXTENSION_KEY,
117};
118use store_api::region_engine::{
119    BatchResponses, MitoCopyRegionFromRequest, MitoCopyRegionFromResponse, RegionEngine,
120    RegionManifestInfo, RegionRole, RegionScannerRef, RegionStatistic, RemapManifestsRequest,
121    RemapManifestsResponse, SetRegionRoleStateResponse, SettableRegionRoleState,
122    SyncRegionFromRequest, SyncRegionFromResponse,
123};
124use store_api::region_info::RegionInfoEntry;
125use store_api::region_request::{
126    AffectedRows, RegionCatchupRequest, RegionOpenRequest, RegionRequest,
127};
128use store_api::sst_entry::{ManifestSstEntry, PuffinIndexMetaEntry, StorageSstEntry};
129use store_api::storage::{FileId, FileRefsManifest, RegionId, ScanRequest, SequenceNumber};
130use tokio::sync::{Semaphore, oneshot};
131
132use crate::access_layer::RegionFilePathFactory;
133use crate::cache::{CacheManagerRef, CacheStrategy};
134use crate::config::MitoConfig;
135use crate::engine::puffin_index::{IndexEntryContext, collect_index_entries_from_puffin};
136use crate::error::{
137    IncrementalQueryStaleSnafu, InvalidRequestSnafu, JoinSnafu, MitoManifestInfoSnafu, RecvSnafu,
138    RegionNotFoundSnafu, Result, SerdeJsonSnafu, SerializeColumnMetadataSnafu,
139    SnapshotFenceStaleSnafu,
140};
141#[cfg(feature = "enterprise")]
142use crate::extension::BoxedExtensionRangeProviderFactory;
143use crate::gc::GcLimiterRef;
144use crate::manifest::action::RegionEdit;
145use crate::memtable::MemtableStats;
146use crate::metrics::{
147    HANDLE_REQUEST_ELAPSED, SCAN_MEMORY_EXHAUSTED_TOTAL, SCAN_MEMORY_USAGE_BYTES,
148    SCAN_REQUESTS_REJECTED_TOTAL,
149};
150use crate::read::scan_region::{ScanRegion, Scanner};
151use crate::read::stream::ScanBatchStream;
152use crate::region::MitoRegionRef;
153use crate::region::opener::PartitionExprFetcherRef;
154use crate::region::options::parse_wal_options;
155use crate::request::{RegionEditRequest, WorkerRequest};
156use crate::sst::file::{FileMeta, RegionFileId, RegionIndexId};
157use crate::sst::file_ref::FileReferenceManagerRef;
158use crate::sst::index::intermediate::IntermediateManager;
159use crate::sst::index::puffin_manager::PuffinManagerFactory;
160use crate::wal::entry_distributor::{
161    DEFAULT_ENTRY_RECEIVER_BUFFER_SIZE, build_wal_entry_distributor_and_receivers,
162};
163use crate::wal::raw_entry_reader::{LogStoreRawEntryReader, RawEntryReader};
164use crate::worker::WorkerGroup;
165
166pub const MITO_ENGINE_NAME: &str = "mito";
167
168pub struct MitoEngineBuilder<'a, S: LogStore> {
169    data_home: &'a str,
170    config: MitoConfig,
171    log_store: Arc<S>,
172    object_store_manager: ObjectStoreManagerRef,
173    schema_metadata_manager: SchemaMetadataManagerRef,
174    file_ref_manager: FileReferenceManagerRef,
175    partition_expr_fetcher: PartitionExprFetcherRef,
176    plugins: Plugins,
177    #[cfg(feature = "enterprise")]
178    extension_range_provider_factory: Option<BoxedExtensionRangeProviderFactory>,
179}
180
181impl<'a, S: LogStore> MitoEngineBuilder<'a, S> {
182    #[allow(clippy::too_many_arguments)]
183    pub fn new(
184        data_home: &'a str,
185        config: MitoConfig,
186        log_store: Arc<S>,
187        object_store_manager: ObjectStoreManagerRef,
188        schema_metadata_manager: SchemaMetadataManagerRef,
189        file_ref_manager: FileReferenceManagerRef,
190        partition_expr_fetcher: PartitionExprFetcherRef,
191        plugins: Plugins,
192    ) -> Self {
193        Self {
194            data_home,
195            config,
196            log_store,
197            object_store_manager,
198            schema_metadata_manager,
199            file_ref_manager,
200            plugins,
201            partition_expr_fetcher,
202            #[cfg(feature = "enterprise")]
203            extension_range_provider_factory: None,
204        }
205    }
206
207    #[cfg(feature = "enterprise")]
208    #[must_use]
209    pub fn with_extension_range_provider_factory(
210        self,
211        extension_range_provider_factory: Option<BoxedExtensionRangeProviderFactory>,
212    ) -> Self {
213        Self {
214            extension_range_provider_factory,
215            ..self
216        }
217    }
218
219    pub async fn try_build(mut self) -> Result<MitoEngine> {
220        self.config.sanitize(self.data_home)?;
221
222        let config = Arc::new(self.config);
223        // Extract the region hook before `plugins` is moved into the WorkerGroup,
224        // so the engine (and thus the GC worker) can fire `on_region_gc`.
225        let region_hook = self.plugins.get::<RegionHookRef>();
226        let workers = WorkerGroup::start(
227            config.clone(),
228            self.log_store.clone(),
229            self.object_store_manager,
230            self.schema_metadata_manager,
231            self.file_ref_manager,
232            self.partition_expr_fetcher.clone(),
233            self.plugins,
234        )
235        .await?;
236        let wal_raw_entry_reader = Arc::new(LogStoreRawEntryReader::new(self.log_store));
237        let total_memory = get_total_memory_bytes().max(0) as u64;
238        let scan_memory_limit = config.scan_memory_limit.resolve(total_memory) as usize;
239        let scan_memory_tracker =
240            QueryMemoryTracker::builder(scan_memory_limit, config.scan_memory_on_exhausted)
241                .on_update(|usage| {
242                    SCAN_MEMORY_USAGE_BYTES.set(usage as i64);
243                })
244                .on_exhausted(|| {
245                    SCAN_MEMORY_EXHAUSTED_TOTAL.inc();
246                })
247                .on_reject(|| {
248                    SCAN_REQUESTS_REJECTED_TOTAL.inc();
249                })
250                .build();
251
252        let inner = EngineInner {
253            workers,
254            config,
255            wal_raw_entry_reader,
256            scan_memory_tracker,
257            region_hook,
258            #[cfg(feature = "enterprise")]
259            extension_range_provider_factory: None,
260        };
261
262        #[cfg(feature = "enterprise")]
263        let inner =
264            inner.with_extension_range_provider_factory(self.extension_range_provider_factory);
265
266        Ok(MitoEngine {
267            inner: Arc::new(inner),
268        })
269    }
270}
271
272/// Region engine implementation for timeseries data.
273#[derive(Clone)]
274pub struct MitoEngine {
275    inner: Arc<EngineInner>,
276}
277
278impl MitoEngine {
279    /// Returns a new [MitoEngine] with specific `config`, `log_store` and `object_store`.
280    #[allow(clippy::too_many_arguments)]
281    pub async fn new<S: LogStore>(
282        data_home: &str,
283        config: MitoConfig,
284        log_store: Arc<S>,
285        object_store_manager: ObjectStoreManagerRef,
286        schema_metadata_manager: SchemaMetadataManagerRef,
287        file_ref_manager: FileReferenceManagerRef,
288        partition_expr_fetcher: PartitionExprFetcherRef,
289        plugins: Plugins,
290    ) -> Result<MitoEngine> {
291        let builder = MitoEngineBuilder::new(
292            data_home,
293            config,
294            log_store,
295            object_store_manager,
296            schema_metadata_manager,
297            file_ref_manager,
298            partition_expr_fetcher,
299            plugins,
300        );
301        builder.try_build().await
302    }
303
304    pub fn mito_config(&self) -> &MitoConfig {
305        &self.inner.config
306    }
307
308    pub fn cache_manager(&self) -> CacheManagerRef {
309        self.inner.workers.cache_manager()
310    }
311
312    pub fn file_ref_manager(&self) -> FileReferenceManagerRef {
313        self.inner.workers.file_ref_manager()
314    }
315
316    pub fn gc_limiter(&self) -> GcLimiterRef {
317        self.inner.workers.gc_limiter()
318    }
319
320    pub fn object_store_manager(&self) -> &ObjectStoreManagerRef {
321        self.inner.workers.object_store_manager()
322    }
323
324    pub fn puffin_manager_factory(&self) -> &PuffinManagerFactory {
325        self.inner.workers.puffin_manager_factory()
326    }
327
328    pub fn intermediate_manager(&self) -> &IntermediateManager {
329        self.inner.workers.intermediate_manager()
330    }
331
332    pub fn schema_metadata_manager(&self) -> &SchemaMetadataManagerRef {
333        self.inner.workers.schema_metadata_manager()
334    }
335
336    /// Returns the registered region hook (if any), for the GC worker to fire
337    /// [`RegionHook::on_region_gc`].
338    pub fn region_hook(&self) -> Option<RegionHookRef> {
339        self.inner.region_hook.clone()
340    }
341
342    /// Get all tmp ref files for given region ids, excluding files that's already in manifest.
343    pub async fn get_snapshot_of_file_refs(
344        &self,
345        file_handle_regions: impl IntoIterator<Item = RegionId>,
346        related_regions: HashMap<RegionId, HashSet<RegionId>>,
347    ) -> Result<FileRefsManifest> {
348        let file_ref_mgr = self.file_ref_manager();
349
350        let file_handle_regions = file_handle_regions.into_iter().collect::<Vec<_>>();
351        // Convert region IDs to MitoRegionRef objects, ignore regions that do not exist on current datanode
352        // as regions on other datanodes are not managed by this engine.
353        let query_regions: Vec<MitoRegionRef> = file_handle_regions
354            .into_iter()
355            .filter_map(|region_id| self.find_region(region_id))
356            .collect();
357
358        let dst_region_to_src_regions: Vec<(MitoRegionRef, HashSet<RegionId>)> = {
359            let dst2src = related_regions
360                .into_iter()
361                .flat_map(|(src, dsts)| dsts.into_iter().map(move |dst| (dst, src)))
362                .fold(
363                    HashMap::<RegionId, HashSet<RegionId>>::new(),
364                    |mut acc, (k, v)| {
365                        let entry = acc.entry(k).or_default();
366                        entry.insert(v);
367                        acc
368                    },
369                );
370            let mut dst_region_to_src_regions = Vec::with_capacity(dst2src.len());
371            for (dst_region, srcs) in dst2src {
372                let Some(region) = self.find_region(dst_region) else {
373                    return RegionNotFoundSnafu {
374                        region_id: dst_region,
375                    }
376                    .fail();
377                };
378                dst_region_to_src_regions.push((region, srcs));
379            }
380            dst_region_to_src_regions
381        };
382
383        file_ref_mgr
384            .get_snapshot_of_file_refs(query_regions, dst_region_to_src_regions)
385            .await
386    }
387
388    /// Returns true if the specific region exists.
389    pub fn is_region_exists(&self, region_id: RegionId) -> bool {
390        self.inner.workers.is_region_exists(region_id)
391    }
392
393    /// Returns true if the specific region exists.
394    pub fn is_region_opening(&self, region_id: RegionId) -> bool {
395        self.inner.workers.is_region_opening(region_id)
396    }
397
398    /// Returns true if the specific region is catching up.
399    pub fn is_region_catching_up(&self, region_id: RegionId) -> bool {
400        self.inner.workers.is_region_catching_up(region_id)
401    }
402
403    /// Returns the region disk/memory statistic.
404    pub fn get_region_statistic(&self, region_id: RegionId) -> Option<RegionStatistic> {
405        self.find_region(region_id)
406            .map(|region| region.region_statistic())
407    }
408
409    /// Returns primary key encoding of the region.
410    pub fn get_primary_key_encoding(&self, region_id: RegionId) -> Option<PrimaryKeyEncoding> {
411        self.find_region(region_id)
412            .map(|r| r.primary_key_encoding())
413    }
414
415    /// Handle substrait query and return a stream of record batches
416    ///
417    /// Notice that the output stream's ordering is not guranateed. If order
418    /// matter, please use [`scanner`] to build a [`Scanner`] to consume.
419    #[tracing::instrument(skip_all)]
420    pub async fn scan_to_stream(
421        &self,
422        region_id: RegionId,
423        request: ScanRequest,
424    ) -> Result<SendableRecordBatchStream, BoxedError> {
425        self.scanner(region_id, request)
426            .await
427            .map_err(BoxedError::new)?
428            .scan()
429            .await
430    }
431
432    /// Scan [`Batch`]es by [`ScanRequest`].
433    pub async fn scan_batch(
434        &self,
435        region_id: RegionId,
436        request: ScanRequest,
437        filter_deleted: bool,
438    ) -> Result<ScanBatchStream> {
439        let mut scan_region = self.scan_region(region_id, request)?;
440        scan_region.set_filter_deleted(filter_deleted);
441        scan_region.scanner().await?.scan_batch()
442    }
443
444    /// Returns a scanner to scan for `request`.
445    pub(crate) async fn scanner(
446        &self,
447        region_id: RegionId,
448        request: ScanRequest,
449    ) -> Result<Scanner> {
450        self.scan_region(region_id, request)?.scanner().await
451    }
452
453    /// Scans a region.
454    #[tracing::instrument(skip_all, fields(region_id = %region_id))]
455    fn scan_region(&self, region_id: RegionId, request: ScanRequest) -> Result<ScanRegion> {
456        self.inner.scan_region(region_id, request)
457    }
458
459    /// Edit region's metadata by [RegionEdit] directly. Use with care.
460    /// Now we only allow adding files or removing files from region (the [RegionEdit] struct can only contain a non-empty "files_to_add" or "files_to_remove" field).
461    /// Other region editing intention will result in an "invalid request" error.
462    /// Also note that if a region is to be edited directly, we MUST not write data to it thereafter.
463    pub async fn edit_region(&self, region_id: RegionId, edit: RegionEdit) -> Result<()> {
464        let _timer = HANDLE_REQUEST_ELAPSED
465            .with_label_values(&["edit_region"])
466            .start_timer();
467
468        ensure!(
469            is_valid_region_edit(&edit),
470            InvalidRequestSnafu {
471                region_id,
472                reason: "invalid region edit"
473            }
474        );
475
476        let (tx, rx) = oneshot::channel();
477        let request = WorkerRequest::EditRegion(RegionEditRequest::new(region_id, edit, true, tx));
478        self.inner
479            .workers
480            .submit_to_worker(region_id, request)
481            .await?;
482        rx.await.context(RecvSnafu)?
483    }
484
485    /// Handles copy region from request.
486    ///
487    /// This method is only supported for internal use and is not exposed in the trait implementation.
488    pub async fn copy_region_from(
489        &self,
490        region_id: RegionId,
491        request: MitoCopyRegionFromRequest,
492    ) -> Result<MitoCopyRegionFromResponse> {
493        self.inner.copy_region_from(region_id, request).await
494    }
495
496    #[cfg(test)]
497    pub(crate) fn get_region(&self, id: RegionId) -> Option<crate::region::MitoRegionRef> {
498        self.find_region(id)
499    }
500
501    pub fn find_region(&self, region_id: RegionId) -> Option<MitoRegionRef> {
502        self.inner.workers.get_region(region_id)
503    }
504
505    /// Returns all regions.
506    pub fn regions(&self) -> Vec<MitoRegionRef> {
507        self.inner.workers.all_regions().collect()
508    }
509
510    fn encode_manifest_info_to_extensions(
511        region_id: &RegionId,
512        manifest_info: RegionManifestInfo,
513        extensions: &mut HashMap<String, Vec<u8>>,
514    ) -> Result<()> {
515        let region_manifest_info = vec![(*region_id, manifest_info)];
516
517        extensions.insert(
518            MANIFEST_INFO_EXTENSION_KEY.to_string(),
519            RegionManifestInfo::encode_list(&region_manifest_info).context(SerdeJsonSnafu)?,
520        );
521        info!(
522            "Added manifest info: {:?} to extensions, region_id: {:?}",
523            region_manifest_info, region_id
524        );
525        Ok(())
526    }
527
528    fn encode_column_metadatas_to_extensions(
529        region_id: &RegionId,
530        column_metadatas: Vec<ColumnMetadata>,
531        extensions: &mut HashMap<String, Vec<u8>>,
532    ) -> Result<()> {
533        extensions.insert(
534            TABLE_COLUMN_METADATA_EXTENSION_KEY.to_string(),
535            ColumnMetadata::encode_list(&column_metadatas).context(SerializeColumnMetadataSnafu)?,
536        );
537        info!(
538            "Added column metadatas: {:?} to extensions, region_id: {:?}",
539            column_metadatas, region_id
540        );
541        Ok(())
542    }
543
544    /// Find the current version's memtables and SSTs stats by region_id.
545    /// The stats must be collected in one place one time to ensure data consistency.
546    pub fn find_memtable_and_sst_stats(
547        &self,
548        region_id: RegionId,
549    ) -> Result<(Vec<MemtableStats>, Vec<FileMeta>)> {
550        let region = self
551            .find_region(region_id)
552            .context(RegionNotFoundSnafu { region_id })?;
553
554        let version = region.version();
555        let memtable_stats = version
556            .memtables
557            .list_memtables()
558            .iter()
559            .map(|x| x.stats())
560            .collect::<Vec<_>>();
561
562        let sst_stats = version
563            .ssts
564            .levels()
565            .iter()
566            .flat_map(|level| level.files().map(|x| x.meta_ref()))
567            .cloned()
568            .collect::<Vec<_>>();
569        Ok((memtable_stats, sst_stats))
570    }
571
572    /// Lists all SSTs from the manifest of all regions in the engine.
573    pub async fn all_ssts_from_manifest(&self) -> Vec<ManifestSstEntry> {
574        let node_id = self.inner.workers.file_ref_manager().node_id();
575        let regions = self.inner.workers.all_regions();
576
577        let mut results = Vec::new();
578        for region in regions {
579            let mut entries = region.manifest_sst_entries().await;
580            for e in &mut entries {
581                e.node_id = node_id;
582            }
583            results.extend(entries);
584        }
585
586        results
587    }
588
589    /// Lists metadata about all puffin index targets stored in the engine.
590    pub async fn all_index_metas(&self) -> Vec<PuffinIndexMetaEntry> {
591        let node_id = self.inner.workers.file_ref_manager().node_id();
592        let cache_manager = self.inner.workers.cache_manager();
593        let puffin_metadata_cache = cache_manager.puffin_metadata_cache().cloned();
594        let bloom_filter_cache = cache_manager.bloom_filter_index_cache().cloned();
595        let inverted_index_cache = cache_manager.inverted_index_cache().cloned();
596
597        let mut results = Vec::new();
598
599        for region in self.inner.workers.all_regions() {
600            let manifest_entries = region.manifest_sst_entries().await;
601            let access_layer = region.access_layer.clone();
602            let table_dir = access_layer.table_dir().to_string();
603            let path_type = access_layer.path_type();
604            let object_store = access_layer.object_store().clone();
605            let puffin_factory = access_layer.puffin_manager_factory().clone();
606            let path_factory = RegionFilePathFactory::new(table_dir, path_type);
607
608            let entry_futures = manifest_entries.into_iter().map(|entry| {
609                let object_store = object_store.clone();
610                let path_factory = path_factory.clone();
611                let puffin_factory = puffin_factory.clone();
612                let puffin_metadata_cache = puffin_metadata_cache.clone();
613                let bloom_filter_cache = bloom_filter_cache.clone();
614                let inverted_index_cache = inverted_index_cache.clone();
615
616                async move {
617                    let Some(index_file_path) = entry.index_file_path.as_ref() else {
618                        return Vec::new();
619                    };
620
621                    let index_version = entry.index_version;
622                    let file_id = match FileId::parse_str(&entry.file_id) {
623                        Ok(file_id) => file_id,
624                        Err(err) => {
625                            warn!(
626                                err;
627                                "Failed to parse puffin index file id, table_dir: {}, file_id: {}",
628                                entry.table_dir,
629                                entry.file_id
630                            );
631                            return Vec::new();
632                        }
633                    };
634                    // The index file path is derived from the physical file owner. After
635                    // repartition, `entry.region_id` is only the referring region.
636                    let region_index_id = RegionIndexId::new(
637                        RegionFileId::new(entry.origin_region_id, file_id),
638                        index_version,
639                    );
640                    let context = IndexEntryContext {
641                        table_dir: &entry.table_dir,
642                        index_file_path: index_file_path.as_str(),
643                        region_id: entry.region_id,
644                        table_id: entry.table_id,
645                        region_number: entry.region_number,
646                        region_group: entry.region_group,
647                        region_sequence: entry.region_sequence,
648                        file_id: &entry.file_id,
649                        index_file_size: entry.index_file_size,
650                        node_id,
651                    };
652
653                    let manager = puffin_factory
654                        .build(object_store, path_factory)
655                        .with_puffin_metadata_cache(puffin_metadata_cache);
656
657                    collect_index_entries_from_puffin(
658                        manager,
659                        region_index_id,
660                        context,
661                        bloom_filter_cache,
662                        inverted_index_cache,
663                    )
664                    .await
665                }
666            });
667
668            let mut meta_stream = stream::iter(entry_futures).buffer_unordered(8); // Parallelism is 8.
669            while let Some(mut metas) = meta_stream.next().await {
670                results.append(&mut metas);
671            }
672        }
673
674        results
675    }
676
677    /// Lists region info entries of all regions in the engine.
678    pub async fn all_region_infos(&self) -> Vec<RegionInfoEntry> {
679        let node_id = self.inner.workers.file_ref_manager().node_id();
680        self.inner
681            .workers
682            .all_regions()
683            .map(|region| region.region_info_entry(node_id))
684            .collect()
685    }
686
687    /// Lists all SSTs from the storage layer of all regions in the engine.
688    pub fn all_ssts_from_storage(&self) -> impl Stream<Item = Result<StorageSstEntry>> {
689        let node_id = self.inner.workers.file_ref_manager().node_id();
690        let regions = self.inner.workers.all_regions();
691
692        let mut layers_distinct_table_dirs = HashMap::new();
693        for region in regions {
694            let table_dir = region.access_layer.table_dir();
695            if !layers_distinct_table_dirs.contains_key(table_dir) {
696                layers_distinct_table_dirs
697                    .insert(table_dir.to_string(), region.access_layer.clone());
698            }
699        }
700
701        stream::iter(layers_distinct_table_dirs)
702            .map(|(_, access_layer)| access_layer.storage_sst_entries())
703            .flatten()
704            .map(move |entry| {
705                entry.map(move |mut entry| {
706                    entry.node_id = node_id;
707                    entry
708                })
709            })
710    }
711}
712
713/// Check whether the region edit is valid.
714///
715/// Only adding or removing files to region is considered valid now.
716fn is_valid_region_edit(edit: &RegionEdit) -> bool {
717    (!edit.files_to_add.is_empty() || !edit.files_to_remove.is_empty())
718        && matches!(
719            edit,
720            RegionEdit {
721                files_to_add: _,
722                files_to_remove: _,
723                timestamp_ms: _,
724                compaction_time_window: None,
725                flushed_entry_id: None,
726                flushed_sequence: None,
727                ..
728            }
729        )
730}
731
732/// Inner struct of [MitoEngine].
733struct EngineInner {
734    /// Region workers group.
735    workers: WorkerGroup,
736    /// Config of the engine.
737    config: Arc<MitoConfig>,
738    /// The Wal raw entry reader.
739    wal_raw_entry_reader: Arc<dyn RawEntryReader>,
740    /// Memory tracker for table scans.
741    scan_memory_tracker: QueryMemoryTracker,
742    /// The region hook (if any) registered via plugins; exposed for the GC worker
743    /// to fire [`RegionHook::on_region_gc`].
744    region_hook: Option<RegionHookRef>,
745    #[cfg(feature = "enterprise")]
746    extension_range_provider_factory: Option<BoxedExtensionRangeProviderFactory>,
747}
748
749type TopicGroupedRegionOpenRequests = HashMap<String, Vec<(RegionId, RegionOpenRequest)>>;
750
751/// Returns requests([TopicGroupedRegionOpenRequests]) grouped by topic and remaining requests.
752fn prepare_batch_open_requests(
753    requests: Vec<(RegionId, RegionOpenRequest)>,
754) -> Result<(
755    TopicGroupedRegionOpenRequests,
756    Vec<(RegionId, RegionOpenRequest)>,
757)> {
758    let mut topic_to_regions: HashMap<String, Vec<(RegionId, RegionOpenRequest)>> = HashMap::new();
759    let mut remaining_regions: Vec<(RegionId, RegionOpenRequest)> = Vec::new();
760    for (region_id, request) in requests {
761        match parse_wal_options(&request.options).context(SerdeJsonSnafu)? {
762            WalOptions::Kafka(options) => {
763                topic_to_regions
764                    .entry(options.topic)
765                    .or_default()
766                    .push((region_id, request));
767            }
768            WalOptions::RaftEngine | WalOptions::Noop => {
769                remaining_regions.push((region_id, request));
770            }
771        }
772    }
773
774    Ok((topic_to_regions, remaining_regions))
775}
776
777impl EngineInner {
778    #[cfg(feature = "enterprise")]
779    #[must_use]
780    fn with_extension_range_provider_factory(
781        self,
782        extension_range_provider_factory: Option<BoxedExtensionRangeProviderFactory>,
783    ) -> Self {
784        Self {
785            extension_range_provider_factory,
786            ..self
787        }
788    }
789
790    /// Stop the inner engine.
791    async fn stop(&self) -> Result<()> {
792        self.workers.stop().await
793    }
794
795    fn find_region(&self, region_id: RegionId) -> Result<MitoRegionRef> {
796        self.workers
797            .get_region(region_id)
798            .context(RegionNotFoundSnafu { region_id })
799    }
800
801    /// Get metadata of a region.
802    ///
803    /// Returns error if the region doesn't exist.
804    fn get_metadata(&self, region_id: RegionId) -> Result<RegionMetadataRef> {
805        // Reading a region doesn't need to go through the region worker thread.
806        let region = self.find_region(region_id)?;
807        Ok(region.metadata())
808    }
809
810    async fn open_topic_regions(
811        &self,
812        topic: String,
813        region_requests: Vec<(RegionId, RegionOpenRequest)>,
814    ) -> Result<Vec<(RegionId, Result<AffectedRows>)>> {
815        let now = Instant::now();
816        let region_ids = region_requests
817            .iter()
818            .map(|(region_id, _)| *region_id)
819            .collect::<Vec<_>>();
820        let provider = Provider::kafka_provider(topic);
821        let (distributor, entry_receivers) = build_wal_entry_distributor_and_receivers(
822            provider.clone(),
823            self.wal_raw_entry_reader.clone(),
824            &region_ids,
825            DEFAULT_ENTRY_RECEIVER_BUFFER_SIZE,
826        );
827
828        let mut responses = Vec::with_capacity(region_requests.len());
829        for ((region_id, request), entry_receiver) in
830            region_requests.into_iter().zip(entry_receivers)
831        {
832            let (request, receiver) =
833                WorkerRequest::new_open_region_request(region_id, request, Some(entry_receiver));
834            self.workers.submit_to_worker(region_id, request).await?;
835            responses.push(async move { receiver.await.context(RecvSnafu)? });
836        }
837
838        // Waits for entries distribution.
839        let distribution =
840            common_runtime::spawn_global(async move { distributor.distribute().await });
841        // Waits for worker returns.
842        let responses = join_all(responses).await;
843        distribution.await.context(JoinSnafu)??;
844
845        let num_failure = responses.iter().filter(|r| r.is_err()).count();
846        info!(
847            "Opened {} regions for topic '{}', failures: {}, elapsed: {:?}",
848            region_ids.len() - num_failure,
849            // Safety: provider is kafka provider.
850            provider.as_kafka_provider().unwrap(),
851            num_failure,
852            now.elapsed(),
853        );
854        Ok(region_ids.into_iter().zip(responses).collect())
855    }
856
857    async fn handle_batch_open_requests(
858        &self,
859        parallelism: usize,
860        requests: Vec<(RegionId, RegionOpenRequest)>,
861    ) -> Result<Vec<(RegionId, Result<AffectedRows>)>> {
862        let semaphore = Arc::new(Semaphore::new(parallelism));
863        let (topic_to_region_requests, remaining_region_requests) =
864            prepare_batch_open_requests(requests)?;
865        let mut responses =
866            Vec::with_capacity(topic_to_region_requests.len() + remaining_region_requests.len());
867
868        if !topic_to_region_requests.is_empty() {
869            let mut tasks = Vec::with_capacity(topic_to_region_requests.len());
870            for (topic, region_requests) in topic_to_region_requests {
871                let semaphore_moved = semaphore.clone();
872                tasks.push(async move {
873                    // Safety: semaphore must exist
874                    let _permit = semaphore_moved.acquire().await.unwrap();
875                    self.open_topic_regions(topic, region_requests).await
876                })
877            }
878            let r = try_join_all(tasks).await?;
879            responses.extend(r.into_iter().flatten());
880        }
881
882        if !remaining_region_requests.is_empty() {
883            let mut tasks = Vec::with_capacity(remaining_region_requests.len());
884            let mut region_ids = Vec::with_capacity(remaining_region_requests.len());
885            for (region_id, request) in remaining_region_requests {
886                let semaphore_moved = semaphore.clone();
887                region_ids.push(region_id);
888                tasks.push(async move {
889                    // Safety: semaphore must exist
890                    let _permit = semaphore_moved.acquire().await.unwrap();
891                    let (request, receiver) =
892                        WorkerRequest::new_open_region_request(region_id, request, None);
893
894                    self.workers.submit_to_worker(region_id, request).await?;
895
896                    receiver.await.context(RecvSnafu)?
897                })
898            }
899
900            let results = join_all(tasks).await;
901            responses.extend(region_ids.into_iter().zip(results));
902        }
903
904        Ok(responses)
905    }
906
907    async fn catchup_topic_regions(
908        &self,
909        provider: Provider,
910        region_requests: Vec<(RegionId, RegionCatchupRequest)>,
911    ) -> Result<Vec<(RegionId, Result<AffectedRows>)>> {
912        let now = Instant::now();
913        let region_ids = region_requests
914            .iter()
915            .map(|(region_id, _)| *region_id)
916            .collect::<Vec<_>>();
917        let (distributor, entry_receivers) = build_wal_entry_distributor_and_receivers(
918            provider.clone(),
919            self.wal_raw_entry_reader.clone(),
920            &region_ids,
921            DEFAULT_ENTRY_RECEIVER_BUFFER_SIZE,
922        );
923
924        let mut responses = Vec::with_capacity(region_requests.len());
925        for ((region_id, request), entry_receiver) in
926            region_requests.into_iter().zip(entry_receivers)
927        {
928            let (request, receiver) =
929                WorkerRequest::new_catchup_region_request(region_id, request, Some(entry_receiver));
930            self.workers.submit_to_worker(region_id, request).await?;
931            responses.push(async move { receiver.await.context(RecvSnafu)? });
932        }
933
934        // Wait for entries distribution.
935        let distribution =
936            common_runtime::spawn_global(async move { distributor.distribute().await });
937        // Wait for worker returns.
938        let responses = join_all(responses).await;
939        distribution.await.context(JoinSnafu)??;
940
941        let num_failure = responses.iter().filter(|r| r.is_err()).count();
942        info!(
943            "Caught up {} regions for topic '{}', failures: {}, elapsed: {:?}",
944            region_ids.len() - num_failure,
945            // Safety: provider is kafka provider.
946            provider.as_kafka_provider().unwrap(),
947            num_failure,
948            now.elapsed(),
949        );
950
951        Ok(region_ids.into_iter().zip(responses).collect())
952    }
953
954    async fn handle_batch_catchup_requests(
955        &self,
956        parallelism: usize,
957        requests: Vec<(RegionId, RegionCatchupRequest)>,
958    ) -> Result<Vec<(RegionId, Result<AffectedRows>)>> {
959        let mut responses = Vec::with_capacity(requests.len());
960        let mut topic_regions: HashMap<Arc<KafkaProvider>, Vec<_>> = HashMap::new();
961        let mut remaining_region_requests = vec![];
962
963        for (region_id, request) in requests {
964            match self.workers.get_region(region_id) {
965                Some(region) => match region.provider.as_kafka_provider() {
966                    Some(provider) => {
967                        topic_regions
968                            .entry(provider.clone())
969                            .or_default()
970                            .push((region_id, request));
971                    }
972                    None => {
973                        remaining_region_requests.push((region_id, request));
974                    }
975                },
976                None => responses.push((region_id, RegionNotFoundSnafu { region_id }.fail())),
977            }
978        }
979
980        let semaphore = Arc::new(Semaphore::new(parallelism));
981
982        if !topic_regions.is_empty() {
983            let mut tasks = Vec::with_capacity(topic_regions.len());
984            for (provider, region_requests) in topic_regions {
985                let semaphore_moved = semaphore.clone();
986                tasks.push(async move {
987                    // Safety: semaphore must exist
988                    let _permit = semaphore_moved.acquire().await.unwrap();
989                    self.catchup_topic_regions(Provider::Kafka(provider), region_requests)
990                        .await
991                })
992            }
993
994            let r = try_join_all(tasks).await?;
995            responses.extend(r.into_iter().flatten());
996        }
997
998        if !remaining_region_requests.is_empty() {
999            let mut tasks = Vec::with_capacity(remaining_region_requests.len());
1000            let mut region_ids = Vec::with_capacity(remaining_region_requests.len());
1001            for (region_id, request) in remaining_region_requests {
1002                let semaphore_moved = semaphore.clone();
1003                region_ids.push(region_id);
1004                tasks.push(async move {
1005                    // Safety: semaphore must exist
1006                    let _permit = semaphore_moved.acquire().await.unwrap();
1007                    let (request, receiver) =
1008                        WorkerRequest::new_catchup_region_request(region_id, request, None);
1009
1010                    self.workers.submit_to_worker(region_id, request).await?;
1011
1012                    receiver.await.context(RecvSnafu)?
1013                })
1014            }
1015
1016            let results = join_all(tasks).await;
1017            responses.extend(region_ids.into_iter().zip(results));
1018        }
1019
1020        Ok(responses)
1021    }
1022
1023    /// Handles [RegionRequest] and return its executed result.
1024    async fn handle_request(
1025        &self,
1026        region_id: RegionId,
1027        request: RegionRequest,
1028    ) -> Result<AffectedRows> {
1029        let region_metadata = self.get_metadata(region_id).ok();
1030        let (request, receiver) =
1031            WorkerRequest::try_from_region_request(region_id, request, region_metadata)?;
1032        self.workers.submit_to_worker(region_id, request).await?;
1033
1034        receiver.await.context(RecvSnafu)?
1035    }
1036
1037    /// Returns the sequence of latest committed data.
1038    fn get_committed_sequence(&self, region_id: RegionId) -> Result<SequenceNumber> {
1039        // Reading a region doesn't need to go through the region worker thread.
1040        self.find_region(region_id)
1041            .map(|r| r.find_committed_sequence())
1042    }
1043
1044    /// Handles the scan `request` and returns a [ScanRegion].
1045    #[tracing::instrument(skip_all, fields(region_id = %region_id))]
1046    fn scan_region(&self, region_id: RegionId, mut request: ScanRequest) -> Result<ScanRegion> {
1047        let query_start = Instant::now();
1048        // Reading a region doesn't need to go through the region worker thread.
1049        let region = self.find_region(region_id)?;
1050        let version_data = region.version_control.current();
1051        let version = version_data.version;
1052
1053        if request.snapshot_on_scan && request.memtable_max_sequence.is_none() {
1054            request.memtable_max_sequence = Some(version_data.committed_sequence);
1055        }
1056
1057        if let Some(given_seq) = request.memtable_min_sequence {
1058            let min_readable_seq = version.flushed_sequence;
1059            ensure!(
1060                given_seq >= min_readable_seq,
1061                IncrementalQueryStaleSnafu {
1062                    region_id,
1063                    given_seq,
1064                    min_readable_seq,
1065                }
1066            );
1067        }
1068
1069        if let Some(given_seq) = request.memtable_max_sequence
1070            && !request.skip_sst_files
1071        {
1072            // Explicit snapshot fences that include SST reads are enforceable
1073            // only while the requested upper bound is not older than the
1074            // region's flushed frontier. If H has already been flushed into SST,
1075            // mito cannot apply a memtable-only sequence upper bound to that
1076            // SST scan, so fail and let Flow rebind the fenced repair instead
1077            // of reading rows beyond H.
1078            let min_enforceable_seq = version.flushed_sequence;
1079            ensure!(
1080                given_seq >= min_enforceable_seq,
1081                SnapshotFenceStaleSnafu {
1082                    region_id,
1083                    given_seq,
1084                    min_enforceable_seq,
1085                }
1086            );
1087        }
1088
1089        // Get cache.
1090        let cache_manager = self.workers.cache_manager();
1091
1092        let scan_region = ScanRegion::new(
1093            version,
1094            region.access_layer.clone(),
1095            request,
1096            CacheStrategy::EnableAll(cache_manager),
1097        )
1098        .with_query_stat_counters(region.region_stats.query_stat_counters())
1099        .with_max_concurrent_scan_files(self.config.max_concurrent_scan_files)
1100        .with_ignore_inverted_index(self.config.inverted_index.apply_on_query.disabled())
1101        .with_ignore_fulltext_index(self.config.fulltext_index.apply_on_query.disabled())
1102        .with_ignore_bloom_filter(self.config.bloom_filter_index.apply_on_query.disabled())
1103        .with_start_time(query_start);
1104
1105        #[cfg(feature = "enterprise")]
1106        let scan_region = self.maybe_fill_extension_range_provider(scan_region, region);
1107
1108        Ok(scan_region)
1109    }
1110
1111    #[cfg(feature = "enterprise")]
1112    fn maybe_fill_extension_range_provider(
1113        &self,
1114        mut scan_region: ScanRegion,
1115        region: MitoRegionRef,
1116    ) -> ScanRegion {
1117        if region.is_follower()
1118            && let Some(factory) = self.extension_range_provider_factory.as_ref()
1119        {
1120            scan_region
1121                .set_extension_range_provider(factory.create_extension_range_provider(region));
1122        }
1123        scan_region
1124    }
1125
1126    /// Converts the [`RegionRole`].
1127    fn set_region_role(&self, region_id: RegionId, role: RegionRole) -> Result<()> {
1128        let region = self.find_region(region_id)?;
1129        region.set_role(role);
1130        Ok(())
1131    }
1132
1133    /// Sets read-only for a region and ensures no more writes in the region after it returns.
1134    async fn set_region_role_state_gracefully(
1135        &self,
1136        region_id: RegionId,
1137        region_role_state: SettableRegionRoleState,
1138    ) -> Result<SetRegionRoleStateResponse> {
1139        // Notes: It acquires the mutable ownership to ensure no other threads,
1140        // Therefore, we submit it to the worker.
1141        let (request, receiver) =
1142            WorkerRequest::new_set_readonly_gracefully(region_id, region_role_state);
1143        self.workers.submit_to_worker(region_id, request).await?;
1144
1145        receiver.await.context(RecvSnafu)
1146    }
1147
1148    async fn sync_region(
1149        &self,
1150        region_id: RegionId,
1151        manifest_info: RegionManifestInfo,
1152    ) -> Result<(ManifestVersion, bool)> {
1153        ensure!(manifest_info.is_mito(), MitoManifestInfoSnafu);
1154        let manifest_version = manifest_info.data_manifest_version();
1155        let (request, receiver) =
1156            WorkerRequest::new_sync_region_request(region_id, manifest_version);
1157        self.workers.submit_to_worker(region_id, request).await?;
1158
1159        receiver.await.context(RecvSnafu)?
1160    }
1161
1162    async fn remap_manifests(
1163        &self,
1164        request: RemapManifestsRequest,
1165    ) -> Result<RemapManifestsResponse> {
1166        let region_id = request.region_id;
1167        let (request, receiver) = WorkerRequest::try_from_remap_manifests_request(request)?;
1168        self.workers.submit_to_worker(region_id, request).await?;
1169        let manifest_paths = receiver.await.context(RecvSnafu)??;
1170        Ok(RemapManifestsResponse { manifest_paths })
1171    }
1172
1173    async fn copy_region_from(
1174        &self,
1175        region_id: RegionId,
1176        request: MitoCopyRegionFromRequest,
1177    ) -> Result<MitoCopyRegionFromResponse> {
1178        let (request, receiver) =
1179            WorkerRequest::try_from_copy_region_from_request(region_id, request)?;
1180        self.workers.submit_to_worker(region_id, request).await?;
1181        let response = receiver.await.context(RecvSnafu)??;
1182        Ok(response)
1183    }
1184
1185    fn role(&self, region_id: RegionId) -> Option<RegionRole> {
1186        self.workers
1187            .get_region(region_id)
1188            .map(|region| region.region_role())
1189    }
1190}
1191
1192fn map_batch_responses(responses: Vec<(RegionId, Result<AffectedRows>)>) -> BatchResponses {
1193    responses
1194        .into_iter()
1195        .map(|(region_id, response)| {
1196            (
1197                region_id,
1198                response.map(RegionResponse::new).map_err(BoxedError::new),
1199            )
1200        })
1201        .collect()
1202}
1203
1204#[async_trait]
1205impl RegionEngine for MitoEngine {
1206    fn name(&self) -> &str {
1207        MITO_ENGINE_NAME
1208    }
1209
1210    #[tracing::instrument(skip_all)]
1211    async fn handle_batch_open_requests(
1212        &self,
1213        parallelism: usize,
1214        requests: Vec<(RegionId, RegionOpenRequest)>,
1215    ) -> Result<BatchResponses, BoxedError> {
1216        // TODO(weny): add metrics.
1217        self.inner
1218            .handle_batch_open_requests(parallelism, requests)
1219            .await
1220            .map(map_batch_responses)
1221            .map_err(BoxedError::new)
1222    }
1223
1224    #[tracing::instrument(skip_all)]
1225    async fn handle_batch_catchup_requests(
1226        &self,
1227        parallelism: usize,
1228        requests: Vec<(RegionId, RegionCatchupRequest)>,
1229    ) -> Result<BatchResponses, BoxedError> {
1230        self.inner
1231            .handle_batch_catchup_requests(parallelism, requests)
1232            .await
1233            .map(map_batch_responses)
1234            .map_err(BoxedError::new)
1235    }
1236
1237    #[tracing::instrument(skip_all)]
1238    async fn handle_request(
1239        &self,
1240        region_id: RegionId,
1241        request: RegionRequest,
1242    ) -> Result<RegionResponse, BoxedError> {
1243        let _timer = HANDLE_REQUEST_ELAPSED
1244            .with_label_values(&[request.request_type()])
1245            .start_timer();
1246
1247        let is_alter = matches!(request, RegionRequest::Alter(_));
1248        let is_create = matches!(request, RegionRequest::Create(_));
1249        let mut response = self
1250            .inner
1251            .handle_request(region_id, request)
1252            .await
1253            .map(RegionResponse::new)
1254            .map_err(BoxedError::new)?;
1255
1256        if is_alter {
1257            self.handle_alter_response(region_id, &mut response)
1258                .map_err(BoxedError::new)?;
1259        } else if is_create {
1260            self.handle_create_response(region_id, &mut response)
1261                .map_err(BoxedError::new)?;
1262        }
1263
1264        Ok(response)
1265    }
1266
1267    #[tracing::instrument(skip_all)]
1268    async fn handle_query(
1269        &self,
1270        region_id: RegionId,
1271        request: ScanRequest,
1272    ) -> Result<RegionScannerRef, BoxedError> {
1273        self.scan_region(region_id, request)
1274            .map_err(BoxedError::new)?
1275            .region_scanner()
1276            .await
1277            .map_err(BoxedError::new)
1278    }
1279
1280    fn query_memory_tracker(&self) -> Option<QueryMemoryTracker> {
1281        Some(self.inner.scan_memory_tracker.clone())
1282    }
1283
1284    async fn get_committed_sequence(
1285        &self,
1286        region_id: RegionId,
1287    ) -> Result<SequenceNumber, BoxedError> {
1288        self.inner
1289            .get_committed_sequence(region_id)
1290            .map_err(BoxedError::new)
1291    }
1292
1293    /// Retrieve region's metadata.
1294    async fn get_metadata(
1295        &self,
1296        region_id: RegionId,
1297    ) -> std::result::Result<RegionMetadataRef, BoxedError> {
1298        self.inner.get_metadata(region_id).map_err(BoxedError::new)
1299    }
1300
1301    /// Stop the engine.
1302    ///
1303    /// Stopping the engine doesn't stop the underlying log store as other components might
1304    /// still use it. (When no other components are referencing the log store, it will
1305    /// automatically shutdown.)
1306    async fn stop(&self) -> std::result::Result<(), BoxedError> {
1307        self.inner.stop().await.map_err(BoxedError::new)
1308    }
1309
1310    fn region_statistic(&self, region_id: RegionId) -> Option<RegionStatistic> {
1311        self.get_region_statistic(region_id)
1312    }
1313
1314    fn set_region_role(&self, region_id: RegionId, role: RegionRole) -> Result<(), BoxedError> {
1315        self.inner
1316            .set_region_role(region_id, role)
1317            .map_err(BoxedError::new)
1318    }
1319
1320    async fn set_region_role_state_gracefully(
1321        &self,
1322        region_id: RegionId,
1323        region_role_state: SettableRegionRoleState,
1324    ) -> Result<SetRegionRoleStateResponse, BoxedError> {
1325        let _timer = HANDLE_REQUEST_ELAPSED
1326            .with_label_values(&["set_region_role_state_gracefully"])
1327            .start_timer();
1328
1329        self.inner
1330            .set_region_role_state_gracefully(region_id, region_role_state)
1331            .await
1332            .map_err(BoxedError::new)
1333    }
1334
1335    async fn sync_region(
1336        &self,
1337        region_id: RegionId,
1338        request: SyncRegionFromRequest,
1339    ) -> Result<SyncRegionFromResponse, BoxedError> {
1340        let manifest_info = request
1341            .into_region_manifest_info()
1342            .context(UnexpectedSnafu {
1343                err_msg: "Expected a manifest info request",
1344            })
1345            .map_err(BoxedError::new)?;
1346        let (_, synced) = self
1347            .inner
1348            .sync_region(region_id, manifest_info)
1349            .await
1350            .map_err(BoxedError::new)?;
1351
1352        Ok(SyncRegionFromResponse::Mito { synced })
1353    }
1354
1355    async fn remap_manifests(
1356        &self,
1357        request: RemapManifestsRequest,
1358    ) -> Result<RemapManifestsResponse, BoxedError> {
1359        self.inner
1360            .remap_manifests(request)
1361            .await
1362            .map_err(BoxedError::new)
1363    }
1364
1365    fn role(&self, region_id: RegionId) -> Option<RegionRole> {
1366        self.inner.role(region_id)
1367    }
1368
1369    fn as_any(&self) -> &dyn Any {
1370        self
1371    }
1372}
1373
1374impl MitoEngine {
1375    fn handle_alter_response(
1376        &self,
1377        region_id: RegionId,
1378        response: &mut RegionResponse,
1379    ) -> Result<()> {
1380        if let Some(statistic) = self.region_statistic(region_id) {
1381            Self::encode_manifest_info_to_extensions(
1382                &region_id,
1383                statistic.manifest,
1384                &mut response.extensions,
1385            )?;
1386        }
1387        let column_metadatas = self
1388            .inner
1389            .find_region(region_id)
1390            .ok()
1391            .map(|r| r.metadata().column_metadatas.clone());
1392        if let Some(column_metadatas) = column_metadatas {
1393            Self::encode_column_metadatas_to_extensions(
1394                &region_id,
1395                column_metadatas,
1396                &mut response.extensions,
1397            )?;
1398        }
1399        Ok(())
1400    }
1401
1402    fn handle_create_response(
1403        &self,
1404        region_id: RegionId,
1405        response: &mut RegionResponse,
1406    ) -> Result<()> {
1407        let column_metadatas = self
1408            .inner
1409            .find_region(region_id)
1410            .ok()
1411            .map(|r| r.metadata().column_metadatas.clone());
1412        if let Some(column_metadatas) = column_metadatas {
1413            Self::encode_column_metadatas_to_extensions(
1414                &region_id,
1415                column_metadatas,
1416                &mut response.extensions,
1417            )?;
1418        }
1419        Ok(())
1420    }
1421}
1422
1423// Tests methods.
1424#[cfg(any(test, feature = "test"))]
1425#[allow(clippy::too_many_arguments)]
1426impl MitoEngine {
1427    /// Returns a new [MitoEngine] for tests.
1428    pub async fn new_for_test<S: LogStore>(
1429        data_home: &str,
1430        mut config: MitoConfig,
1431        log_store: Arc<S>,
1432        object_store_manager: ObjectStoreManagerRef,
1433        write_buffer_manager: Option<crate::flush::WriteBufferManagerRef>,
1434        listener: Option<crate::engine::listener::EventListenerRef>,
1435        time_provider: crate::time_provider::TimeProviderRef,
1436        schema_metadata_manager: SchemaMetadataManagerRef,
1437        file_ref_manager: FileReferenceManagerRef,
1438        partition_expr_fetcher: PartitionExprFetcherRef,
1439    ) -> Result<MitoEngine> {
1440        config.sanitize(data_home)?;
1441
1442        let config = Arc::new(config);
1443        let wal_raw_entry_reader = Arc::new(LogStoreRawEntryReader::new(log_store.clone()));
1444        let total_memory = get_total_memory_bytes().max(0) as u64;
1445        let scan_memory_limit = config.scan_memory_limit.resolve(total_memory) as usize;
1446        let scan_memory_tracker =
1447            QueryMemoryTracker::builder(scan_memory_limit, config.scan_memory_on_exhausted)
1448                .on_update(|usage| {
1449                    SCAN_MEMORY_USAGE_BYTES.set(usage as i64);
1450                })
1451                .on_exhausted(|| {
1452                    SCAN_MEMORY_EXHAUSTED_TOTAL.inc();
1453                })
1454                .on_reject(|| {
1455                    SCAN_REQUESTS_REJECTED_TOTAL.inc();
1456                })
1457                .build();
1458        Ok(MitoEngine {
1459            inner: Arc::new(EngineInner {
1460                workers: WorkerGroup::start_for_test(
1461                    config.clone(),
1462                    log_store,
1463                    object_store_manager,
1464                    write_buffer_manager,
1465                    listener,
1466                    schema_metadata_manager,
1467                    file_ref_manager,
1468                    time_provider,
1469                    partition_expr_fetcher,
1470                )
1471                .await?,
1472                config,
1473                wal_raw_entry_reader,
1474                scan_memory_tracker,
1475                region_hook: None,
1476                #[cfg(feature = "enterprise")]
1477                extension_range_provider_factory: None,
1478            }),
1479        })
1480    }
1481
1482    /// Returns the purge scheduler.
1483    pub fn purge_scheduler(&self) -> &crate::schedule::scheduler::SchedulerRef {
1484        self.inner.workers.purge_scheduler()
1485    }
1486}
1487
1488#[cfg(test)]
1489mod tests {
1490    use std::time::Duration;
1491
1492    use super::*;
1493    use crate::sst::file::FileMeta;
1494
1495    #[test]
1496    fn test_is_valid_region_edit() {
1497        // Valid: has only "files_to_add"
1498        let edit = RegionEdit {
1499            files_to_add: vec![FileMeta::default()],
1500            files_to_remove: vec![],
1501            timestamp_ms: None,
1502            compaction_time_window: None,
1503            flushed_entry_id: None,
1504            flushed_sequence: None,
1505            committed_sequence: None,
1506        };
1507        assert!(is_valid_region_edit(&edit));
1508
1509        // Invalid: "files_to_add" and "files_to_remove" are both empty
1510        let edit = RegionEdit {
1511            files_to_add: vec![],
1512            files_to_remove: vec![],
1513            timestamp_ms: None,
1514            compaction_time_window: None,
1515            flushed_entry_id: None,
1516            flushed_sequence: None,
1517            committed_sequence: None,
1518        };
1519        assert!(!is_valid_region_edit(&edit));
1520
1521        // Valid: has only "files_to_remove"
1522        let edit = RegionEdit {
1523            files_to_add: vec![],
1524            files_to_remove: vec![FileMeta::default()],
1525            timestamp_ms: None,
1526            compaction_time_window: None,
1527            flushed_entry_id: None,
1528            flushed_sequence: None,
1529            committed_sequence: None,
1530        };
1531        assert!(is_valid_region_edit(&edit));
1532
1533        // Valid: both "files_to_add" and "files_to_remove" are not empty
1534        let edit = RegionEdit {
1535            files_to_add: vec![FileMeta::default()],
1536            files_to_remove: vec![FileMeta::default()],
1537            timestamp_ms: None,
1538            compaction_time_window: None,
1539            flushed_entry_id: None,
1540            flushed_sequence: None,
1541            committed_sequence: None,
1542        };
1543        assert!(is_valid_region_edit(&edit));
1544
1545        // Invalid: other fields are not all "None"s
1546        let edit = RegionEdit {
1547            files_to_add: vec![FileMeta::default()],
1548            files_to_remove: vec![],
1549            timestamp_ms: None,
1550            compaction_time_window: Some(Duration::from_secs(1)),
1551            flushed_entry_id: None,
1552            flushed_sequence: None,
1553            committed_sequence: None,
1554        };
1555        assert!(!is_valid_region_edit(&edit));
1556        let edit = RegionEdit {
1557            files_to_add: vec![FileMeta::default()],
1558            files_to_remove: vec![],
1559            timestamp_ms: None,
1560            compaction_time_window: None,
1561            flushed_entry_id: Some(1),
1562            flushed_sequence: None,
1563            committed_sequence: None,
1564        };
1565        assert!(!is_valid_region_edit(&edit));
1566        let edit = RegionEdit {
1567            files_to_add: vec![FileMeta::default()],
1568            files_to_remove: vec![],
1569            timestamp_ms: None,
1570            compaction_time_window: None,
1571            flushed_entry_id: None,
1572            flushed_sequence: Some(1),
1573            committed_sequence: None,
1574        };
1575        assert!(!is_valid_region_edit(&edit));
1576    }
1577}