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