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