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 => 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            access_layer: access_layer.clone(),
444            // Region is writable after it is created.
445            manifest_ctx: Arc::new(ManifestContext::new(
446                manifest_manager,
447                RegionRoleState::Leader(RegionLeaderState::Writable),
448                self.hook.clone(),
449            )),
450            file_purger: create_file_purger(
451                config.gc.enable,
452                self.path_type,
453                self.purge_scheduler,
454                access_layer,
455                self.cache_manager,
456                self.file_ref_manager.clone(),
457                self.series_index_store
458                    .map(|store| RangeIndexDeleter::new(store, region_id)),
459            ),
460            provider,
461            last_flush_millis: AtomicI64::new(now),
462            last_schedule_compaction_millis: AtomicI64::new(now),
463            time_provider: self.time_provider.clone(),
464            topic_latest_entry_id: AtomicU64::new(flushed_entry_id),
465            region_stats: RegionStats::new(),
466            stats: self.stats,
467        }))
468    }
469
470    /// Opens an existing region in read only mode.
471    ///
472    /// Returns error if the region doesn't exist.
473    pub(crate) async fn open<S: LogStore>(
474        mut self,
475        config: &MitoConfig,
476        wal: &Wal<S>,
477    ) -> Result<MitoRegionRef> {
478        let region_id = self.region_id;
479        let region_dir = self.region_dir();
480        let region = self
481            .maybe_open(config, wal)
482            .await?
483            .with_context(|| EmptyRegionDirSnafu {
484                region_id,
485                region_dir: &region_dir,
486            })?;
487
488        ensure!(
489            region.region_id == self.region_id,
490            RegionCorruptedSnafu {
491                region_id: self.region_id,
492                reason: format!(
493                    "recovered region has different region id {}",
494                    region.region_id
495                ),
496            }
497        );
498
499        Ok(region)
500    }
501
502    fn provider<S: LogStore>(&self, wal_options: &WalOptions) -> Result<Provider> {
503        provider_from_wal_options::<S>(self.region_id, wal_options)
504    }
505
506    /// Tries to open the region and returns `None` if the region directory is empty.
507    async fn maybe_open<S: LogStore>(
508        &mut self,
509        config: &MitoConfig,
510        wal: &Wal<S>,
511    ) -> Result<Option<MitoRegionRef>> {
512        let now = Instant::now();
513        let mut region_options = self.options.as_ref().unwrap().clone();
514        let object_storage = get_object_store(&region_options.storage, &self.object_store_manager)?;
515        let mut region_manifest_options =
516            RegionManifestOptions::new(config, &self.region_dir(), &object_storage);
517        // Set manifest cache if available
518        region_manifest_options.manifest_cache = self
519            .cache_manager
520            .as_ref()
521            .and_then(|cm| cm.write_cache())
522            .and_then(|wc| wc.manifest_cache());
523        let Some(manifest_manager) =
524            RegionManifestManager::open(region_manifest_options, &self.stats).await?
525        else {
526            return Ok(None);
527        };
528
529        // Backfill `partition_expr` if missing. Use the backfilled metadata in-memory during this open.
530        let manifest = manifest_manager.manifest();
531        let metadata = if manifest.metadata.partition_expr.is_none()
532            && let Some(expr_json) = self.partition_expr_fetcher.fetch_expr(self.region_id).await
533        {
534            let metadata = manifest.metadata.as_ref().clone();
535            let mut builder = RegionMetadataBuilder::from_existing(metadata);
536            builder.partition_expr_json(Some(expr_json));
537            Arc::new(builder.build().context(InvalidMetadataSnafu)?)
538        } else {
539            manifest.metadata.clone()
540        };
541        let metadata = maybe_upgrade_json2_layout(metadata)?;
542        // Updates the region options with the manifest.
543        sanitize_region_options(&manifest, &mut region_options);
544        ensure_json2_not_use_time_series_memtable(&metadata, &region_options)?;
545
546        let region_id = self.region_id;
547        let provider = self.provider::<S>(&region_options.wal_options)?;
548        let wal_entry_reader = self
549            .wal_entry_reader
550            .take()
551            .unwrap_or_else(|| wal.wal_entry_reader(&provider, region_id, None));
552        let on_region_opened = wal.on_region_opened();
553        let object_store = get_object_store(&region_options.storage, &self.object_store_manager)?;
554
555        debug!(
556            "Open region {} at {} with options: {:?}",
557            region_id, self.table_dir, self.options
558        );
559
560        let access_layer = Arc::new(AccessLayer::new(
561            self.table_dir.clone(),
562            self.path_type,
563            object_store,
564            self.puffin_manager_factory.clone(),
565            self.intermediate_manager.clone(),
566        ));
567        let file_purger = create_file_purger(
568            config.gc.enable,
569            self.path_type,
570            self.purge_scheduler.clone(),
571            access_layer.clone(),
572            self.cache_manager.clone(),
573            self.file_ref_manager.clone(),
574            self.series_index_store
575                .clone()
576                .map(|store| RangeIndexDeleter::new(store, region_id)),
577        );
578        // We should sanitize the region options before creating a new memtable.
579        let memtable_builder = self
580            .memtable_builder_provider
581            .builder_for_options(&region_options);
582        // Use compaction time window in the manifest if region doesn't provide
583        // the time window option.
584        let part_duration = region_options
585            .compaction
586            .time_window()
587            .or(manifest.compaction_time_window);
588        // Initial memtable id is 0.
589        let mutable = Arc::new(TimePartitions::new(
590            metadata.clone(),
591            memtable_builder.clone(),
592            0,
593            part_duration,
594        ));
595
596        // Updates region options by manifest before creating version.
597        let version_builder = version_builder_from_manifest(
598            &manifest,
599            metadata,
600            file_purger.clone(),
601            mutable,
602            region_options,
603        );
604        let version = version_builder.build();
605        let flushed_entry_id = version.flushed_entry_id;
606        let version_control = Arc::new(VersionControl::new(version));
607
608        let replay_from_entry_id = self
609            .replay_checkpoint
610            .unwrap_or_default()
611            .max(flushed_entry_id);
612        let topic_latest_entry_id = if !self.skip_wal_replay {
613            info!(
614                "Start replaying memtable at replay_from_entry_id: {} for region {}, manifest version: {}, flushed entry id: {}, elapsed: {:?}",
615                replay_from_entry_id,
616                region_id,
617                manifest.manifest_version,
618                flushed_entry_id,
619                now.elapsed()
620            );
621            replay_memtable(
622                &provider,
623                wal_entry_reader,
624                region_id,
625                replay_from_entry_id,
626                &version_control,
627                config.allow_stale_entries,
628                on_region_opened,
629            )
630            .await?;
631            // For remote WAL, we need to set topic_latest_entry_id to current topic's latest entry id.
632            // Only set after the WAL replay is completed.
633
634            if provider.is_remote_wal() && version_control.current().version.memtables.is_empty() {
635                wal.store()
636                    .latest_entry_id(&provider)
637                    .unwrap_or(replay_from_entry_id)
638            } else if provider.is_remote_wal() {
639                replay_from_entry_id
640            } else {
641                0
642            }
643        } else {
644            info!(
645                "Skip the WAL replay for region: {}, manifest version: {}, flushed_entry_id: {}, elapsed: {:?}",
646                region_id,
647                manifest.manifest_version,
648                flushed_entry_id,
649                now.elapsed()
650            );
651
652            if provider.is_remote_wal() {
653                replay_from_entry_id
654            } else {
655                0
656            }
657        };
658
659        if let Some(committed_in_manifest) = manifest.committed_sequence {
660            let committed_after_replay = version_control.committed_sequence();
661            if committed_in_manifest > committed_after_replay {
662                info!(
663                    "Overriding committed sequence, region: {}, flushed_sequence: {}, committed_sequence: {} -> {}",
664                    self.region_id,
665                    version_control.current().version.flushed_sequence,
666                    version_control.committed_sequence(),
667                    committed_in_manifest
668                );
669                version_control.set_committed_sequence(committed_in_manifest);
670            }
671        }
672
673        let now = self.time_provider.current_time_millis();
674
675        let series_index_version_control =
676            match (&self.series_index_store, &self.series_index_purger) {
677                (Some(store), Some(purger))
678                    if is_sparse_metric_metadata(&version_control.current().version.metadata) =>
679                {
680                    load_version_control(store, self.region_id, purger).await
681                }
682                _ => Default::default(),
683            };
684
685        let region = MitoRegion {
686            region_id: self.region_id,
687            version_control: version_control.clone(),
688            series_index_version_control,
689            access_layer: access_layer.clone(),
690            // Region is always opened in read only mode.
691            manifest_ctx: Arc::new(ManifestContext::new(
692                manifest_manager,
693                RegionRoleState::Follower,
694                self.hook.clone(),
695            )),
696            file_purger,
697            provider: provider.clone(),
698            last_flush_millis: AtomicI64::new(now),
699            last_schedule_compaction_millis: AtomicI64::new(now),
700            time_provider: self.time_provider.clone(),
701            topic_latest_entry_id: AtomicU64::new(topic_latest_entry_id),
702            region_stats: RegionStats::new(),
703            stats: self.stats.clone(),
704        };
705
706        let region = Arc::new(region);
707
708        maybe_load_cache(&region, config, &self.cache_manager);
709        maybe_preload_parquet_meta_cache(&region, config, &self.cache_manager);
710
711        Ok(Some(region))
712    }
713}
714
715pub(crate) fn provider_from_wal_options<S: LogStore>(
716    region_id: RegionId,
717    wal_options: &WalOptions,
718) -> Result<Provider> {
719    match wal_options {
720        WalOptions::RaftEngine => {
721            ensure!(
722                TypeId::of::<RaftEngineLogStore>() == TypeId::of::<S>()
723                    || TypeId::of::<NoopLogStore>() == TypeId::of::<S>(),
724                error::IncompatibleWalProviderChangeSnafu {
725                    global: "`kafka`",
726                    region: "`raft_engine`",
727                }
728            );
729            Ok(Provider::raft_engine_provider(region_id.as_u64()))
730        }
731        WalOptions::Kafka(options) => {
732            ensure!(
733                TypeId::of::<KafkaLogStore>() == TypeId::of::<S>()
734                    || TypeId::of::<NoopLogStore>() == TypeId::of::<S>(),
735                error::IncompatibleWalProviderChangeSnafu {
736                    global: "`raft_engine`",
737                    region: "`kafka`",
738                }
739            );
740            Ok(Provider::kafka_provider(options.topic.clone()))
741        }
742        WalOptions::Noop => Ok(Provider::noop_provider()),
743    }
744}
745
746#[cfg(not(feature = "test-shared-fs-region-migration"))]
747fn supports_open_region_object_storage_requirement(object_store: &ObjectStore) -> bool {
748    is_object_storage(object_store)
749}
750
751#[cfg(feature = "test-shared-fs-region-migration")]
752fn supports_open_region_object_storage_requirement(object_store: &ObjectStore) -> bool {
753    // Integration tests can configure multiple datanodes to share the same
754    // temporary home dir. That makes file storage accessible to all test
755    // datanodes, but production file storage still does not satisfy this
756    // requirement.
757    is_object_storage(object_store)
758        || object_store.info().scheme() == object_store::services::FS_SCHEME
759}
760
761/// Creates a version builder from a region manifest.
762pub(crate) fn version_builder_from_manifest(
763    manifest: &RegionManifest,
764    metadata: RegionMetadataRef,
765    file_purger: FilePurgerRef,
766    mutable: TimePartitionsRef,
767    region_options: RegionOptions,
768) -> VersionBuilder {
769    VersionBuilder::new(metadata, mutable)
770        .add_files(file_purger, manifest.files.values().cloned())
771        .flushed_entry_id(manifest.flushed_entry_id)
772        .flushed_sequence(manifest.flushed_sequence)
773        .truncated_entry_id(manifest.truncated_entry_id)
774        .compaction_time_window(manifest.compaction_time_window)
775        .options(region_options)
776}
777
778/// Updates region options by persistent options.
779pub(crate) fn sanitize_region_options(manifest: &RegionManifest, options: &mut RegionOptions) {
780    // The caller-supplied options win when they specify an explicit `sst_format`,
781    // so re-parsed options (e.g. bulk memtable forcing flat format) take effect on
782    // the running region instead of being clobbered by the persisted manifest value.
783    // Fall back to the manifest only when the caller did not provide a value.
784    match options.sst_format {
785        Some(format) if format != manifest.sst_format => {
786            common_telemetry::warn!(
787                "Overriding SST format from {:?} (manifest) to {:?} (options) for region {}",
788                manifest.sst_format,
789                format,
790                manifest.metadata.region_id,
791            );
792        }
793        Some(_) => {}
794        None => {
795            options.sst_format = Some(manifest.sst_format);
796        }
797    }
798    if let Some(manifest_append_mode) = manifest.append_mode
799        && options.append_mode != manifest_append_mode
800    {
801        common_telemetry::warn!(
802            "Overriding append_mode from {} to {} for region {}",
803            options.append_mode,
804            manifest_append_mode,
805            manifest.metadata.region_id,
806        );
807        options.append_mode = manifest_append_mode;
808    }
809    if options.append_mode && options.merge_mode.take().is_some() {
810        common_telemetry::warn!(
811            "Ignoring merge_mode because append_mode is enabled for region {}",
812            manifest.metadata.region_id,
813        );
814    }
815}
816
817/// Sanitizes open request options before parsing.
818pub(crate) fn sanitize_open_request_options(options: &mut HashMap<String, String>) {
819    let append_mode_enabled = options
820        .get("append_mode")
821        .is_some_and(|v| matches!(v.trim().to_ascii_lowercase().as_str(), "true" | "1"));
822
823    if append_mode_enabled && options.remove("merge_mode").is_some() {
824        common_telemetry::warn!(
825            "Ignoring merge_mode in open request options because append_mode is enabled"
826        );
827    }
828}
829
830/// Returns an object store corresponding to `name`. If `name` is `None`, this method returns the default object store.
831pub fn get_object_store(
832    name: &Option<String>,
833    object_store_manager: &ObjectStoreManagerRef,
834) -> Result<object_store::ObjectStore> {
835    if let Some(name) = name {
836        Ok(object_store_manager
837            .find(name)
838            .with_context(|| ObjectStoreNotFoundSnafu {
839                object_store: name.clone(),
840            })?
841            .clone())
842    } else {
843        Ok(object_store_manager.default_object_store().clone())
844    }
845}
846
847/// Checks whether the recovered region has the same schema as region to create.
848pub(crate) fn check_recovered_region(
849    recovered: &RegionMetadata,
850    region_id: RegionId,
851    column_metadatas: &[ColumnMetadata],
852    primary_key: &[ColumnId],
853) -> Result<()> {
854    if recovered.region_id != region_id {
855        error!(
856            "Recovered region {}, expect region {}",
857            recovered.region_id, region_id
858        );
859        return RegionCorruptedSnafu {
860            region_id,
861            reason: format!(
862                "recovered metadata has different region id {}",
863                recovered.region_id
864            ),
865        }
866        .fail();
867    }
868    if recovered.column_metadatas != column_metadatas {
869        error!(
870            "Unexpected schema in recovered region {}, recovered: {:?}, expect: {:?}",
871            recovered.region_id, recovered.column_metadatas, column_metadatas
872        );
873
874        return RegionCorruptedSnafu {
875            region_id,
876            reason: "recovered metadata has different schema",
877        }
878        .fail();
879    }
880    if recovered.primary_key != primary_key {
881        error!(
882            "Unexpected primary key in recovered region {}, recovered: {:?}, expect: {:?}",
883            recovered.region_id, recovered.primary_key, primary_key
884        );
885
886        return RegionCorruptedSnafu {
887            region_id,
888            reason: "recovered metadata has different primary key",
889        }
890        .fail();
891    }
892
893    Ok(())
894}
895
896/// Replays the mutations from WAL and inserts mutations to memtable of given region.
897pub(crate) async fn replay_memtable<F>(
898    provider: &Provider,
899    mut wal_entry_reader: Box<dyn WalEntryReader>,
900    region_id: RegionId,
901    flushed_entry_id: EntryId,
902    version_control: &VersionControlRef,
903    allow_stale_entries: bool,
904    on_region_opened: F,
905) -> Result<EntryId>
906where
907    F: FnOnce(RegionId, EntryId, &Provider) -> BoxFuture<Result<()>> + Send,
908{
909    let now = Instant::now();
910    let mut rows_replayed = 0;
911    // Last entry id should start from flushed entry id since there might be no
912    // data in the WAL.
913    let mut last_entry_id = flushed_entry_id;
914    let replay_from_entry_id = flushed_entry_id + 1;
915    let region_metadata = version_control.current().version.metadata.clone();
916
917    let mut wal_stream = wal_entry_reader.read(provider, replay_from_entry_id)?;
918    while let Some(res) = wal_stream.next().await {
919        let (entry_id, entry) = res?;
920        if entry_id <= flushed_entry_id {
921            warn!(
922                "Stale WAL entries read during replay, region id: {}, flushed entry id: {}, entry id read: {}",
923                region_id, flushed_entry_id, entry_id
924            );
925            ensure!(
926                allow_stale_entries,
927                StaleLogEntrySnafu {
928                    region_id,
929                    flushed_entry_id,
930                    unexpected_entry_id: entry_id,
931                }
932            );
933        }
934        last_entry_id = last_entry_id.max(entry_id);
935
936        let mut region_write_ctx = RegionWriteCtx::new(
937            region_id,
938            version_control,
939            provider.clone(),
940            // For WAL replay, we don't need to track the write bytes rate.
941            None,
942        );
943        for mutation in entry.mutations {
944            rows_replayed += mutation
945                .rows
946                .as_ref()
947                .map(|rows| rows.rows.len())
948                .unwrap_or(0);
949            region_write_ctx.push_mutation(
950                mutation.op_type,
951                mutation.rows,
952                mutation.write_hint,
953                OptionOutputTx::none(),
954                // We should respect the sequence in WAL during replay.
955                Some(mutation.sequence),
956            );
957        }
958
959        for bulk_entry in entry.bulk_entries {
960            let mut part = BulkPart::try_from(bulk_entry)?;
961            // The entry may miss columns added by a concurrent alter if it was
962            // written by a writer with a stale schema (older versions kept the
963            // stale raw data when filling missing columns). Fills missing
964            // columns like the write path, otherwise the memtable rejects the
965            // batch. Skips sparse batches as they don't carry tag columns by
966            // design.
967            if region_metadata.primary_key_encoding != PrimaryKeyEncoding::Sparse {
968                part.fill_missing_columns(&region_metadata)?;
969            }
970            rows_replayed += part.num_rows();
971            // During replay, we should adopt the sequence from WAL.
972            let bulk_sequence_from_wal = part.sequence;
973            ensure!(
974                region_write_ctx.push_bulk(
975                    OptionOutputTx::none(),
976                    part,
977                    Some(bulk_sequence_from_wal)
978                ),
979                RegionCorruptedSnafu {
980                    region_id,
981                    reason: "unable to replay memtable with bulk entries",
982                }
983            );
984        }
985
986        // set next_entry_id and write to memtable.
987        region_write_ctx.set_next_entry_id(last_entry_id + 1);
988        region_write_ctx.write_memtable().await;
989        region_write_ctx.write_bulk().await;
990        // Publish the replayed sequences only after all rows (including bulk
991        // parts) are installed, matching the write path ordering.
992        region_write_ctx.publish_sequence_and_entry_id();
993    }
994
995    // TODO(weny): We need to update `flushed_entry_id` in the region manifest
996    // to avoid reading potentially incomplete entries in the future.
997    (on_region_opened)(region_id, flushed_entry_id, provider).await?;
998
999    let series_count = version_control.current().series_count();
1000    info!(
1001        "Replay WAL for region: {}, provider: {:?}, rows recovered: {}, replay from entry id: {}, last entry id: {}, total timeseries replayed: {}, elapsed: {:?}",
1002        region_id,
1003        provider,
1004        rows_replayed,
1005        replay_from_entry_id,
1006        last_entry_id,
1007        series_count,
1008        now.elapsed()
1009    );
1010    Ok(last_entry_id)
1011}
1012
1013/// A task to load and fill the region file cache.
1014pub(crate) struct RegionLoadCacheTask {
1015    region: MitoRegionRef,
1016}
1017
1018impl RegionLoadCacheTask {
1019    pub(crate) fn new(region: MitoRegionRef) -> Self {
1020        Self { region }
1021    }
1022
1023    /// Fills the file cache with index files from the region.
1024    pub(crate) async fn fill_cache(&self, file_cache: &FileCache) {
1025        let region_id = self.region.region_id;
1026        let table_dir = self.region.access_layer.table_dir();
1027        let path_type = self.region.access_layer.path_type();
1028        let object_store = self.region.access_layer.object_store();
1029        let version_control = &self.region.version_control;
1030
1031        // Collects IndexKeys, file sizes, and max timestamps for files that need to be downloaded
1032        let mut files_to_download = Vec::new();
1033        let mut files_already_cached = 0;
1034
1035        {
1036            let version = version_control.current().version;
1037            for level in version.ssts.levels() {
1038                for file_handle in level.files.values() {
1039                    let file_meta = file_handle.meta_ref();
1040                    if file_meta.exists_index() {
1041                        let puffin_key = IndexKey::new(
1042                            file_meta.region_id,
1043                            file_meta.file_id,
1044                            FileType::Puffin(file_meta.index_version),
1045                        );
1046
1047                        if !file_cache.contains_key(&puffin_key) {
1048                            files_to_download.push((
1049                                puffin_key,
1050                                file_meta.index_file_size,
1051                                file_meta.time_range.1, // max timestamp
1052                            ));
1053                        } else {
1054                            files_already_cached += 1;
1055                        }
1056                    }
1057                }
1058            }
1059            // Releases the Version after the scope to avoid holding the memtables and file handles
1060            // for a long time.
1061        }
1062
1063        // Sorts files by max timestamp in descending order to loads latest files first
1064        files_to_download.sort_by_key(|b| std::cmp::Reverse(b.2));
1065
1066        let total_files = files_to_download.len() as i64;
1067
1068        info!(
1069            "Starting background index cache preload for region {}, total_files_to_download: {}, files_already_cached: {}",
1070            region_id, total_files, files_already_cached
1071        );
1072
1073        CACHE_FILL_PENDING_FILES.add(total_files);
1074
1075        let mut files_downloaded = 0;
1076        let mut files_skipped = 0;
1077
1078        for (puffin_key, file_size, max_timestamp) in files_to_download {
1079            let current_size = file_cache.puffin_cache_size();
1080            let capacity = file_cache.puffin_cache_capacity();
1081            let region_state = self.region.state();
1082            if !can_load_cache(region_state) {
1083                info!(
1084                    "Stopping index cache by state: {:?}, region: {}, current_size: {}, capacity: {}",
1085                    region_state, region_id, current_size, capacity
1086                );
1087                break;
1088            }
1089
1090            // Checks if adding this file would exceed capacity
1091            if current_size + file_size > capacity {
1092                info!(
1093                    "Stopping index cache preload due to capacity limit, region: {}, file_id: {}, current_size: {}, file_size: {}, capacity: {}, file_timestamp: {:?}",
1094                    region_id, puffin_key.file_id, current_size, file_size, capacity, max_timestamp
1095                );
1096                files_skipped = (total_files - files_downloaded) as usize;
1097                CACHE_FILL_PENDING_FILES.sub(total_files - files_downloaded);
1098                break;
1099            }
1100
1101            let index_version = if let FileType::Puffin(version) = puffin_key.file_type {
1102                version
1103            } else {
1104                unreachable!("`files_to_download` should only contains Puffin files");
1105            };
1106            let index_id = RegionIndexId::new(
1107                RegionFileId::new(puffin_key.region_id, puffin_key.file_id),
1108                index_version,
1109            );
1110
1111            let index_remote_path = location::index_file_path(table_dir, index_id, path_type);
1112
1113            match file_cache
1114                .download(puffin_key, &index_remote_path, object_store, file_size)
1115                .await
1116            {
1117                Ok(_) => {
1118                    debug!(
1119                        "Downloaded index file to write cache, region: {}, file_id: {}",
1120                        region_id, puffin_key.file_id
1121                    );
1122                    files_downloaded += 1;
1123                    CACHE_FILL_DOWNLOADED_FILES.inc_by(1);
1124                    CACHE_FILL_PENDING_FILES.dec();
1125                }
1126                Err(e) => {
1127                    warn!(
1128                        e; "Failed to download index file to write cache, region: {}, file_id: {}",
1129                        region_id, puffin_key.file_id
1130                    );
1131                    CACHE_FILL_PENDING_FILES.dec();
1132                }
1133            }
1134        }
1135
1136        info!(
1137            "Completed background cache fill task for region {}, total_files: {}, files_downloaded: {}, files_already_cached: {}, files_skipped: {}",
1138            region_id, total_files, files_downloaded, files_already_cached, files_skipped
1139        );
1140    }
1141}
1142
1143/// Loads all index (Puffin) files from the version into the write cache.
1144fn maybe_load_cache(
1145    region: &MitoRegionRef,
1146    config: &MitoConfig,
1147    cache_manager: &Option<CacheManagerRef>,
1148) {
1149    let Some(cache_manager) = cache_manager else {
1150        return;
1151    };
1152    let Some(write_cache) = cache_manager.write_cache() else {
1153        return;
1154    };
1155
1156    let preload_enabled = config.preload_index_cache;
1157    if !preload_enabled {
1158        return;
1159    }
1160
1161    let task = RegionLoadCacheTask::new(region.clone());
1162    write_cache.load_region_cache(task);
1163}
1164
1165/// Preloads Parquet metadata into the in-memory SST meta cache on region open.
1166///
1167/// This improves the latency of the first query after server start by avoiding large Parquet
1168/// metadata reads on demand.
1169///
1170/// The preload is best-effort:
1171/// - Always tries to warm from the local write cache (file cache) first.
1172/// - If the region storage backend is local filesystem (`Scheme::Fs`), it may also load metadata
1173///   directly from the local store.
1174/// - It will not fetch metadata from remote object stores (S3/GCS/OSS/...).
1175#[allow(clippy::too_many_arguments)]
1176async fn preload_parquet_meta_cache_for_files(
1177    region_id: RegionId,
1178    cache_manager: CacheManagerRef,
1179    sst_meta_cache_capacity: u64,
1180    table_dir: String,
1181    path_type: PathType,
1182    object_store: object_store::ObjectStore,
1183    region_metadata: RegionMetadataRef,
1184    mut files: Vec<FileHandle>,
1185) -> usize {
1186    if !cache_manager.sst_meta_cache_enabled()
1187        || sst_meta_cache_capacity == 0
1188        || cache_manager.sst_meta_cache_is_full()
1189    {
1190        return 0;
1191    }
1192
1193    let allow_direct_load = object_store.info().scheme() == object_store::services::FS_SCHEME;
1194
1195    // Sort by time range so we can prefer preloading newer files first.
1196    files.sort_by_key(|b| std::cmp::Reverse(b.meta_ref().time_range.1));
1197
1198    let mut loaded = 0usize;
1199    for file_handle in files {
1200        // Stop when the protected authoritative tier is full.
1201        if cache_manager.sst_meta_cache_is_full() {
1202            break;
1203        }
1204
1205        let file_id = file_handle.file_id();
1206        let mut cache_metrics = MetadataCacheMetrics::default();
1207        if let Some(metadata) = cache_manager
1208            .get_compact_sst_meta_data(file_id, PageIndexPolicy::Optional)
1209            .await
1210        {
1211            if file_handle.primary_key_range().is_none()
1212                && let Some(primary_key_range) = extract_primary_key_range(
1213                    metadata.parquet_metadata().as_ref(),
1214                    &region_metadata,
1215                )
1216            {
1217                file_handle.set_primary_key_range(primary_key_range);
1218            }
1219            continue;
1220        }
1221
1222        if let Some(write_cache) = cache_manager.write_cache() {
1223            let key = IndexKey::new(file_id.region_id(), file_id.file_id(), FileType::Parquet);
1224            if let Some(metadata) = write_cache
1225                .file_cache()
1226                .get_sst_meta_data(key, &mut cache_metrics, PageIndexPolicy::Optional)
1227                .await
1228            {
1229                let decoded = metadata.decoded();
1230                if file_handle.primary_key_range().is_none()
1231                    && let Some(primary_key_range) = extract_primary_key_range(
1232                        decoded.parquet_metadata().as_ref(),
1233                        &region_metadata,
1234                    )
1235                {
1236                    file_handle.set_primary_key_range(primary_key_range);
1237                }
1238                match metadata {
1239                    SstMetaPreparation::Prepared(metadata) => {
1240                        cache_manager.put_prepared_sst_meta(file_id, metadata, false);
1241                        loaded += 1;
1242                    }
1243                    SstMetaPreparation::DecodedOnly { encoding_error, .. } => warn!(
1244                        encoding_error;
1245                        "Failed to encode file-cached SST metadata during preload, region: {}, file: {}",
1246                        region_id,
1247                        file_id.file_id()
1248                    ),
1249                }
1250                continue;
1251            }
1252        }
1253
1254        if !allow_direct_load {
1255            continue;
1256        }
1257
1258        let file_size = file_handle.meta_ref().file_size;
1259        let file_path = file_handle.file_path(&table_dir, path_type);
1260        let mut loader = MetadataLoader::new(object_store.clone(), &file_path, file_size);
1261        loader.with_page_index_policy(PageIndexPolicy::Optional);
1262        match loader.load(&mut cache_metrics).await {
1263            Ok(metadata) => {
1264                if let Some(primary_key_range) =
1265                    extract_primary_key_range(&metadata, &region_metadata)
1266                {
1267                    file_handle.set_primary_key_range(primary_key_range);
1268                }
1269                match prepare_sst_meta(
1270                    &file_path,
1271                    metadata,
1272                    // The SST can predate schema changes, so decode its embedded region metadata
1273                    // instead of substituting the region's current schema.
1274                    None,
1275                    PageIndexPolicy::Optional,
1276                )
1277                .await
1278                {
1279                    Ok(SstMetaPreparation::Prepared(metadata)) => {
1280                        // Preload the compact canonical entry only. Query traffic promotes decoded
1281                        // entries into the separate acceleration tier according to actual demand.
1282                        cache_manager.put_prepared_sst_meta(file_id, metadata, false);
1283                        loaded += 1;
1284                    }
1285                    Ok(SstMetaPreparation::DecodedOnly { encoding_error, .. }) => warn!(
1286                        encoding_error;
1287                        "Failed to encode preloaded parquet metadata, region: {}, file: {}",
1288                        region_id,
1289                        file_path
1290                    ),
1291                    Err(err) => {
1292                        warn!(
1293                            err; "Failed to prepare preloaded parquet metadata, region: {}, file: {}",
1294                            region_id, file_path
1295                        );
1296                    }
1297                }
1298            }
1299            Err(err) => {
1300                // Preloading is best-effort. Failure shouldn't affect region open.
1301                warn!(
1302                    err; "Failed to preload parquet metadata from local store, region: {}, file: {}",
1303                    region_id, file_path
1304                );
1305            }
1306        }
1307    }
1308
1309    loaded
1310}
1311
1312fn maybe_preload_parquet_meta_cache(
1313    region: &MitoRegionRef,
1314    config: &MitoConfig,
1315    cache_manager: &Option<CacheManagerRef>,
1316) {
1317    let Some(cache_manager) = cache_manager else {
1318        return;
1319    };
1320    if !cache_manager.sst_meta_cache_enabled() {
1321        return;
1322    }
1323
1324    // Skip if SST meta cache is disabled.
1325    if config.sst_meta_cache_size.as_bytes() == 0 {
1326        return;
1327    }
1328    if !config.preload_index_cache {
1329        return;
1330    }
1331
1332    let region = region.clone();
1333    let cache_manager = cache_manager.clone();
1334    let sst_meta_cache_capacity = config.sst_meta_cache_size.as_bytes();
1335
1336    tokio::spawn(async move {
1337        // Safety: semaphore must exist.
1338        let _permit = PARQUET_META_PRELOAD_SEMAPHORE.acquire().await.unwrap();
1339
1340        let region_id = region.region_id;
1341        let table_dir = region.access_layer.table_dir().to_string();
1342        let path_type = region.access_layer.path_type();
1343        let object_store = region.access_layer.object_store().clone();
1344        let region_metadata = region.version_control.current().version.metadata.clone();
1345
1346        // Collect SST files. Do not hold the version longer than needed.
1347        let mut files = Vec::new();
1348        {
1349            let version = region.version_control.current().version;
1350            for level in version.ssts.levels() {
1351                for file_handle in level.files.values() {
1352                    files.push(file_handle.clone());
1353                }
1354            }
1355        }
1356        let preloading_start = Instant::now();
1357        let loaded = preload_parquet_meta_cache_for_files(
1358            region_id,
1359            cache_manager,
1360            sst_meta_cache_capacity,
1361            table_dir,
1362            path_type,
1363            object_store,
1364            region_metadata,
1365            files,
1366        )
1367        .await;
1368        let preloading_cost = preloading_start.elapsed();
1369
1370        if loaded > 0 {
1371            info!(
1372                "Preloaded parquet metadata for region {}, loaded_files: {}, elapsed_ms: {}",
1373                region_id,
1374                loaded,
1375                preloading_cost.as_millis()
1376            );
1377        }
1378    });
1379}
1380
1381fn can_load_cache(state: RegionRoleState) -> bool {
1382    match state {
1383        RegionRoleState::Leader(RegionLeaderState::Writable)
1384        | RegionRoleState::Leader(RegionLeaderState::Staging)
1385        | RegionRoleState::Leader(RegionLeaderState::Altering)
1386        | RegionRoleState::Leader(RegionLeaderState::EnteringStaging)
1387        | RegionRoleState::Leader(RegionLeaderState::Editing)
1388        | RegionRoleState::Follower => true,
1389        // The region will be closed soon if it is downgrading.
1390        RegionRoleState::Leader(RegionLeaderState::Downgrading)
1391        | RegionRoleState::Leader(RegionLeaderState::Dropping)
1392        | RegionRoleState::Leader(RegionLeaderState::Truncating) => false,
1393    }
1394}
1395
1396#[cfg(test)]
1397mod tests {
1398    use std::collections::HashMap;
1399    use std::sync::Arc;
1400
1401    use arrow_schema::extension::ExtensionType;
1402    use common_base::readable_size::ReadableSize;
1403    use common_error::ext::WhateverResult;
1404    use common_test_util::temp_dir::create_temp_dir;
1405    use common_time::Timestamp;
1406    use common_wal::options::{KafkaWalOptions, WalOptions};
1407    use datatypes::arrow::array::{ArrayRef, BinaryArray, Int64Array};
1408    use datatypes::arrow::record_batch::RecordBatch;
1409    use datatypes::extension::json::{Json2ExtensionType, JsonMetadata};
1410    use datatypes::json::JsonSettings;
1411    use datatypes::prelude::ConcreteDataType;
1412    use datatypes::schema::ColumnSchema;
1413    use datatypes::types::json_type::{JsonNativeType, JsonObjectType};
1414    use object_store::ObjectStore;
1415    use object_store::services::{Fs, Memory, S3};
1416    use parquet::arrow::ArrowWriter;
1417    use parquet::file::metadata::{KeyValue, PageIndexPolicy};
1418    use parquet::file::properties::WriterProperties;
1419    use store_api::metadata::RegionMetadataBuilder;
1420    use store_api::region_request::PathType;
1421    use store_api::storage::{FileId, RegionId};
1422
1423    use super::{
1424        initial_pruned_entry_id, maybe_upgrade_json2_layout, preload_parquet_meta_cache_for_files,
1425        sanitize_region_options, supports_open_region_object_storage_requirement,
1426    };
1427    use crate::cache::CacheManager;
1428    use crate::cache::file_cache::{FileType, IndexKey};
1429    use crate::manifest::action::{RegionManifest, RemovedFilesRecord};
1430    use crate::region::options::RegionOptions;
1431    use crate::sst::FormatType;
1432    use crate::sst::file::{FileHandle, FileMeta};
1433    use crate::sst::file_purger::NoopFilePurger;
1434    use crate::sst::parquet::PARQUET_METADATA_KEY;
1435    use crate::test_util::TestEnv;
1436    use crate::test_util::sst_util::sst_region_metadata;
1437
1438    fn build_test_manifest(sst_format: FormatType) -> RegionManifest {
1439        RegionManifest {
1440            metadata: Arc::new(sst_region_metadata()),
1441            files: HashMap::new(),
1442            removed_files: RemovedFilesRecord::default(),
1443            flushed_entry_id: 0,
1444            flushed_sequence: 0,
1445            committed_sequence: None,
1446            manifest_version: 0,
1447            truncated_entry_id: None,
1448            compaction_time_window: None,
1449            sst_format,
1450            append_mode: None,
1451        }
1452    }
1453
1454    fn build_fs_object_store() -> ObjectStore {
1455        ObjectStore::new(Fs::default().root("/tmp")).unwrap()
1456    }
1457
1458    #[test]
1459    fn test_initial_pruned_entry_id() {
1460        assert_eq!(0, initial_pruned_entry_id(&WalOptions::RaftEngine));
1461        assert_eq!(0, initial_pruned_entry_id(&WalOptions::Noop));
1462        assert_eq!(
1463            0,
1464            initial_pruned_entry_id(&WalOptions::Kafka(KafkaWalOptions::new(
1465                "test_topic".to_string()
1466            )))
1467        );
1468        assert_eq!(
1469            42,
1470            initial_pruned_entry_id(&WalOptions::Kafka(KafkaWalOptions {
1471                topic: "test_topic".to_string(),
1472                initial_pruned_entry_id: Some(42),
1473            }))
1474        );
1475    }
1476
1477    #[test]
1478    fn test_upgrade_json2_layout() -> WhateverResult<()> {
1479        let settings = JsonSettings::try_new(vec![], Some(3))?;
1480        let extension = Json2ExtensionType::new(Arc::new(JsonMetadata::new_v1(settings.clone())));
1481        let mut column = ColumnSchema::new(
1482            "field_0",
1483            ConcreteDataType::json2(JsonNativeType::Object(JsonObjectType::new())),
1484            true,
1485        );
1486        column.with_extension_type(&extension);
1487
1488        let mut metadata = sst_region_metadata();
1489        metadata.column_metadatas[2].column_schema = column;
1490        let builder = RegionMetadataBuilder::from_existing(metadata);
1491        let metadata = Arc::new(builder.build()?);
1492
1493        let upgraded = maybe_upgrade_json2_layout(metadata)?;
1494        let column = &upgraded.column_metadatas[2].column_schema;
1495        let extension = column.extension_type::<Json2ExtensionType>()?.unwrap();
1496        assert!(extension.metadata().is_version_2());
1497        assert_eq!(&settings, extension.metadata().json_settings());
1498
1499        let arrow_schema = upgraded.schema.arrow_schema();
1500        let field = arrow_schema.field_with_name("field_0").unwrap();
1501        let extension = field.try_extension_type::<Json2ExtensionType>().unwrap();
1502        assert!(extension.metadata().is_version_2());
1503
1504        let unchanged = maybe_upgrade_json2_layout(upgraded.clone())?;
1505        assert!(Arc::ptr_eq(&upgraded, &unchanged));
1506        Ok(())
1507    }
1508
1509    #[test]
1510    #[cfg(not(feature = "test-shared-fs-region-migration"))]
1511    fn test_open_requirement_rejects_fs_object_store() {
1512        let object_store = build_fs_object_store();
1513
1514        assert!(!supports_open_region_object_storage_requirement(
1515            &object_store
1516        ));
1517    }
1518
1519    #[test]
1520    #[cfg(feature = "test-shared-fs-region-migration")]
1521    fn test_open_requirement_accepts_shared_fs_object_store_for_tests() {
1522        let object_store = build_fs_object_store();
1523
1524        assert!(supports_open_region_object_storage_requirement(
1525            &object_store
1526        ));
1527    }
1528
1529    #[test]
1530    fn test_open_requirement_accepts_s3_object_store() {
1531        let object_store = ObjectStore::new(
1532            S3::default()
1533                .bucket("test-bucket")
1534                .region("us-east-1")
1535                .disable_ec2_metadata(),
1536        )
1537        .unwrap();
1538
1539        assert!(supports_open_region_object_storage_requirement(
1540            &object_store
1541        ));
1542    }
1543
1544    #[test]
1545    fn test_sanitize_region_options_options_format_wins() {
1546        // Manifest persisted PrimaryKey, but the re-parsed options now request Flat
1547        // (e.g., bulk memtable). The options' value must win.
1548        let manifest = build_test_manifest(FormatType::PrimaryKey);
1549        let mut options = RegionOptions {
1550            sst_format: Some(FormatType::Flat),
1551            ..Default::default()
1552        };
1553        sanitize_region_options(&manifest, &mut options);
1554        assert_eq!(options.sst_format, Some(FormatType::Flat));
1555    }
1556
1557    #[test]
1558    fn test_sanitize_region_options_fills_from_manifest_when_unset() {
1559        // When the caller didn't specify sst_format, fall back to the manifest value
1560        // so downstream code always sees a concrete format.
1561        let manifest = build_test_manifest(FormatType::Flat);
1562        let mut options = RegionOptions {
1563            sst_format: None,
1564            ..Default::default()
1565        };
1566        sanitize_region_options(&manifest, &mut options);
1567        assert_eq!(options.sst_format, Some(FormatType::Flat));
1568    }
1569
1570    fn sst_parquet_bytes(batch: &RecordBatch) -> Vec<u8> {
1571        let key_value_meta = KeyValue::new(
1572            PARQUET_METADATA_KEY.to_string(),
1573            sst_region_metadata().to_json().unwrap(),
1574        );
1575        let props = WriterProperties::builder()
1576            .set_key_value_metadata(Some(vec![key_value_meta]))
1577            .build();
1578
1579        let mut parquet_bytes = Vec::new();
1580        let mut writer =
1581            ArrowWriter::try_new(&mut parquet_bytes, batch.schema(), Some(props)).unwrap();
1582        writer.write(batch).unwrap();
1583        writer.close().unwrap();
1584
1585        parquet_bytes
1586    }
1587
1588    #[tokio::test]
1589    async fn test_preload_parquet_meta_cache_uses_file_cache() {
1590        let env = TestEnv::new().await;
1591
1592        let local_store = ObjectStore::new(Memory::default()).unwrap();
1593        let write_cache = env
1594            .create_write_cache(local_store, ReadableSize::mb(1024))
1595            .await;
1596        let cache_manager = Arc::new(
1597            CacheManager::builder()
1598                .sst_meta_cache_size(1024 * 1024)
1599                .write_cache(Some(write_cache.clone()))
1600                .build(),
1601        );
1602
1603        let region_id = RegionId::new(1, 1);
1604        let file_id = FileId::random();
1605
1606        let col = Arc::new(Int64Array::from_iter_values([1, 2, 3])) as ArrayRef;
1607        let primary_key = Arc::new(BinaryArray::from_iter_values([b"a", b"b", b"c"])) as ArrayRef;
1608        let batch = RecordBatch::try_from_iter([
1609            ("col", col),
1610            (
1611                store_api::storage::consts::PRIMARY_KEY_COLUMN_NAME,
1612                primary_key,
1613            ),
1614        ])
1615        .unwrap();
1616        let parquet_bytes = sst_parquet_bytes(&batch);
1617        let file_size = parquet_bytes.len() as u64;
1618
1619        let file_meta = FileMeta {
1620            region_id,
1621            file_id,
1622            time_range: (Timestamp::new_millisecond(0), Timestamp::new_millisecond(1)),
1623            level: 0,
1624            file_size,
1625            max_row_group_uncompressed_size: 0,
1626            available_indexes: Default::default(),
1627            indexes: vec![],
1628            index_file_size: 0,
1629            index_version: 0,
1630            num_rows: 3,
1631            num_row_groups: 1,
1632            sequence: None,
1633            partition_expr: None,
1634            num_series: 0,
1635            ..Default::default()
1636        };
1637        let file_handle = FileHandle::new(file_meta, Arc::new(NoopFilePurger));
1638
1639        let table_dir = "test_table";
1640        let path_type = PathType::Bare;
1641        let remote_path = file_handle.file_path(table_dir, path_type);
1642
1643        let source_store = ObjectStore::new(Memory::default()).unwrap();
1644        source_store
1645            .write(&remote_path, parquet_bytes)
1646            .await
1647            .unwrap();
1648
1649        // Put the parquet file into the write cache, so file cache contains metadata.
1650        let index_key = IndexKey::new(region_id, file_id, FileType::Parquet);
1651        write_cache
1652            .file_cache()
1653            .download(index_key, &remote_path, &source_store, file_size)
1654            .await
1655            .unwrap();
1656
1657        let region_file_id = file_handle.file_id();
1658        assert!(
1659            cache_manager
1660                .get_parquet_meta_data_from_mem_cache(region_file_id)
1661                .is_none()
1662        );
1663
1664        let loaded = preload_parquet_meta_cache_for_files(
1665            region_id,
1666            cache_manager.clone(),
1667            1024 * 1024,
1668            table_dir.to_string(),
1669            path_type,
1670            source_store.clone(),
1671            Arc::new(sst_region_metadata()),
1672            vec![file_handle.clone()],
1673        )
1674        .await;
1675
1676        // Should warm the in-memory cache from the local file cache.
1677        assert_eq!(loaded, 1);
1678        assert!(
1679            cache_manager
1680                .get_compact_sst_meta_data(
1681                    region_file_id,
1682                    parquet::file::metadata::PageIndexPolicy::Optional,
1683                )
1684                .await
1685                .is_some()
1686        );
1687        assert!(file_handle.primary_key_range().is_some());
1688    }
1689
1690    #[tokio::test]
1691    async fn test_preload_parquet_meta_cache_skips_files_not_in_file_cache() {
1692        let cache_manager = Arc::new(
1693            CacheManager::builder()
1694                .sst_meta_cache_size(1024 * 1024)
1695                .build(),
1696        );
1697
1698        let region_id = RegionId::new(1, 1);
1699        let file_id = FileId::random();
1700
1701        // Without a local file cache entry, preloading should skip the file.
1702        let file_meta = FileMeta {
1703            region_id,
1704            file_id,
1705            time_range: (Timestamp::new_millisecond(0), Timestamp::new_millisecond(1)),
1706            level: 0,
1707            file_size: 0,
1708            max_row_group_uncompressed_size: 0,
1709            available_indexes: Default::default(),
1710            indexes: vec![],
1711            index_file_size: 0,
1712            index_version: 0,
1713            num_rows: 3,
1714            num_row_groups: 1,
1715            sequence: None,
1716            partition_expr: None,
1717            num_series: 0,
1718            ..Default::default()
1719        };
1720        let file_handle = FileHandle::new(file_meta, Arc::new(NoopFilePurger));
1721
1722        let table_dir = "test_table";
1723        let path_type = PathType::Bare;
1724        let remote_path = file_handle.file_path(table_dir, path_type);
1725
1726        // Even if the remote object store has the file, we should not preload from it.
1727        let object_store = ObjectStore::new(Memory::default()).unwrap();
1728        object_store
1729            .write(&remote_path, b"noop".as_slice())
1730            .await
1731            .unwrap();
1732
1733        let region_file_id = file_handle.file_id();
1734        assert!(
1735            cache_manager
1736                .get_parquet_meta_data_from_mem_cache(region_file_id)
1737                .is_none()
1738        );
1739
1740        let loaded = preload_parquet_meta_cache_for_files(
1741            region_id,
1742            cache_manager.clone(),
1743            1024 * 1024,
1744            table_dir.to_string(),
1745            path_type,
1746            object_store,
1747            Arc::new(sst_region_metadata()),
1748            vec![file_handle],
1749        )
1750        .await;
1751
1752        assert_eq!(loaded, 0);
1753        assert!(
1754            cache_manager
1755                .get_parquet_meta_data_from_mem_cache(region_file_id)
1756                .is_none()
1757        );
1758    }
1759
1760    #[tokio::test]
1761    async fn test_preload_parquet_meta_cache_loads_from_local_fs() {
1762        let cache_manager = Arc::new(
1763            CacheManager::builder()
1764                .sst_meta_cache_size(1024 * 1024)
1765                .build(),
1766        );
1767
1768        let region_id = RegionId::new(1, 1);
1769        let file_id = FileId::random();
1770
1771        let col = Arc::new(Int64Array::from_iter_values([1, 2, 3])) as ArrayRef;
1772        let batch = RecordBatch::try_from_iter([("col", col)]).unwrap();
1773        let parquet_bytes = sst_parquet_bytes(&batch);
1774
1775        // file_size is 0 when it's missing/defaulted in manifests; MetadataLoader::load will stat
1776        // the local filesystem to retrieve it.
1777        let file_meta = FileMeta {
1778            region_id,
1779            file_id,
1780            time_range: (Timestamp::new_millisecond(0), Timestamp::new_millisecond(1)),
1781            level: 0,
1782            file_size: 0,
1783            max_row_group_uncompressed_size: 0,
1784            available_indexes: Default::default(),
1785            indexes: vec![],
1786            index_file_size: 0,
1787            index_version: 0,
1788            num_rows: 3,
1789            num_row_groups: 1,
1790            sequence: None,
1791            partition_expr: None,
1792            num_series: 0,
1793            ..Default::default()
1794        };
1795        let file_handle = FileHandle::new(file_meta, Arc::new(NoopFilePurger));
1796
1797        let table_dir = "test_table";
1798        let path_type = PathType::Bare;
1799        let file_path = file_handle.file_path(table_dir, path_type);
1800
1801        let root = create_temp_dir("parquet-meta-preload");
1802        let object_store =
1803            ObjectStore::new(Fs::default().root(root.path().to_str().unwrap())).unwrap();
1804        object_store.write(&file_path, parquet_bytes).await.unwrap();
1805
1806        let region_file_id = file_handle.file_id();
1807        assert!(
1808            cache_manager
1809                .get_parquet_meta_data_from_mem_cache(region_file_id)
1810                .is_none()
1811        );
1812
1813        let loaded = preload_parquet_meta_cache_for_files(
1814            region_id,
1815            cache_manager.clone(),
1816            1024 * 1024,
1817            table_dir.to_string(),
1818            path_type,
1819            object_store,
1820            Arc::new(sst_region_metadata()),
1821            vec![file_handle],
1822        )
1823        .await;
1824
1825        assert_eq!(loaded, 1);
1826        assert!(
1827            cache_manager
1828                .get_compact_sst_meta_data(region_file_id, PageIndexPolicy::Optional)
1829                .await
1830                .is_some()
1831        );
1832    }
1833}