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