Skip to main content

mito2/region/
opener.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//! Region opener.
16
17use std::any::TypeId;
18use std::collections::HashMap;
19use std::sync::atomic::{AtomicI64, AtomicU64};
20use std::sync::{Arc, LazyLock};
21use std::time::Instant;
22
23use arrow_schema::extension::ExtensionType;
24use common_telemetry::{debug, error, info, warn};
25use common_wal::options::WalOptions;
26use datatypes::extension::json::{Json2ExtensionType, JsonMetadata};
27use futures::StreamExt;
28use futures::future::BoxFuture;
29use log_store::kafka::log_store::KafkaLogStore;
30use log_store::noop::log_store::NoopLogStore;
31use log_store::raft_engine::log_store::RaftEngineLogStore;
32use object_store::ObjectStore;
33use object_store::manager::ObjectStoreManagerRef;
34use object_store::util::{is_object_storage, normalize_dir};
35use parquet::file::metadata::PageIndexPolicy;
36use snafu::{OptionExt, ResultExt, ensure};
37use store_api::codec::PrimaryKeyEncoding;
38use store_api::logstore::LogStore;
39use store_api::logstore::provider::Provider;
40use store_api::metadata::{
41    ColumnMetadata, RegionMetadata, RegionMetadataBuilder, RegionMetadataRef,
42};
43use store_api::region_engine::RegionRole;
44use store_api::region_request::{PathType, RegionRequirements};
45use store_api::storage::{ColumnId, RegionId};
46use tokio::sync::Semaphore;
47
48use crate::access_layer::AccessLayer;
49use crate::cache::file_cache::{FileCache, FileType, IndexKey};
50use crate::cache::{CacheManagerRef, SstMetaPreparation, prepare_sst_meta};
51use crate::config::MitoConfig;
52use crate::engine::region_hook::RegionHookRef;
53use crate::error;
54use crate::error::{
55    DataTypeMismatchSnafu, EmptyRegionDirSnafu, InvalidMetadataSnafu, InvalidRegionOptionsSnafu,
56    ObjectStoreNotFoundSnafu, RegionCorruptedSnafu, Result, StaleLogEntrySnafu,
57};
58use crate::manifest::action::RegionManifest;
59use crate::manifest::manager::{RegionManifestManager, RegionManifestOptions};
60use crate::memtable::bulk::part::BulkPart;
61use crate::memtable::time_partition::{TimePartitions, TimePartitionsRef};
62use crate::memtable::{MemtableBuilderProvider, ensure_json2_not_use_time_series_memtable};
63use crate::metrics::{CACHE_FILL_DOWNLOADED_FILES, CACHE_FILL_PENDING_FILES};
64use crate::read::series_candidate::is_sparse_metric_metadata;
65use crate::region::options::RegionOptions;
66use crate::region::version::{VersionBuilder, VersionControl, VersionControlRef};
67use crate::region::{
68    ManifestContext, ManifestStats, MitoRegion, MitoRegionRef, RegionLeaderState, RegionRoleState,
69    RegionStats,
70};
71use crate::region_write_ctx::RegionWriteCtx;
72use crate::request::OptionOutputTx;
73use crate::schedule::scheduler::SchedulerRef;
74use crate::series_index::{IndexFilePurger, load_version_control};
75use crate::sst::FormatType;
76use crate::sst::file::{FileHandle, RegionFileId, RegionIndexId};
77use crate::sst::file_purger::{FilePurgerRef, create_file_purger};
78use crate::sst::file_ref::FileReferenceManagerRef;
79use crate::sst::index::intermediate::IntermediateManager;
80use crate::sst::index::puffin_manager::PuffinManagerFactory;
81use crate::sst::location::{self, region_dir_from_table_dir};
82use crate::sst::parquet::metadata::{MetadataLoader, extract_primary_key_range};
83use crate::sst::parquet::reader::MetadataCacheMetrics;
84use crate::sst::range_index::RangeIndexDeleter;
85use crate::time_provider::TimeProviderRef;
86use crate::wal::entry_reader::WalEntryReader;
87use crate::wal::{EntryId, Wal};
88
89const PARQUET_META_PRELOAD_CONCURRENCY: usize = 8;
90
91static PARQUET_META_PRELOAD_SEMAPHORE: LazyLock<Semaphore> =
92    LazyLock::new(|| Semaphore::new(PARQUET_META_PRELOAD_CONCURRENCY));
93
94fn initial_pruned_entry_id(wal_options: &WalOptions) -> EntryId {
95    match wal_options {
96        WalOptions::Kafka(options) => options.initial_pruned_entry_id.unwrap_or(0),
97        WalOptions::RaftEngine | WalOptions::Noop | WalOptions::ObjectStore(_) => 0,
98    }
99}
100
101fn maybe_upgrade_json2_layout(metadata: RegionMetadataRef) -> Result<RegionMetadataRef> {
102    let mut upgrades = Vec::new();
103    for (index, column) in metadata.column_metadatas.iter().enumerate() {
104        if !column.column_schema.data_type.is_json2() {
105            continue;
106        }
107        let Some(extension) = column
108            .column_schema
109            .extension_type::<Json2ExtensionType>()
110            .context(DataTypeMismatchSnafu)?
111        else {
112            continue;
113        };
114        if extension.metadata().is_version_2() {
115            continue;
116        }
117        upgrades.push((index, extension.metadata().json_settings().clone()));
118    }
119
120    if upgrades.is_empty() {
121        return Ok(metadata);
122    }
123
124    let mut upgraded = metadata.as_ref().clone();
125    for (index, settings) in upgrades {
126        let extension = Json2ExtensionType::new(Arc::new(JsonMetadata::new(settings)));
127        upgraded.column_metadatas[index]
128            .column_schema
129            .with_extension_type(&extension);
130    }
131
132    let builder = RegionMetadataBuilder::from_existing(upgraded);
133    Ok(Arc::new(builder.build().context(InvalidMetadataSnafu)?))
134}
135
136/// A fetcher to retrieve partition expr for a region.
137///
138/// Compatibility: older regions didn't persist `partition_expr` in engine metadata,
139/// while newer ones do. On open, we backfill it via this fetcher and persist it
140/// to the manifest so future opens don't need refetching.
141#[async_trait::async_trait]
142pub trait PartitionExprFetcher {
143    async fn fetch_expr(&self, region_id: RegionId) -> Option<String>;
144}
145
146pub type PartitionExprFetcherRef = Arc<dyn PartitionExprFetcher + Send + Sync>;
147
148/// Builder to create a new [MitoRegion] or open an existing one.
149pub(crate) struct RegionOpener {
150    region_id: RegionId,
151    metadata_builder: Option<RegionMetadataBuilder>,
152    memtable_builder_provider: MemtableBuilderProvider,
153    object_store_manager: ObjectStoreManagerRef,
154    table_dir: String,
155    path_type: PathType,
156    purge_scheduler: SchedulerRef,
157    options: Option<RegionOptions>,
158    cache_manager: Option<CacheManagerRef>,
159    skip_wal_replay: bool,
160    puffin_manager_factory: PuffinManagerFactory,
161    intermediate_manager: IntermediateManager,
162    time_provider: TimeProviderRef,
163    stats: ManifestStats,
164    wal_entry_reader: Option<Box<dyn WalEntryReader>>,
165    replay_checkpoint: Option<u64>,
166    file_ref_manager: FileReferenceManagerRef,
167    partition_expr_fetcher: PartitionExprFetcherRef,
168    hook: Option<RegionHookRef>,
169    series_index_store: Option<ObjectStore>,
170    series_index_purger: Option<IndexFilePurger>,
171}
172
173impl RegionOpener {
174    /// Returns a new opener.
175    // TODO(LFC): Reduce the number of arguments.
176    #[allow(clippy::too_many_arguments)]
177    pub(crate) fn new(
178        region_id: RegionId,
179        table_dir: &str,
180        path_type: PathType,
181        memtable_builder_provider: MemtableBuilderProvider,
182        object_store_manager: ObjectStoreManagerRef,
183        purge_scheduler: SchedulerRef,
184        puffin_manager_factory: PuffinManagerFactory,
185        intermediate_manager: IntermediateManager,
186        time_provider: TimeProviderRef,
187        file_ref_manager: FileReferenceManagerRef,
188        partition_expr_fetcher: PartitionExprFetcherRef,
189    ) -> RegionOpener {
190        RegionOpener {
191            region_id,
192            metadata_builder: None,
193            memtable_builder_provider,
194            object_store_manager,
195            table_dir: normalize_dir(table_dir),
196            path_type,
197            purge_scheduler,
198            options: None,
199            cache_manager: None,
200            skip_wal_replay: false,
201            puffin_manager_factory,
202            intermediate_manager,
203            time_provider,
204            stats: Default::default(),
205            wal_entry_reader: None,
206            replay_checkpoint: None,
207            file_ref_manager,
208            partition_expr_fetcher,
209            hook: None,
210            series_index_store: None,
211            series_index_purger: None,
212        }
213    }
214
215    /// Sets the store for companion range indexes.
216    pub(crate) fn series_index_store(mut self, store: Option<ObjectStore>) -> Self {
217        self.series_index_store = store;
218        self
219    }
220
221    /// Sets the purger shared with the worker's series-index maintenance task.
222    pub(crate) fn series_index_purger(mut self, purger: Option<IndexFilePurger>) -> Self {
223        self.series_index_purger = purger;
224        self
225    }
226
227    /// Sets the region hook for observing manifest mutations.
228    pub(crate) fn hook(mut self, hook: Option<RegionHookRef>) -> Self {
229        self.hook = hook;
230        self
231    }
232
233    /// Sets metadata builder of the region to create.
234    pub(crate) fn metadata_builder(mut self, builder: RegionMetadataBuilder) -> Self {
235        self.metadata_builder = Some(builder);
236        self
237    }
238
239    /// Computes the region directory from table_dir and region_id.
240    fn region_dir(&self) -> String {
241        region_dir_from_table_dir(&self.table_dir, self.region_id, self.path_type)
242    }
243
244    /// Builds the region metadata.
245    ///
246    /// # Panics
247    /// - Panics if `options` is not set.
248    /// - Panics if `metadata_builder` is not set.
249    fn build_metadata(&mut self) -> Result<RegionMetadata> {
250        let options = self.options.as_ref().unwrap();
251        let mut metadata_builder = self.metadata_builder.take().unwrap();
252        metadata_builder.primary_key_encoding(options.primary_key_encoding());
253        metadata_builder.build().context(InvalidMetadataSnafu)
254    }
255
256    /// Parses and sets options for the region.
257    pub(crate) fn parse_options(self, options: HashMap<String, String>) -> Result<Self> {
258        let region_id = self.region_id;
259        let options = self
260            .memtable_builder_provider
261            .parse_options(region_id, &options)?;
262        self.options(options)
263    }
264
265    /// Sets the replay checkpoint for the region.
266    pub(crate) fn replay_checkpoint(mut self, replay_checkpoint: Option<u64>) -> Self {
267        self.replay_checkpoint = replay_checkpoint;
268        self
269    }
270
271    /// If a [WalEntryReader] is set, the [RegionOpener] will use [WalEntryReader] instead of
272    /// constructing a new one from scratch.
273    pub(crate) fn wal_entry_reader(
274        mut self,
275        wal_entry_reader: Option<Box<dyn WalEntryReader>>,
276    ) -> Self {
277        self.wal_entry_reader = wal_entry_reader;
278        self
279    }
280
281    /// Sets options for the region.
282    pub(crate) fn options(mut self, options: RegionOptions) -> Result<Self> {
283        options.validate()?;
284        self.options = Some(options);
285        Ok(self)
286    }
287
288    /// Ensures the current region request satisfies its requirements.
289    pub(crate) fn ensure_region_requirements(
290        &self,
291        requirements: RegionRequirements,
292    ) -> Result<()> {
293        if !requirements.object_storage {
294            return Ok(());
295        }
296
297        let options = self.options.as_ref().context(InvalidRegionOptionsSnafu {
298            reason: "missing region options before requirement check".to_string(),
299        })?;
300        let object_store = get_object_store(&options.storage, &self.object_store_manager)?;
301
302        ensure!(
303            supports_open_region_object_storage_requirement(&object_store),
304            error::RegionRequirementSnafu {
305                region_id: self.region_id,
306                requirement: "object storage",
307                reason: "region data must be accessible from another datanode",
308            }
309        );
310
311        Ok(())
312    }
313
314    /// Sets the cache manager for the region.
315    pub(crate) fn cache(mut self, cache_manager: Option<CacheManagerRef>) -> Self {
316        self.cache_manager = cache_manager;
317        self
318    }
319
320    /// Sets the `skip_wal_replay`.
321    pub(crate) fn skip_wal_replay(mut self, skip: bool) -> Self {
322        self.skip_wal_replay = skip;
323        self
324    }
325
326    /// Writes region manifest and creates a new region if it does not exist.
327    /// Opens the region if it already exists.
328    ///
329    /// # Panics
330    /// - Panics if `metadata_builder` is not set.
331    /// - Panics if `options` is not set.
332    pub(crate) async fn create_or_open<S: LogStore>(
333        mut self,
334        config: &MitoConfig,
335        wal: &Wal<S>,
336    ) -> Result<MitoRegionRef> {
337        let region_id = self.region_id;
338        let region_dir = self.region_dir();
339        let metadata = self.build_metadata()?;
340        // Tries to open the region.
341        match self.maybe_open(config, wal).await {
342            Ok(Some(region)) => {
343                let recovered = region.metadata();
344                // Checks the schema of the region.
345                let expect = &metadata;
346                check_recovered_region(
347                    &recovered,
348                    expect.region_id,
349                    &expect.column_metadatas,
350                    &expect.primary_key,
351                )?;
352                // To keep consistency with Create behavior, set the opened Region to RegionRole::Leader.
353                region.set_role(RegionRole::Leader);
354
355                return Ok(region);
356            }
357            Ok(None) => {
358                debug!(
359                    "No data under directory {}, region_id: {}",
360                    region_dir, self.region_id
361                );
362            }
363            Err(e) => {
364                warn!(e;
365                    "Failed to open region {} before creating it, region_dir: {}",
366                    self.region_id, region_dir
367                );
368            }
369        }
370        // Safety: must be set before calling this method.
371        let mut options = self.options.take().unwrap();
372        let object_store = get_object_store(&options.storage, &self.object_store_manager)?;
373        let provider = self.provider::<S>(&options.wal_options)?;
374        let metadata = Arc::new(metadata);
375        // Sets the sst_format based on options or flat_format flag
376        let sst_format = if let Some(format) = options.sst_format {
377            format
378        } else if config.default_flat_format {
379            options.sst_format = Some(FormatType::Flat);
380            FormatType::Flat
381        } else {
382            // Default to PrimaryKeyParquet for newly created regions
383            options.sst_format = Some(FormatType::PrimaryKey);
384            FormatType::PrimaryKey
385        };
386        ensure_json2_not_use_time_series_memtable(&metadata, &options)?;
387        // Create a manifest manager for this region and writes regions to the manifest file.
388        let mut region_manifest_options =
389            RegionManifestOptions::new(config, &region_dir, &object_store);
390        // Set manifest cache if available
391        region_manifest_options.manifest_cache = self
392            .cache_manager
393            .as_ref()
394            .and_then(|cm| cm.write_cache())
395            .and_then(|wc| wc.manifest_cache());
396        // For remote WAL, we need to set flushed_entry_id to current topic's latest entry id.
397        // Kafka WAL allocation also carries the topic's pruned entry id as a create-time hint.
398        let flushed_entry_id = provider
399            .initial_flushed_entry_id::<S>(wal.store())
400            .max(initial_pruned_entry_id(&options.wal_options));
401        let manifest_manager = RegionManifestManager::new(
402            metadata.clone(),
403            flushed_entry_id,
404            region_manifest_options,
405            sst_format,
406            &self.stats,
407        )
408        .await?;
409
410        let memtable_builder = self.memtable_builder_provider.builder_for_options(&options);
411        let part_duration = options.compaction.time_window();
412        // Initial memtable id is 0.
413        let mutable = Arc::new(TimePartitions::new(
414            metadata.clone(),
415            memtable_builder.clone(),
416            0,
417            part_duration,
418        ));
419
420        debug!(
421            "Create region {} with options: {:?}, default_flat_format: {}",
422            region_id, options, config.default_flat_format
423        );
424
425        let version = VersionBuilder::new(metadata, mutable)
426            .options(options)
427            .flushed_entry_id(flushed_entry_id)
428            .build();
429        let version_control = Arc::new(VersionControl::new(version));
430        let access_layer = Arc::new(AccessLayer::new(
431            self.table_dir.clone(),
432            self.path_type,
433            object_store,
434            self.puffin_manager_factory,
435            self.intermediate_manager,
436        ));
437        let now = self.time_provider.current_time_millis();
438
439        Ok(Arc::new(MitoRegion {
440            region_id,
441            version_control,
442            series_index_version_control: Default::default(),
443            series_index_store: self.series_index_store.clone(),
444            access_layer: access_layer.clone(),
445            // Region is writable after it is created.
446            manifest_ctx: Arc::new(ManifestContext::new(
447                manifest_manager,
448                RegionRoleState::Leader(RegionLeaderState::Writable),
449                self.hook.clone(),
450            )),
451            file_purger: create_file_purger(
452                config.gc.enable,
453                self.path_type,
454                self.purge_scheduler,
455                access_layer,
456                self.cache_manager,
457                self.file_ref_manager.clone(),
458                self.series_index_store
459                    .map(|store| RangeIndexDeleter::new(store, region_id)),
460            ),
461            provider,
462            last_flush_millis: AtomicI64::new(now),
463            last_schedule_compaction_millis: AtomicI64::new(now),
464            time_provider: self.time_provider.clone(),
465            topic_latest_entry_id: AtomicU64::new(flushed_entry_id),
466            region_stats: RegionStats::new(),
467            stats: self.stats,
468        }))
469    }
470
471    /// Opens an existing region in read only mode.
472    ///
473    /// Returns error if the region doesn't exist.
474    pub(crate) async fn open<S: LogStore>(
475        mut self,
476        config: &MitoConfig,
477        wal: &Wal<S>,
478    ) -> Result<MitoRegionRef> {
479        let region_id = self.region_id;
480        let region_dir = self.region_dir();
481        let region = self
482            .maybe_open(config, wal)
483            .await?
484            .with_context(|| EmptyRegionDirSnafu {
485                region_id,
486                region_dir: &region_dir,
487            })?;
488
489        ensure!(
490            region.region_id == self.region_id,
491            RegionCorruptedSnafu {
492                region_id: self.region_id,
493                reason: format!(
494                    "recovered region has different region id {}",
495                    region.region_id
496                ),
497            }
498        );
499
500        Ok(region)
501    }
502
503    fn provider<S: LogStore>(&self, wal_options: &WalOptions) -> Result<Provider> {
504        provider_from_wal_options::<S>(self.region_id, wal_options)
505    }
506
507    /// Tries to open the region and returns `None` if the region directory is empty.
508    async fn maybe_open<S: LogStore>(
509        &mut self,
510        config: &MitoConfig,
511        wal: &Wal<S>,
512    ) -> Result<Option<MitoRegionRef>> {
513        let now = Instant::now();
514        let mut region_options = self.options.as_ref().unwrap().clone();
515        let object_storage = get_object_store(&region_options.storage, &self.object_store_manager)?;
516        let mut region_manifest_options =
517            RegionManifestOptions::new(config, &self.region_dir(), &object_storage);
518        // Set manifest cache if available
519        region_manifest_options.manifest_cache = self
520            .cache_manager
521            .as_ref()
522            .and_then(|cm| cm.write_cache())
523            .and_then(|wc| wc.manifest_cache());
524        let Some(manifest_manager) =
525            RegionManifestManager::open(region_manifest_options, &self.stats).await?
526        else {
527            return Ok(None);
528        };
529
530        // Backfill `partition_expr` if missing. Use the backfilled metadata in-memory during this open.
531        let manifest = manifest_manager.manifest();
532        let metadata = if manifest.metadata.partition_expr.is_none()
533            && let Some(expr_json) = self.partition_expr_fetcher.fetch_expr(self.region_id).await
534        {
535            let metadata = manifest.metadata.as_ref().clone();
536            let mut builder = RegionMetadataBuilder::from_existing(metadata);
537            builder.partition_expr_json(Some(expr_json));
538            Arc::new(builder.build().context(InvalidMetadataSnafu)?)
539        } else {
540            manifest.metadata.clone()
541        };
542        let metadata = maybe_upgrade_json2_layout(metadata)?;
543        // Updates the region options with the manifest.
544        sanitize_region_options(&manifest, &mut region_options);
545        ensure_json2_not_use_time_series_memtable(&metadata, &region_options)?;
546
547        let region_id = self.region_id;
548        let provider = self.provider::<S>(&region_options.wal_options)?;
549        let wal_entry_reader = self
550            .wal_entry_reader
551            .take()
552            .unwrap_or_else(|| wal.wal_entry_reader(&provider, region_id, None));
553        let on_region_opened = wal.on_region_opened();
554        let object_store = get_object_store(&region_options.storage, &self.object_store_manager)?;
555
556        debug!(
557            "Open region {} at {} with options: {:?}",
558            region_id, self.table_dir, self.options
559        );
560
561        let access_layer = Arc::new(AccessLayer::new(
562            self.table_dir.clone(),
563            self.path_type,
564            object_store,
565            self.puffin_manager_factory.clone(),
566            self.intermediate_manager.clone(),
567        ));
568        let file_purger = create_file_purger(
569            config.gc.enable,
570            self.path_type,
571            self.purge_scheduler.clone(),
572            access_layer.clone(),
573            self.cache_manager.clone(),
574            self.file_ref_manager.clone(),
575            self.series_index_store
576                .clone()
577                .map(|store| RangeIndexDeleter::new(store, region_id)),
578        );
579        // We should sanitize the region options before creating a new memtable.
580        let memtable_builder = self
581            .memtable_builder_provider
582            .builder_for_options(&region_options);
583        // Use compaction time window in the manifest if region doesn't provide
584        // the time window option.
585        let part_duration = region_options
586            .compaction
587            .time_window()
588            .or(manifest.compaction_time_window);
589        // Initial memtable id is 0.
590        let mutable = Arc::new(TimePartitions::new(
591            metadata.clone(),
592            memtable_builder.clone(),
593            0,
594            part_duration,
595        ));
596
597        // Updates region options by manifest before creating version.
598        let version_builder = version_builder_from_manifest(
599            &manifest,
600            metadata,
601            file_purger.clone(),
602            mutable,
603            region_options,
604        );
605        let version = version_builder.build();
606        let flushed_entry_id = version.flushed_entry_id;
607        let version_control = Arc::new(VersionControl::new(version));
608
609        let replay_from_entry_id = self
610            .replay_checkpoint
611            .unwrap_or_default()
612            .max(flushed_entry_id);
613        let topic_latest_entry_id = if !self.skip_wal_replay {
614            info!(
615                "Start replaying memtable at replay_from_entry_id: {} for region {}, manifest version: {}, flushed entry id: {}, elapsed: {:?}",
616                replay_from_entry_id,
617                region_id,
618                manifest.manifest_version,
619                flushed_entry_id,
620                now.elapsed()
621            );
622            replay_memtable(
623                &provider,
624                wal_entry_reader,
625                region_id,
626                replay_from_entry_id,
627                &version_control,
628                config.allow_stale_entries,
629                on_region_opened,
630            )
631            .await?;
632            // For remote WAL, we need to set topic_latest_entry_id to current topic's latest entry id.
633            // Only set after the WAL replay is completed.
634
635            if provider.is_remote_wal() && version_control.current().version.memtables.is_empty() {
636                wal.store()
637                    .latest_entry_id(&provider)
638                    .unwrap_or(replay_from_entry_id)
639            } else if provider.is_remote_wal() {
640                replay_from_entry_id
641            } else {
642                0
643            }
644        } else {
645            info!(
646                "Skip the WAL replay for region: {}, manifest version: {}, flushed_entry_id: {}, elapsed: {:?}",
647                region_id,
648                manifest.manifest_version,
649                flushed_entry_id,
650                now.elapsed()
651            );
652
653            if provider.is_remote_wal() {
654                replay_from_entry_id
655            } else {
656                0
657            }
658        };
659
660        if let Some(committed_in_manifest) = manifest.committed_sequence {
661            let committed_after_replay = version_control.committed_sequence();
662            if committed_in_manifest > committed_after_replay {
663                info!(
664                    "Overriding committed sequence, region: {}, flushed_sequence: {}, committed_sequence: {} -> {}",
665                    self.region_id,
666                    version_control.current().version.flushed_sequence,
667                    version_control.committed_sequence(),
668                    committed_in_manifest
669                );
670                version_control.set_committed_sequence(committed_in_manifest);
671            }
672        }
673
674        let now = self.time_provider.current_time_millis();
675
676        let series_index_version_control =
677            match (&self.series_index_store, &self.series_index_purger) {
678                (Some(store), Some(purger))
679                    if is_sparse_metric_metadata(&version_control.current().version.metadata) =>
680                {
681                    load_version_control(store, self.region_id, purger).await
682                }
683                _ => Default::default(),
684            };
685
686        let region = MitoRegion {
687            region_id: self.region_id,
688            version_control: version_control.clone(),
689            series_index_version_control,
690            series_index_store: self.series_index_store.clone(),
691            access_layer: access_layer.clone(),
692            // Region is always opened in read only mode.
693            manifest_ctx: Arc::new(ManifestContext::new(
694                manifest_manager,
695                RegionRoleState::Follower,
696                self.hook.clone(),
697            )),
698            file_purger,
699            provider: provider.clone(),
700            last_flush_millis: AtomicI64::new(now),
701            last_schedule_compaction_millis: AtomicI64::new(now),
702            time_provider: self.time_provider.clone(),
703            topic_latest_entry_id: AtomicU64::new(topic_latest_entry_id),
704            region_stats: RegionStats::new(),
705            stats: self.stats.clone(),
706        };
707
708        let region = Arc::new(region);
709
710        maybe_load_cache(&region, config, &self.cache_manager);
711        maybe_preload_parquet_meta_cache(&region, config, &self.cache_manager);
712
713        Ok(Some(region))
714    }
715}
716
717pub(crate) fn provider_from_wal_options<S: LogStore>(
718    region_id: RegionId,
719    wal_options: &WalOptions,
720) -> Result<Provider> {
721    match wal_options {
722        WalOptions::RaftEngine => {
723            ensure!(
724                TypeId::of::<RaftEngineLogStore>() == TypeId::of::<S>()
725                    || TypeId::of::<NoopLogStore>() == TypeId::of::<S>(),
726                error::IncompatibleWalProviderChangeSnafu {
727                    global: "`kafka`",
728                    region: "`raft_engine`",
729                }
730            );
731            Ok(Provider::raft_engine_provider(region_id.as_u64()))
732        }
733        WalOptions::Kafka(options) => {
734            ensure!(
735                TypeId::of::<KafkaLogStore>() == TypeId::of::<S>()
736                    || TypeId::of::<NoopLogStore>() == TypeId::of::<S>(),
737                error::IncompatibleWalProviderChangeSnafu {
738                    global: "`raft_engine`",
739                    region: "`kafka`",
740                }
741            );
742            Ok(Provider::kafka_provider(options.topic.clone()))
743        }
744        WalOptions::ObjectStore(options) => {
745            ensure!(
746                TypeId::of::<RaftEngineLogStore>() != TypeId::of::<S>(),
747                error::IncompatibleWalProviderChangeSnafu {
748                    global: "`raft_engine`",
749                    region: "`object_store`",
750                }
751            );
752            ensure!(
753                TypeId::of::<KafkaLogStore>() != TypeId::of::<S>(),
754                error::IncompatibleWalProviderChangeSnafu {
755                    global: "`kafka`",
756                    region: "`object_store`",
757                }
758            );
759            Ok(Provider::object_store_provider(
760                region_id,
761                options.prefix.clone(),
762            ))
763        }
764        WalOptions::Noop => Ok(Provider::noop_provider()),
765    }
766}
767
768#[cfg(not(feature = "test-shared-fs-region-migration"))]
769fn supports_open_region_object_storage_requirement(object_store: &ObjectStore) -> bool {
770    is_object_storage(object_store)
771}
772
773#[cfg(feature = "test-shared-fs-region-migration")]
774fn supports_open_region_object_storage_requirement(object_store: &ObjectStore) -> bool {
775    // Integration tests can configure multiple datanodes to share the same
776    // temporary home dir. That makes file storage accessible to all test
777    // datanodes, but production file storage still does not satisfy this
778    // requirement.
779    is_object_storage(object_store)
780        || object_store.info().scheme() == object_store::services::FS_SCHEME
781}
782
783/// Creates a version builder from a region manifest.
784pub(crate) fn version_builder_from_manifest(
785    manifest: &RegionManifest,
786    metadata: RegionMetadataRef,
787    file_purger: FilePurgerRef,
788    mutable: TimePartitionsRef,
789    region_options: RegionOptions,
790) -> VersionBuilder {
791    VersionBuilder::new(metadata, mutable)
792        .add_files(file_purger, manifest.files.values().cloned())
793        .flushed_entry_id(manifest.flushed_entry_id)
794        .flushed_sequence(manifest.flushed_sequence)
795        .truncated_entry_id(manifest.truncated_entry_id)
796        .compaction_time_window(manifest.compaction_time_window)
797        .options(region_options)
798}
799
800/// Updates region options by persistent options.
801pub(crate) fn sanitize_region_options(manifest: &RegionManifest, options: &mut RegionOptions) {
802    // The caller-supplied options win when they specify an explicit `sst_format`,
803    // so re-parsed options (e.g. bulk memtable forcing flat format) take effect on
804    // the running region instead of being clobbered by the persisted manifest value.
805    // Fall back to the manifest only when the caller did not provide a value.
806    match options.sst_format {
807        Some(format) if format != manifest.sst_format => {
808            common_telemetry::warn!(
809                "Overriding SST format from {:?} (manifest) to {:?} (options) for region {}",
810                manifest.sst_format,
811                format,
812                manifest.metadata.region_id,
813            );
814        }
815        Some(_) => {}
816        None => {
817            options.sst_format = Some(manifest.sst_format);
818        }
819    }
820    if let Some(manifest_append_mode) = manifest.append_mode
821        && options.append_mode != manifest_append_mode
822    {
823        common_telemetry::warn!(
824            "Overriding append_mode from {} to {} for region {}",
825            options.append_mode,
826            manifest_append_mode,
827            manifest.metadata.region_id,
828        );
829        options.append_mode = manifest_append_mode;
830    }
831    if options.append_mode && options.merge_mode.take().is_some() {
832        common_telemetry::warn!(
833            "Ignoring merge_mode because append_mode is enabled for region {}",
834            manifest.metadata.region_id,
835        );
836    }
837}
838
839/// Sanitizes open request options before parsing.
840pub(crate) fn sanitize_open_request_options(options: &mut HashMap<String, String>) {
841    let append_mode_enabled = options
842        .get("append_mode")
843        .is_some_and(|v| matches!(v.trim().to_ascii_lowercase().as_str(), "true" | "1"));
844
845    if append_mode_enabled && options.remove("merge_mode").is_some() {
846        common_telemetry::warn!(
847            "Ignoring merge_mode in open request options because append_mode is enabled"
848        );
849    }
850}
851
852/// Returns an object store corresponding to `name`. If `name` is `None`, this method returns the default object store.
853pub fn get_object_store(
854    name: &Option<String>,
855    object_store_manager: &ObjectStoreManagerRef,
856) -> Result<object_store::ObjectStore> {
857    if let Some(name) = name {
858        Ok(object_store_manager
859            .find(name)
860            .with_context(|| ObjectStoreNotFoundSnafu {
861                object_store: name.clone(),
862            })?
863            .clone())
864    } else {
865        Ok(object_store_manager.default_object_store().clone())
866    }
867}
868
869/// Checks whether the recovered region has the same schema as region to create.
870pub(crate) fn check_recovered_region(
871    recovered: &RegionMetadata,
872    region_id: RegionId,
873    column_metadatas: &[ColumnMetadata],
874    primary_key: &[ColumnId],
875) -> Result<()> {
876    if recovered.region_id != region_id {
877        error!(
878            "Recovered region {}, expect region {}",
879            recovered.region_id, region_id
880        );
881        return RegionCorruptedSnafu {
882            region_id,
883            reason: format!(
884                "recovered metadata has different region id {}",
885                recovered.region_id
886            ),
887        }
888        .fail();
889    }
890    if recovered.column_metadatas != column_metadatas {
891        error!(
892            "Unexpected schema in recovered region {}, recovered: {:?}, expect: {:?}",
893            recovered.region_id, recovered.column_metadatas, column_metadatas
894        );
895
896        return RegionCorruptedSnafu {
897            region_id,
898            reason: "recovered metadata has different schema",
899        }
900        .fail();
901    }
902    if recovered.primary_key != primary_key {
903        error!(
904            "Unexpected primary key in recovered region {}, recovered: {:?}, expect: {:?}",
905            recovered.region_id, recovered.primary_key, primary_key
906        );
907
908        return RegionCorruptedSnafu {
909            region_id,
910            reason: "recovered metadata has different primary key",
911        }
912        .fail();
913    }
914
915    Ok(())
916}
917
918/// Replays the mutations from WAL and inserts mutations to memtable of given region.
919pub(crate) async fn replay_memtable<F>(
920    provider: &Provider,
921    mut wal_entry_reader: Box<dyn WalEntryReader>,
922    region_id: RegionId,
923    flushed_entry_id: EntryId,
924    version_control: &VersionControlRef,
925    allow_stale_entries: bool,
926    on_region_opened: F,
927) -> Result<EntryId>
928where
929    F: FnOnce(RegionId, EntryId, &Provider) -> BoxFuture<Result<()>> + Send,
930{
931    let now = Instant::now();
932    let mut rows_replayed = 0;
933    // Last entry id should start from flushed entry id since there might be no
934    // data in the WAL.
935    let mut last_entry_id = flushed_entry_id;
936    let replay_from_entry_id = flushed_entry_id + 1;
937    let region_metadata = version_control.current().version.metadata.clone();
938
939    let mut wal_stream = wal_entry_reader.read(provider, replay_from_entry_id)?;
940    while let Some(res) = wal_stream.next().await {
941        let (entry_id, entry) = res?;
942        if entry_id <= flushed_entry_id {
943            warn!(
944                "Stale WAL entries read during replay, region id: {}, flushed entry id: {}, entry id read: {}",
945                region_id, flushed_entry_id, entry_id
946            );
947            ensure!(
948                allow_stale_entries,
949                StaleLogEntrySnafu {
950                    region_id,
951                    flushed_entry_id,
952                    unexpected_entry_id: entry_id,
953                }
954            );
955        }
956        last_entry_id = last_entry_id.max(entry_id);
957
958        let mut region_write_ctx = RegionWriteCtx::new(
959            region_id,
960            version_control,
961            provider.clone(),
962            // For WAL replay, we don't need to track the write bytes rate.
963            None,
964        );
965        for mutation in entry.mutations {
966            rows_replayed += mutation
967                .rows
968                .as_ref()
969                .map(|rows| rows.rows.len())
970                .unwrap_or(0);
971            region_write_ctx.push_mutation(
972                mutation.op_type,
973                mutation.rows,
974                mutation.write_hint,
975                OptionOutputTx::none(),
976                // We should respect the sequence in WAL during replay.
977                Some(mutation.sequence),
978                false,
979            );
980        }
981
982        for bulk_entry in entry.bulk_entries {
983            let mut part = BulkPart::try_from(bulk_entry)?;
984            // The entry may miss columns added by a concurrent alter if it was
985            // written by a writer with a stale schema (older versions kept the
986            // stale raw data when filling missing columns). Fills missing
987            // columns like the write path, otherwise the memtable rejects the
988            // batch. Skips sparse batches as they don't carry tag columns by
989            // design.
990            if region_metadata.primary_key_encoding != PrimaryKeyEncoding::Sparse {
991                part.fill_missing_columns(&region_metadata)?;
992            }
993            rows_replayed += part.num_rows();
994            // During replay, we should adopt the sequence from WAL.
995            let bulk_sequence_from_wal = part.sequence;
996            ensure!(
997                region_write_ctx.push_bulk(
998                    OptionOutputTx::none(),
999                    part,
1000                    Some(bulk_sequence_from_wal),
1001                    false
1002                ),
1003                RegionCorruptedSnafu {
1004                    region_id,
1005                    reason: "unable to replay memtable with bulk entries",
1006                }
1007            );
1008        }
1009
1010        // set next_entry_id and write to memtable.
1011        region_write_ctx.set_next_entry_id(last_entry_id + 1);
1012        region_write_ctx.write_memtable().await;
1013        region_write_ctx.write_bulk().await;
1014        // Publish the replayed sequences only after all rows (including bulk
1015        // parts) are installed, matching the write path ordering.
1016        region_write_ctx.publish_sequence_and_entry_id();
1017    }
1018
1019    // TODO(weny): We need to update `flushed_entry_id` in the region manifest
1020    // to avoid reading potentially incomplete entries in the future.
1021    (on_region_opened)(region_id, flushed_entry_id, provider).await?;
1022
1023    let series_count = version_control.current().series_count();
1024    info!(
1025        "Replay WAL for region: {}, provider: {:?}, rows recovered: {}, replay from entry id: {}, last entry id: {}, total timeseries replayed: {}, elapsed: {:?}",
1026        region_id,
1027        provider,
1028        rows_replayed,
1029        replay_from_entry_id,
1030        last_entry_id,
1031        series_count,
1032        now.elapsed()
1033    );
1034    Ok(last_entry_id)
1035}
1036
1037/// A task to load and fill the region file cache.
1038pub(crate) struct RegionLoadCacheTask {
1039    region: MitoRegionRef,
1040}
1041
1042impl RegionLoadCacheTask {
1043    pub(crate) fn new(region: MitoRegionRef) -> Self {
1044        Self { region }
1045    }
1046
1047    /// Fills the file cache with index files from the region.
1048    pub(crate) async fn fill_cache(&self, file_cache: &FileCache) {
1049        let region_id = self.region.region_id;
1050        let table_dir = self.region.access_layer.table_dir();
1051        let path_type = self.region.access_layer.path_type();
1052        let object_store = self.region.access_layer.object_store();
1053        let version_control = &self.region.version_control;
1054
1055        // Collects IndexKeys, file sizes, and max timestamps for files that need to be downloaded
1056        let mut files_to_download = Vec::new();
1057        let mut files_already_cached = 0;
1058
1059        {
1060            let version = version_control.current().version;
1061            for level in version.ssts.levels() {
1062                for file_handle in level.files.values() {
1063                    let file_meta = file_handle.meta_ref();
1064                    if file_meta.exists_index() {
1065                        let puffin_key = IndexKey::new(
1066                            file_meta.region_id,
1067                            file_meta.file_id,
1068                            FileType::Puffin(file_meta.index_version),
1069                        );
1070
1071                        if !file_cache.contains_key(&puffin_key) {
1072                            files_to_download.push((
1073                                puffin_key,
1074                                file_meta.index_file_size,
1075                                file_meta.time_range.1, // max timestamp
1076                            ));
1077                        } else {
1078                            files_already_cached += 1;
1079                        }
1080                    }
1081                }
1082            }
1083            // Releases the Version after the scope to avoid holding the memtables and file handles
1084            // for a long time.
1085        }
1086
1087        // Sorts files by max timestamp in descending order to loads latest files first
1088        files_to_download.sort_by_key(|b| std::cmp::Reverse(b.2));
1089
1090        let total_files = files_to_download.len() as i64;
1091
1092        info!(
1093            "Starting background index cache preload for region {}, total_files_to_download: {}, files_already_cached: {}",
1094            region_id, total_files, files_already_cached
1095        );
1096
1097        CACHE_FILL_PENDING_FILES.add(total_files);
1098
1099        let mut files_downloaded = 0;
1100        let mut files_skipped = 0;
1101
1102        for (puffin_key, file_size, max_timestamp) in files_to_download {
1103            let current_size = file_cache.puffin_cache_size();
1104            let capacity = file_cache.puffin_cache_capacity();
1105            let region_state = self.region.state();
1106            if !can_load_cache(region_state) {
1107                info!(
1108                    "Stopping index cache by state: {:?}, region: {}, current_size: {}, capacity: {}",
1109                    region_state, region_id, current_size, capacity
1110                );
1111                break;
1112            }
1113
1114            // Checks if adding this file would exceed capacity
1115            if current_size + file_size > capacity {
1116                info!(
1117                    "Stopping index cache preload due to capacity limit, region: {}, file_id: {}, current_size: {}, file_size: {}, capacity: {}, file_timestamp: {:?}",
1118                    region_id, puffin_key.file_id, current_size, file_size, capacity, max_timestamp
1119                );
1120                files_skipped = (total_files - files_downloaded) as usize;
1121                CACHE_FILL_PENDING_FILES.sub(total_files - files_downloaded);
1122                break;
1123            }
1124
1125            let index_version = if let FileType::Puffin(version) = puffin_key.file_type {
1126                version
1127            } else {
1128                unreachable!("`files_to_download` should only contains Puffin files");
1129            };
1130            let index_id = RegionIndexId::new(
1131                RegionFileId::new(puffin_key.region_id, puffin_key.file_id),
1132                index_version,
1133            );
1134
1135            let index_remote_path = location::index_file_path(table_dir, index_id, path_type);
1136
1137            match file_cache
1138                .download(puffin_key, &index_remote_path, object_store, file_size)
1139                .await
1140            {
1141                Ok(_) => {
1142                    debug!(
1143                        "Downloaded index file to write cache, region: {}, file_id: {}",
1144                        region_id, puffin_key.file_id
1145                    );
1146                    files_downloaded += 1;
1147                    CACHE_FILL_DOWNLOADED_FILES.inc_by(1);
1148                    CACHE_FILL_PENDING_FILES.dec();
1149                }
1150                Err(e) => {
1151                    warn!(
1152                        e; "Failed to download index file to write cache, region: {}, file_id: {}",
1153                        region_id, puffin_key.file_id
1154                    );
1155                    CACHE_FILL_PENDING_FILES.dec();
1156                }
1157            }
1158        }
1159
1160        info!(
1161            "Completed background cache fill task for region {}, total_files: {}, files_downloaded: {}, files_already_cached: {}, files_skipped: {}",
1162            region_id, total_files, files_downloaded, files_already_cached, files_skipped
1163        );
1164    }
1165}
1166
1167/// Loads all index (Puffin) files from the version into the write cache.
1168fn maybe_load_cache(
1169    region: &MitoRegionRef,
1170    config: &MitoConfig,
1171    cache_manager: &Option<CacheManagerRef>,
1172) {
1173    let Some(cache_manager) = cache_manager else {
1174        return;
1175    };
1176    let Some(write_cache) = cache_manager.write_cache() else {
1177        return;
1178    };
1179
1180    let preload_enabled = config.preload_index_cache;
1181    if !preload_enabled {
1182        return;
1183    }
1184
1185    let task = RegionLoadCacheTask::new(region.clone());
1186    write_cache.load_region_cache(task);
1187}
1188
1189/// Preloads Parquet metadata into the in-memory SST meta cache on region open.
1190///
1191/// This improves the latency of the first query after server start by avoiding large Parquet
1192/// metadata reads on demand.
1193///
1194/// The preload is best-effort:
1195/// - Always tries to warm from the local write cache (file cache) first.
1196/// - If the region storage backend is local filesystem (`Scheme::Fs`), it may also load metadata
1197///   directly from the local store.
1198/// - It will not fetch metadata from remote object stores (S3/GCS/OSS/...).
1199#[allow(clippy::too_many_arguments)]
1200async fn preload_parquet_meta_cache_for_files(
1201    region_id: RegionId,
1202    cache_manager: CacheManagerRef,
1203    sst_meta_cache_capacity: u64,
1204    table_dir: String,
1205    path_type: PathType,
1206    object_store: object_store::ObjectStore,
1207    region_metadata: RegionMetadataRef,
1208    mut files: Vec<FileHandle>,
1209) -> usize {
1210    if !cache_manager.sst_meta_cache_enabled()
1211        || sst_meta_cache_capacity == 0
1212        || cache_manager.sst_meta_cache_is_full()
1213    {
1214        return 0;
1215    }
1216
1217    let allow_direct_load = object_store.info().scheme() == object_store::services::FS_SCHEME;
1218
1219    // Sort by time range so we can prefer preloading newer files first.
1220    files.sort_by_key(|b| std::cmp::Reverse(b.meta_ref().time_range.1));
1221
1222    let mut loaded = 0usize;
1223    for file_handle in files {
1224        // Stop when the protected authoritative tier is full.
1225        if cache_manager.sst_meta_cache_is_full() {
1226            break;
1227        }
1228
1229        let file_id = file_handle.file_id();
1230        let mut cache_metrics = MetadataCacheMetrics::default();
1231        if let Some(metadata) = cache_manager
1232            .get_compact_sst_meta_data(file_id, PageIndexPolicy::Optional)
1233            .await
1234        {
1235            if file_handle.raw_primary_key_range().is_none()
1236                && let Some(primary_key_range) = extract_primary_key_range(
1237                    metadata.parquet_metadata().as_ref(),
1238                    &region_metadata,
1239                )
1240            {
1241                file_handle.set_primary_key_range(primary_key_range);
1242            }
1243            continue;
1244        }
1245
1246        if let Some(write_cache) = cache_manager.write_cache() {
1247            let key = IndexKey::new(file_id.region_id(), file_id.file_id(), FileType::Parquet);
1248            if let Some(metadata) = write_cache
1249                .file_cache()
1250                .get_sst_meta_data(
1251                    key,
1252                    &mut cache_metrics,
1253                    PageIndexPolicy::Optional,
1254                    &common_runtime::global_runtime(),
1255                )
1256                .await
1257            {
1258                let decoded = metadata.decoded();
1259                if file_handle.raw_primary_key_range().is_none()
1260                    && let Some(primary_key_range) = extract_primary_key_range(
1261                        decoded.parquet_metadata().as_ref(),
1262                        &region_metadata,
1263                    )
1264                {
1265                    file_handle.set_primary_key_range(primary_key_range);
1266                }
1267                match metadata {
1268                    SstMetaPreparation::Prepared(metadata) => {
1269                        cache_manager.put_prepared_sst_meta(file_id, metadata, false);
1270                        loaded += 1;
1271                    }
1272                    SstMetaPreparation::DecodedOnly { encoding_error, .. } => warn!(
1273                        encoding_error;
1274                        "Failed to encode file-cached SST metadata during preload, region: {}, file: {}",
1275                        region_id,
1276                        file_id.file_id()
1277                    ),
1278                }
1279                continue;
1280            }
1281        }
1282
1283        if !allow_direct_load {
1284            continue;
1285        }
1286
1287        let file_size = file_handle.meta_ref().file_size;
1288        let file_path = file_handle.file_path(&table_dir, path_type);
1289        let mut loader = MetadataLoader::new(object_store.clone(), &file_path, file_size);
1290        loader.with_page_index_policy(PageIndexPolicy::Optional);
1291        match loader.load(&mut cache_metrics).await {
1292            Ok(metadata) => {
1293                if let Some(primary_key_range) =
1294                    extract_primary_key_range(&metadata, &region_metadata)
1295                {
1296                    file_handle.set_primary_key_range(primary_key_range);
1297                }
1298                match prepare_sst_meta(
1299                    &file_path,
1300                    metadata,
1301                    // The SST can predate schema changes, so decode its embedded region metadata
1302                    // instead of substituting the region's current schema.
1303                    None,
1304                    PageIndexPolicy::Optional,
1305                    &common_runtime::global_runtime(),
1306                )
1307                .await
1308                {
1309                    Ok(SstMetaPreparation::Prepared(metadata)) => {
1310                        // Preload the compact canonical entry only. Query traffic promotes decoded
1311                        // entries into the separate acceleration tier according to actual demand.
1312                        cache_manager.put_prepared_sst_meta(file_id, metadata, false);
1313                        loaded += 1;
1314                    }
1315                    Ok(SstMetaPreparation::DecodedOnly { encoding_error, .. }) => warn!(
1316                        encoding_error;
1317                        "Failed to encode preloaded parquet metadata, region: {}, file: {}",
1318                        region_id,
1319                        file_path
1320                    ),
1321                    Err(err) => {
1322                        warn!(
1323                            err; "Failed to prepare preloaded parquet metadata, region: {}, file: {}",
1324                            region_id, file_path
1325                        );
1326                    }
1327                }
1328            }
1329            Err(err) => {
1330                // Preloading is best-effort. Failure shouldn't affect region open.
1331                warn!(
1332                    err; "Failed to preload parquet metadata from local store, region: {}, file: {}",
1333                    region_id, file_path
1334                );
1335            }
1336        }
1337    }
1338
1339    loaded
1340}
1341
1342fn maybe_preload_parquet_meta_cache(
1343    region: &MitoRegionRef,
1344    config: &MitoConfig,
1345    cache_manager: &Option<CacheManagerRef>,
1346) {
1347    let Some(cache_manager) = cache_manager else {
1348        return;
1349    };
1350    if !cache_manager.sst_meta_cache_enabled() {
1351        return;
1352    }
1353
1354    // Skip if SST meta cache is disabled.
1355    if config.sst_meta_cache_size.as_bytes() == 0 {
1356        return;
1357    }
1358    if !config.preload_index_cache {
1359        return;
1360    }
1361
1362    let region = region.clone();
1363    let cache_manager = cache_manager.clone();
1364    let sst_meta_cache_capacity = config.sst_meta_cache_size.as_bytes();
1365
1366    tokio::spawn(async move {
1367        // Safety: semaphore must exist.
1368        let _permit = PARQUET_META_PRELOAD_SEMAPHORE.acquire().await.unwrap();
1369
1370        let region_id = region.region_id;
1371        let table_dir = region.access_layer.table_dir().to_string();
1372        let path_type = region.access_layer.path_type();
1373        let object_store = region.access_layer.object_store().clone();
1374        let region_metadata = region.version_control.current().version.metadata.clone();
1375
1376        // Collect SST files. Do not hold the version longer than needed.
1377        let mut files = Vec::new();
1378        {
1379            let version = region.version_control.current().version;
1380            for level in version.ssts.levels() {
1381                for file_handle in level.files.values() {
1382                    files.push(file_handle.clone());
1383                }
1384            }
1385        }
1386        let preloading_start = Instant::now();
1387        let loaded = preload_parquet_meta_cache_for_files(
1388            region_id,
1389            cache_manager,
1390            sst_meta_cache_capacity,
1391            table_dir,
1392            path_type,
1393            object_store,
1394            region_metadata,
1395            files,
1396        )
1397        .await;
1398        let preloading_cost = preloading_start.elapsed();
1399
1400        if loaded > 0 {
1401            info!(
1402                "Preloaded parquet metadata for region {}, loaded_files: {}, elapsed_ms: {}",
1403                region_id,
1404                loaded,
1405                preloading_cost.as_millis()
1406            );
1407        }
1408    });
1409}
1410
1411fn can_load_cache(state: RegionRoleState) -> bool {
1412    match state {
1413        RegionRoleState::Leader(RegionLeaderState::Writable)
1414        | RegionRoleState::Leader(RegionLeaderState::Staging)
1415        | RegionRoleState::Leader(RegionLeaderState::Altering)
1416        | RegionRoleState::Leader(RegionLeaderState::EnteringStaging)
1417        | RegionRoleState::Leader(RegionLeaderState::Editing)
1418        | RegionRoleState::Follower => true,
1419        // The region will be closed soon if it is downgrading.
1420        RegionRoleState::Leader(RegionLeaderState::Downgrading)
1421        | RegionRoleState::Leader(RegionLeaderState::Dropping)
1422        | RegionRoleState::Leader(RegionLeaderState::Truncating) => false,
1423    }
1424}
1425
1426#[cfg(test)]
1427mod tests {
1428    use std::collections::HashMap;
1429    use std::sync::Arc;
1430
1431    use arrow_schema::extension::ExtensionType;
1432    use common_base::readable_size::ReadableSize;
1433    use common_error::ext::WhateverResult;
1434    use common_test_util::temp_dir::create_temp_dir;
1435    use common_time::Timestamp;
1436    use common_wal::options::{KafkaWalOptions, ObjectStoreWalOptions, WalOptions};
1437    use datatypes::arrow::array::{ArrayRef, BinaryArray, Int64Array};
1438    use datatypes::arrow::record_batch::RecordBatch;
1439    use datatypes::extension::json::{Json2ExtensionType, JsonMetadata};
1440    use datatypes::json::JsonSettings;
1441    use datatypes::prelude::ConcreteDataType;
1442    use datatypes::schema::ColumnSchema;
1443    use datatypes::types::json_type::{JsonNativeType, JsonObjectType};
1444    use log_store::kafka::log_store::KafkaLogStore;
1445    use log_store::noop::log_store::NoopLogStore;
1446    use log_store::raft_engine::log_store::RaftEngineLogStore;
1447    use object_store::ObjectStore;
1448    use object_store::services::{Fs, Memory, S3};
1449    use parquet::arrow::ArrowWriter;
1450    use parquet::file::metadata::{KeyValue, PageIndexPolicy};
1451    use parquet::file::properties::WriterProperties;
1452    use store_api::logstore::provider::Provider;
1453    use store_api::metadata::RegionMetadataBuilder;
1454    use store_api::region_request::PathType;
1455    use store_api::storage::{FileId, RegionId};
1456
1457    use super::{
1458        initial_pruned_entry_id, maybe_upgrade_json2_layout, preload_parquet_meta_cache_for_files,
1459        provider_from_wal_options, sanitize_region_options,
1460        supports_open_region_object_storage_requirement,
1461    };
1462    use crate::cache::CacheManager;
1463    use crate::cache::file_cache::{FileType, IndexKey};
1464    use crate::error;
1465    use crate::manifest::action::{RegionManifest, RemovedFilesRecord};
1466    use crate::region::options::RegionOptions;
1467    use crate::sst::FormatType;
1468    use crate::sst::file::{FileHandle, FileMeta};
1469    use crate::sst::file_purger::NoopFilePurger;
1470    use crate::sst::parquet::PARQUET_METADATA_KEY;
1471    use crate::test_util::TestEnv;
1472    use crate::test_util::sst_util::sst_region_metadata;
1473
1474    fn build_test_manifest(sst_format: FormatType) -> RegionManifest {
1475        RegionManifest {
1476            metadata: Arc::new(sst_region_metadata()),
1477            files: HashMap::new(),
1478            removed_files: RemovedFilesRecord::default(),
1479            flushed_entry_id: 0,
1480            flushed_sequence: 0,
1481            committed_sequence: None,
1482            manifest_version: 0,
1483            truncated_entry_id: None,
1484            compaction_time_window: None,
1485            sst_format,
1486            append_mode: None,
1487        }
1488    }
1489
1490    fn build_fs_object_store() -> ObjectStore {
1491        ObjectStore::new(Fs::default().root("/tmp")).unwrap()
1492    }
1493
1494    #[test]
1495    fn test_provider_from_object_store_wal_options() {
1496        let region_id = RegionId::new(1, 2);
1497        let wal_options = WalOptions::ObjectStore(ObjectStoreWalOptions::new("wal".to_string()));
1498
1499        let provider = provider_from_wal_options::<NoopLogStore>(region_id, &wal_options).unwrap();
1500        assert_eq!(
1501            Provider::object_store_provider(region_id, "wal".to_string()),
1502            provider
1503        );
1504
1505        let err =
1506            provider_from_wal_options::<RaftEngineLogStore>(region_id, &wal_options).unwrap_err();
1507        assert!(matches!(
1508            err,
1509            error::Error::IncompatibleWalProviderChange { .. }
1510        ));
1511        let err = provider_from_wal_options::<KafkaLogStore>(region_id, &wal_options).unwrap_err();
1512        assert!(matches!(
1513            err,
1514            error::Error::IncompatibleWalProviderChange { .. }
1515        ));
1516    }
1517
1518    #[test]
1519    fn test_initial_pruned_entry_id() {
1520        assert_eq!(0, initial_pruned_entry_id(&WalOptions::RaftEngine));
1521        assert_eq!(0, initial_pruned_entry_id(&WalOptions::Noop));
1522        assert_eq!(
1523            0,
1524            initial_pruned_entry_id(&WalOptions::ObjectStore(ObjectStoreWalOptions::new(
1525                "wal".to_string()
1526            )))
1527        );
1528        assert_eq!(
1529            0,
1530            initial_pruned_entry_id(&WalOptions::Kafka(KafkaWalOptions::new(
1531                "test_topic".to_string()
1532            )))
1533        );
1534        assert_eq!(
1535            42,
1536            initial_pruned_entry_id(&WalOptions::Kafka(KafkaWalOptions {
1537                topic: "test_topic".to_string(),
1538                initial_pruned_entry_id: Some(42),
1539            }))
1540        );
1541    }
1542
1543    #[test]
1544    fn test_upgrade_json2_layout() -> WhateverResult<()> {
1545        let settings = JsonSettings::try_new(vec![], Some(3))?;
1546        let extension = Json2ExtensionType::new(Arc::new(JsonMetadata::new_v1(settings.clone())));
1547        let mut column = ColumnSchema::new(
1548            "field_0",
1549            ConcreteDataType::json2(JsonNativeType::Object(JsonObjectType::new())),
1550            true,
1551        );
1552        column.with_extension_type(&extension);
1553
1554        let mut metadata = sst_region_metadata();
1555        metadata.column_metadatas[2].column_schema = column;
1556        let builder = RegionMetadataBuilder::from_existing(metadata);
1557        let metadata = Arc::new(builder.build()?);
1558
1559        let upgraded = maybe_upgrade_json2_layout(metadata)?;
1560        let column = &upgraded.column_metadatas[2].column_schema;
1561        let extension = column.extension_type::<Json2ExtensionType>()?.unwrap();
1562        assert!(extension.metadata().is_version_2());
1563        assert_eq!(&settings, extension.metadata().json_settings());
1564
1565        let arrow_schema = upgraded.schema.arrow_schema();
1566        let field = arrow_schema.field_with_name("field_0").unwrap();
1567        let extension = field.try_extension_type::<Json2ExtensionType>().unwrap();
1568        assert!(extension.metadata().is_version_2());
1569
1570        let unchanged = maybe_upgrade_json2_layout(upgraded.clone())?;
1571        assert!(Arc::ptr_eq(&upgraded, &unchanged));
1572        Ok(())
1573    }
1574
1575    #[test]
1576    #[cfg(not(feature = "test-shared-fs-region-migration"))]
1577    fn test_open_requirement_rejects_fs_object_store() {
1578        let object_store = build_fs_object_store();
1579
1580        assert!(!supports_open_region_object_storage_requirement(
1581            &object_store
1582        ));
1583    }
1584
1585    #[test]
1586    #[cfg(feature = "test-shared-fs-region-migration")]
1587    fn test_open_requirement_accepts_shared_fs_object_store_for_tests() {
1588        let object_store = build_fs_object_store();
1589
1590        assert!(supports_open_region_object_storage_requirement(
1591            &object_store
1592        ));
1593    }
1594
1595    #[test]
1596    fn test_open_requirement_accepts_s3_object_store() {
1597        let object_store = ObjectStore::new(
1598            S3::default()
1599                .bucket("test-bucket")
1600                .region("us-east-1")
1601                .disable_ec2_metadata(),
1602        )
1603        .unwrap();
1604
1605        assert!(supports_open_region_object_storage_requirement(
1606            &object_store
1607        ));
1608    }
1609
1610    #[test]
1611    fn test_sanitize_region_options_options_format_wins() {
1612        // Manifest persisted PrimaryKey, but the re-parsed options now request Flat
1613        // (e.g., bulk memtable). The options' value must win.
1614        let manifest = build_test_manifest(FormatType::PrimaryKey);
1615        let mut options = RegionOptions {
1616            sst_format: Some(FormatType::Flat),
1617            ..Default::default()
1618        };
1619        sanitize_region_options(&manifest, &mut options);
1620        assert_eq!(options.sst_format, Some(FormatType::Flat));
1621    }
1622
1623    #[test]
1624    fn test_sanitize_region_options_fills_from_manifest_when_unset() {
1625        // When the caller didn't specify sst_format, fall back to the manifest value
1626        // so downstream code always sees a concrete format.
1627        let manifest = build_test_manifest(FormatType::Flat);
1628        let mut options = RegionOptions {
1629            sst_format: None,
1630            ..Default::default()
1631        };
1632        sanitize_region_options(&manifest, &mut options);
1633        assert_eq!(options.sst_format, Some(FormatType::Flat));
1634    }
1635
1636    fn sst_parquet_bytes(batch: &RecordBatch) -> Vec<u8> {
1637        let key_value_meta = KeyValue::new(
1638            PARQUET_METADATA_KEY.to_string(),
1639            sst_region_metadata().to_json().unwrap(),
1640        );
1641        let props = WriterProperties::builder()
1642            .set_key_value_metadata(Some(vec![key_value_meta]))
1643            .build();
1644
1645        let mut parquet_bytes = Vec::new();
1646        let mut writer =
1647            ArrowWriter::try_new(&mut parquet_bytes, batch.schema(), Some(props)).unwrap();
1648        writer.write(batch).unwrap();
1649        writer.close().unwrap();
1650
1651        parquet_bytes
1652    }
1653
1654    #[tokio::test]
1655    async fn test_preload_parquet_meta_cache_uses_file_cache() {
1656        let env = TestEnv::new().await;
1657
1658        let local_store = ObjectStore::new(Memory::default()).unwrap();
1659        let write_cache = env
1660            .create_write_cache(local_store, ReadableSize::mb(1024))
1661            .await;
1662        let cache_manager = Arc::new(
1663            CacheManager::builder()
1664                .sst_meta_cache_size(1024 * 1024)
1665                .write_cache(Some(write_cache.clone()))
1666                .build(),
1667        );
1668
1669        let region_id = RegionId::new(1, 1);
1670        let file_id = FileId::random();
1671
1672        let col = Arc::new(Int64Array::from_iter_values([1, 2, 3])) as ArrayRef;
1673        let keys =
1674            ["a", "b", "c"].map(|tag| crate::test_util::sst_util::new_primary_key(&[tag, ""]));
1675        let primary_key = Arc::new(BinaryArray::from_iter_values(&keys)) as ArrayRef;
1676        let batch = RecordBatch::try_from_iter([
1677            ("col", col),
1678            (
1679                store_api::storage::consts::PRIMARY_KEY_COLUMN_NAME,
1680                primary_key,
1681            ),
1682        ])
1683        .unwrap();
1684        let parquet_bytes = sst_parquet_bytes(&batch);
1685        let file_size = parquet_bytes.len() as u64;
1686
1687        let file_meta = FileMeta {
1688            region_id,
1689            file_id,
1690            time_range: (Timestamp::new_millisecond(0), Timestamp::new_millisecond(1)),
1691            level: 0,
1692            file_size,
1693            max_row_group_uncompressed_size: 0,
1694            available_indexes: Default::default(),
1695            indexes: vec![],
1696            index_file_size: 0,
1697            index_version: 0,
1698            num_rows: 3,
1699            num_row_groups: 1,
1700            sequence: None,
1701            partition_expr: None,
1702            num_series: 0,
1703            ..Default::default()
1704        };
1705        let file_handle = FileHandle::new(file_meta, Arc::new(NoopFilePurger));
1706
1707        let table_dir = "test_table";
1708        let path_type = PathType::Bare;
1709        let remote_path = file_handle.file_path(table_dir, path_type);
1710
1711        let source_store = ObjectStore::new(Memory::default()).unwrap();
1712        source_store
1713            .write(&remote_path, parquet_bytes)
1714            .await
1715            .unwrap();
1716
1717        // Put the parquet file into the write cache, so file cache contains metadata.
1718        let index_key = IndexKey::new(region_id, file_id, FileType::Parquet);
1719        write_cache
1720            .file_cache()
1721            .download(index_key, &remote_path, &source_store, file_size)
1722            .await
1723            .unwrap();
1724
1725        let region_file_id = file_handle.file_id();
1726        assert!(
1727            cache_manager
1728                .get_parquet_meta_data_from_mem_cache(region_file_id)
1729                .is_none()
1730        );
1731
1732        let loaded = preload_parquet_meta_cache_for_files(
1733            region_id,
1734            cache_manager.clone(),
1735            1024 * 1024,
1736            table_dir.to_string(),
1737            path_type,
1738            source_store.clone(),
1739            Arc::new(sst_region_metadata()),
1740            vec![file_handle.clone()],
1741        )
1742        .await;
1743
1744        // Should warm the in-memory cache from the local file cache.
1745        assert_eq!(loaded, 1);
1746        assert!(
1747            cache_manager
1748                .get_compact_sst_meta_data(
1749                    region_file_id,
1750                    parquet::file::metadata::PageIndexPolicy::Optional,
1751                )
1752                .await
1753                .is_some()
1754        );
1755        assert!(file_handle.raw_primary_key_range().is_some());
1756    }
1757
1758    #[tokio::test]
1759    async fn test_preload_parquet_meta_cache_skips_files_not_in_file_cache() {
1760        let cache_manager = Arc::new(
1761            CacheManager::builder()
1762                .sst_meta_cache_size(1024 * 1024)
1763                .build(),
1764        );
1765
1766        let region_id = RegionId::new(1, 1);
1767        let file_id = FileId::random();
1768
1769        // Without a local file cache entry, preloading should skip the file.
1770        let file_meta = FileMeta {
1771            region_id,
1772            file_id,
1773            time_range: (Timestamp::new_millisecond(0), Timestamp::new_millisecond(1)),
1774            level: 0,
1775            file_size: 0,
1776            max_row_group_uncompressed_size: 0,
1777            available_indexes: Default::default(),
1778            indexes: vec![],
1779            index_file_size: 0,
1780            index_version: 0,
1781            num_rows: 3,
1782            num_row_groups: 1,
1783            sequence: None,
1784            partition_expr: None,
1785            num_series: 0,
1786            ..Default::default()
1787        };
1788        let file_handle = FileHandle::new(file_meta, Arc::new(NoopFilePurger));
1789
1790        let table_dir = "test_table";
1791        let path_type = PathType::Bare;
1792        let remote_path = file_handle.file_path(table_dir, path_type);
1793
1794        // Even if the remote object store has the file, we should not preload from it.
1795        let object_store = ObjectStore::new(Memory::default()).unwrap();
1796        object_store
1797            .write(&remote_path, b"noop".as_slice())
1798            .await
1799            .unwrap();
1800
1801        let region_file_id = file_handle.file_id();
1802        assert!(
1803            cache_manager
1804                .get_parquet_meta_data_from_mem_cache(region_file_id)
1805                .is_none()
1806        );
1807
1808        let loaded = preload_parquet_meta_cache_for_files(
1809            region_id,
1810            cache_manager.clone(),
1811            1024 * 1024,
1812            table_dir.to_string(),
1813            path_type,
1814            object_store,
1815            Arc::new(sst_region_metadata()),
1816            vec![file_handle],
1817        )
1818        .await;
1819
1820        assert_eq!(loaded, 0);
1821        assert!(
1822            cache_manager
1823                .get_parquet_meta_data_from_mem_cache(region_file_id)
1824                .is_none()
1825        );
1826    }
1827
1828    #[tokio::test]
1829    async fn test_preload_parquet_meta_cache_loads_from_local_fs() {
1830        let cache_manager = Arc::new(
1831            CacheManager::builder()
1832                .sst_meta_cache_size(1024 * 1024)
1833                .build(),
1834        );
1835
1836        let region_id = RegionId::new(1, 1);
1837        let file_id = FileId::random();
1838
1839        let col = Arc::new(Int64Array::from_iter_values([1, 2, 3])) as ArrayRef;
1840        let batch = RecordBatch::try_from_iter([("col", col)]).unwrap();
1841        let parquet_bytes = sst_parquet_bytes(&batch);
1842
1843        // file_size is 0 when it's missing/defaulted in manifests; MetadataLoader::load will stat
1844        // the local filesystem to retrieve it.
1845        let file_meta = FileMeta {
1846            region_id,
1847            file_id,
1848            time_range: (Timestamp::new_millisecond(0), Timestamp::new_millisecond(1)),
1849            level: 0,
1850            file_size: 0,
1851            max_row_group_uncompressed_size: 0,
1852            available_indexes: Default::default(),
1853            indexes: vec![],
1854            index_file_size: 0,
1855            index_version: 0,
1856            num_rows: 3,
1857            num_row_groups: 1,
1858            sequence: None,
1859            partition_expr: None,
1860            num_series: 0,
1861            ..Default::default()
1862        };
1863        let file_handle = FileHandle::new(file_meta, Arc::new(NoopFilePurger));
1864
1865        let table_dir = "test_table";
1866        let path_type = PathType::Bare;
1867        let file_path = file_handle.file_path(table_dir, path_type);
1868
1869        let root = create_temp_dir("parquet-meta-preload");
1870        let object_store =
1871            ObjectStore::new(Fs::default().root(root.path().to_str().unwrap())).unwrap();
1872        object_store.write(&file_path, parquet_bytes).await.unwrap();
1873
1874        let region_file_id = file_handle.file_id();
1875        assert!(
1876            cache_manager
1877                .get_parquet_meta_data_from_mem_cache(region_file_id)
1878                .is_none()
1879        );
1880
1881        let loaded = preload_parquet_meta_cache_for_files(
1882            region_id,
1883            cache_manager.clone(),
1884            1024 * 1024,
1885            table_dir.to_string(),
1886            path_type,
1887            object_store,
1888            Arc::new(sst_region_metadata()),
1889            vec![file_handle],
1890        )
1891        .await;
1892
1893        assert_eq!(loaded, 1);
1894        assert!(
1895            cache_manager
1896                .get_compact_sst_meta_data(region_file_id, PageIndexPolicy::Optional)
1897                .await
1898                .is_some()
1899        );
1900    }
1901}