1use std::any::TypeId;
18use std::collections::HashMap;
19use std::sync::atomic::{AtomicI64, AtomicU64};
20use std::sync::{Arc, LazyLock};
21use std::time::Instant;
22
23use arrow_schema::extension::ExtensionType;
24use common_telemetry::{debug, error, info, warn};
25use common_wal::options::WalOptions;
26use datatypes::extension::json::{Json2ExtensionType, JsonMetadata};
27use futures::StreamExt;
28use futures::future::BoxFuture;
29use log_store::kafka::log_store::KafkaLogStore;
30use log_store::noop::log_store::NoopLogStore;
31use log_store::raft_engine::log_store::RaftEngineLogStore;
32use object_store::ObjectStore;
33use object_store::manager::ObjectStoreManagerRef;
34use object_store::util::{is_object_storage, normalize_dir};
35use parquet::file::metadata::PageIndexPolicy;
36use snafu::{OptionExt, ResultExt, ensure};
37use store_api::codec::PrimaryKeyEncoding;
38use store_api::logstore::LogStore;
39use store_api::logstore::provider::Provider;
40use store_api::metadata::{
41 ColumnMetadata, RegionMetadata, RegionMetadataBuilder, RegionMetadataRef,
42};
43use store_api::region_engine::RegionRole;
44use store_api::region_request::{PathType, RegionRequirements};
45use store_api::storage::{ColumnId, RegionId};
46use tokio::sync::Semaphore;
47
48use crate::access_layer::AccessLayer;
49use crate::cache::file_cache::{FileCache, FileType, IndexKey};
50use crate::cache::{CacheManagerRef, SstMetaPreparation, prepare_sst_meta};
51use crate::config::MitoConfig;
52use crate::engine::region_hook::RegionHookRef;
53use crate::error;
54use crate::error::{
55 DataTypeMismatchSnafu, EmptyRegionDirSnafu, InvalidMetadataSnafu, InvalidRegionOptionsSnafu,
56 ObjectStoreNotFoundSnafu, RegionCorruptedSnafu, Result, StaleLogEntrySnafu,
57};
58use crate::manifest::action::RegionManifest;
59use crate::manifest::manager::{RegionManifestManager, RegionManifestOptions};
60use crate::memtable::bulk::part::BulkPart;
61use crate::memtable::time_partition::{TimePartitions, TimePartitionsRef};
62use crate::memtable::{MemtableBuilderProvider, ensure_json2_not_use_time_series_memtable};
63use crate::metrics::{CACHE_FILL_DOWNLOADED_FILES, CACHE_FILL_PENDING_FILES};
64use crate::read::series_candidate::is_sparse_metric_metadata;
65use crate::region::options::RegionOptions;
66use crate::region::version::{VersionBuilder, VersionControl, VersionControlRef};
67use crate::region::{
68 ManifestContext, ManifestStats, MitoRegion, MitoRegionRef, RegionLeaderState, RegionRoleState,
69 RegionStats,
70};
71use crate::region_write_ctx::RegionWriteCtx;
72use crate::request::OptionOutputTx;
73use crate::schedule::scheduler::SchedulerRef;
74use crate::series_index::{IndexFilePurger, load_version_control};
75use crate::sst::FormatType;
76use crate::sst::file::{FileHandle, RegionFileId, RegionIndexId};
77use crate::sst::file_purger::{FilePurgerRef, create_file_purger};
78use crate::sst::file_ref::FileReferenceManagerRef;
79use crate::sst::index::intermediate::IntermediateManager;
80use crate::sst::index::puffin_manager::PuffinManagerFactory;
81use crate::sst::location::{self, region_dir_from_table_dir};
82use crate::sst::parquet::metadata::{MetadataLoader, extract_primary_key_range};
83use crate::sst::parquet::reader::MetadataCacheMetrics;
84use crate::sst::range_index::RangeIndexDeleter;
85use crate::time_provider::TimeProviderRef;
86use crate::wal::entry_reader::WalEntryReader;
87use crate::wal::{EntryId, Wal};
88
89const PARQUET_META_PRELOAD_CONCURRENCY: usize = 8;
90
91static PARQUET_META_PRELOAD_SEMAPHORE: LazyLock<Semaphore> =
92 LazyLock::new(|| Semaphore::new(PARQUET_META_PRELOAD_CONCURRENCY));
93
94fn initial_pruned_entry_id(wal_options: &WalOptions) -> EntryId {
95 match wal_options {
96 WalOptions::Kafka(options) => options.initial_pruned_entry_id.unwrap_or(0),
97 WalOptions::RaftEngine | WalOptions::Noop | WalOptions::ObjectStore(_) => 0,
98 }
99}
100
101fn maybe_upgrade_json2_layout(metadata: RegionMetadataRef) -> Result<RegionMetadataRef> {
102 let mut upgrades = Vec::new();
103 for (index, column) in metadata.column_metadatas.iter().enumerate() {
104 if !column.column_schema.data_type.is_json2() {
105 continue;
106 }
107 let Some(extension) = column
108 .column_schema
109 .extension_type::<Json2ExtensionType>()
110 .context(DataTypeMismatchSnafu)?
111 else {
112 continue;
113 };
114 if extension.metadata().is_version_2() {
115 continue;
116 }
117 upgrades.push((index, extension.metadata().json_settings().clone()));
118 }
119
120 if upgrades.is_empty() {
121 return Ok(metadata);
122 }
123
124 let mut upgraded = metadata.as_ref().clone();
125 for (index, settings) in upgrades {
126 let extension = Json2ExtensionType::new(Arc::new(JsonMetadata::new(settings)));
127 upgraded.column_metadatas[index]
128 .column_schema
129 .with_extension_type(&extension);
130 }
131
132 let builder = RegionMetadataBuilder::from_existing(upgraded);
133 Ok(Arc::new(builder.build().context(InvalidMetadataSnafu)?))
134}
135
136#[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
148pub(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 #[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 pub(crate) fn series_index_store(mut self, store: Option<ObjectStore>) -> Self {
217 self.series_index_store = store;
218 self
219 }
220
221 pub(crate) fn series_index_purger(mut self, purger: Option<IndexFilePurger>) -> Self {
223 self.series_index_purger = purger;
224 self
225 }
226
227 pub(crate) fn hook(mut self, hook: Option<RegionHookRef>) -> Self {
229 self.hook = hook;
230 self
231 }
232
233 pub(crate) fn metadata_builder(mut self, builder: RegionMetadataBuilder) -> Self {
235 self.metadata_builder = Some(builder);
236 self
237 }
238
239 fn region_dir(&self) -> String {
241 region_dir_from_table_dir(&self.table_dir, self.region_id, self.path_type)
242 }
243
244 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 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 pub(crate) fn replay_checkpoint(mut self, replay_checkpoint: Option<u64>) -> Self {
267 self.replay_checkpoint = replay_checkpoint;
268 self
269 }
270
271 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 pub(crate) fn options(mut self, options: RegionOptions) -> Result<Self> {
283 options.validate()?;
284 self.options = Some(options);
285 Ok(self)
286 }
287
288 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 pub(crate) fn cache(mut self, cache_manager: Option<CacheManagerRef>) -> Self {
316 self.cache_manager = cache_manager;
317 self
318 }
319
320 pub(crate) fn skip_wal_replay(mut self, skip: bool) -> Self {
322 self.skip_wal_replay = skip;
323 self
324 }
325
326 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 match self.maybe_open(config, wal).await {
342 Ok(Some(region)) => {
343 let recovered = region.metadata();
344 let expect = &metadata;
346 check_recovered_region(
347 &recovered,
348 expect.region_id,
349 &expect.column_metadatas,
350 &expect.primary_key,
351 )?;
352 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 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 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 options.sst_format = Some(FormatType::PrimaryKey);
384 FormatType::PrimaryKey
385 };
386 ensure_json2_not_use_time_series_memtable(&metadata, &options)?;
387 let mut region_manifest_options =
389 RegionManifestOptions::new(config, ®ion_dir, &object_store);
390 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 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 let mutable = Arc::new(TimePartitions::new(
414 metadata.clone(),
415 memtable_builder.clone(),
416 0,
417 part_duration,
418 ));
419
420 debug!(
421 "Create region {} with options: {:?}, default_flat_format: {}",
422 region_id, options, config.default_flat_format
423 );
424
425 let version = VersionBuilder::new(metadata, mutable)
426 .options(options)
427 .flushed_entry_id(flushed_entry_id)
428 .build();
429 let version_control = Arc::new(VersionControl::new(version));
430 let access_layer = Arc::new(AccessLayer::new(
431 self.table_dir.clone(),
432 self.path_type,
433 object_store,
434 self.puffin_manager_factory,
435 self.intermediate_manager,
436 ));
437 let now = self.time_provider.current_time_millis();
438
439 Ok(Arc::new(MitoRegion {
440 region_id,
441 version_control,
442 series_index_version_control: Default::default(),
443 series_index_store: self.series_index_store.clone(),
444 access_layer: access_layer.clone(),
445 manifest_ctx: Arc::new(ManifestContext::new(
447 manifest_manager,
448 RegionRoleState::Leader(RegionLeaderState::Writable),
449 self.hook.clone(),
450 )),
451 file_purger: create_file_purger(
452 config.gc.enable,
453 self.path_type,
454 self.purge_scheduler,
455 access_layer,
456 self.cache_manager,
457 self.file_ref_manager.clone(),
458 self.series_index_store
459 .map(|store| RangeIndexDeleter::new(store, region_id)),
460 ),
461 provider,
462 last_flush_millis: AtomicI64::new(now),
463 last_schedule_compaction_millis: AtomicI64::new(now),
464 time_provider: self.time_provider.clone(),
465 topic_latest_entry_id: AtomicU64::new(flushed_entry_id),
466 region_stats: RegionStats::new(),
467 stats: self.stats,
468 }))
469 }
470
471 pub(crate) async fn open<S: LogStore>(
475 mut self,
476 config: &MitoConfig,
477 wal: &Wal<S>,
478 ) -> Result<MitoRegionRef> {
479 let region_id = self.region_id;
480 let region_dir = self.region_dir();
481 let region = self
482 .maybe_open(config, wal)
483 .await?
484 .with_context(|| EmptyRegionDirSnafu {
485 region_id,
486 region_dir: ®ion_dir,
487 })?;
488
489 ensure!(
490 region.region_id == self.region_id,
491 RegionCorruptedSnafu {
492 region_id: self.region_id,
493 reason: format!(
494 "recovered region has different region id {}",
495 region.region_id
496 ),
497 }
498 );
499
500 Ok(region)
501 }
502
503 fn provider<S: LogStore>(&self, wal_options: &WalOptions) -> Result<Provider> {
504 provider_from_wal_options::<S>(self.region_id, wal_options)
505 }
506
507 async fn maybe_open<S: LogStore>(
509 &mut self,
510 config: &MitoConfig,
511 wal: &Wal<S>,
512 ) -> Result<Option<MitoRegionRef>> {
513 let now = Instant::now();
514 let mut region_options = self.options.as_ref().unwrap().clone();
515 let object_storage = get_object_store(®ion_options.storage, &self.object_store_manager)?;
516 let mut region_manifest_options =
517 RegionManifestOptions::new(config, &self.region_dir(), &object_storage);
518 region_manifest_options.manifest_cache = self
520 .cache_manager
521 .as_ref()
522 .and_then(|cm| cm.write_cache())
523 .and_then(|wc| wc.manifest_cache());
524 let Some(manifest_manager) =
525 RegionManifestManager::open(region_manifest_options, &self.stats).await?
526 else {
527 return Ok(None);
528 };
529
530 let manifest = manifest_manager.manifest();
532 let metadata = if manifest.metadata.partition_expr.is_none()
533 && let Some(expr_json) = self.partition_expr_fetcher.fetch_expr(self.region_id).await
534 {
535 let metadata = manifest.metadata.as_ref().clone();
536 let mut builder = RegionMetadataBuilder::from_existing(metadata);
537 builder.partition_expr_json(Some(expr_json));
538 Arc::new(builder.build().context(InvalidMetadataSnafu)?)
539 } else {
540 manifest.metadata.clone()
541 };
542 let metadata = maybe_upgrade_json2_layout(metadata)?;
543 sanitize_region_options(&manifest, &mut region_options);
545 ensure_json2_not_use_time_series_memtable(&metadata, ®ion_options)?;
546
547 let region_id = self.region_id;
548 let provider = self.provider::<S>(®ion_options.wal_options)?;
549 let wal_entry_reader = self
550 .wal_entry_reader
551 .take()
552 .unwrap_or_else(|| wal.wal_entry_reader(&provider, region_id, None));
553 let on_region_opened = wal.on_region_opened();
554 let object_store = get_object_store(®ion_options.storage, &self.object_store_manager)?;
555
556 debug!(
557 "Open region {} at {} with options: {:?}",
558 region_id, self.table_dir, self.options
559 );
560
561 let access_layer = Arc::new(AccessLayer::new(
562 self.table_dir.clone(),
563 self.path_type,
564 object_store,
565 self.puffin_manager_factory.clone(),
566 self.intermediate_manager.clone(),
567 ));
568 let file_purger = create_file_purger(
569 config.gc.enable,
570 self.path_type,
571 self.purge_scheduler.clone(),
572 access_layer.clone(),
573 self.cache_manager.clone(),
574 self.file_ref_manager.clone(),
575 self.series_index_store
576 .clone()
577 .map(|store| RangeIndexDeleter::new(store, region_id)),
578 );
579 let memtable_builder = self
581 .memtable_builder_provider
582 .builder_for_options(®ion_options);
583 let part_duration = region_options
586 .compaction
587 .time_window()
588 .or(manifest.compaction_time_window);
589 let mutable = Arc::new(TimePartitions::new(
591 metadata.clone(),
592 memtable_builder.clone(),
593 0,
594 part_duration,
595 ));
596
597 let version_builder = version_builder_from_manifest(
599 &manifest,
600 metadata,
601 file_purger.clone(),
602 mutable,
603 region_options,
604 );
605 let version = version_builder.build();
606 let flushed_entry_id = version.flushed_entry_id;
607 let version_control = Arc::new(VersionControl::new(version));
608
609 let replay_from_entry_id = self
610 .replay_checkpoint
611 .unwrap_or_default()
612 .max(flushed_entry_id);
613 let topic_latest_entry_id = if !self.skip_wal_replay {
614 info!(
615 "Start replaying memtable at replay_from_entry_id: {} for region {}, manifest version: {}, flushed entry id: {}, elapsed: {:?}",
616 replay_from_entry_id,
617 region_id,
618 manifest.manifest_version,
619 flushed_entry_id,
620 now.elapsed()
621 );
622 replay_memtable(
623 &provider,
624 wal_entry_reader,
625 region_id,
626 replay_from_entry_id,
627 &version_control,
628 config.allow_stale_entries,
629 on_region_opened,
630 )
631 .await?;
632 if provider.is_remote_wal() && version_control.current().version.memtables.is_empty() {
636 wal.store()
637 .latest_entry_id(&provider)
638 .unwrap_or(replay_from_entry_id)
639 } else if provider.is_remote_wal() {
640 replay_from_entry_id
641 } else {
642 0
643 }
644 } else {
645 info!(
646 "Skip the WAL replay for region: {}, manifest version: {}, flushed_entry_id: {}, elapsed: {:?}",
647 region_id,
648 manifest.manifest_version,
649 flushed_entry_id,
650 now.elapsed()
651 );
652
653 if provider.is_remote_wal() {
654 replay_from_entry_id
655 } else {
656 0
657 }
658 };
659
660 if let Some(committed_in_manifest) = manifest.committed_sequence {
661 let committed_after_replay = version_control.committed_sequence();
662 if committed_in_manifest > committed_after_replay {
663 info!(
664 "Overriding committed sequence, region: {}, flushed_sequence: {}, committed_sequence: {} -> {}",
665 self.region_id,
666 version_control.current().version.flushed_sequence,
667 version_control.committed_sequence(),
668 committed_in_manifest
669 );
670 version_control.set_committed_sequence(committed_in_manifest);
671 }
672 }
673
674 let now = self.time_provider.current_time_millis();
675
676 let series_index_version_control =
677 match (&self.series_index_store, &self.series_index_purger) {
678 (Some(store), Some(purger))
679 if is_sparse_metric_metadata(&version_control.current().version.metadata) =>
680 {
681 load_version_control(store, self.region_id, purger).await
682 }
683 _ => Default::default(),
684 };
685
686 let region = MitoRegion {
687 region_id: self.region_id,
688 version_control: version_control.clone(),
689 series_index_version_control,
690 series_index_store: self.series_index_store.clone(),
691 access_layer: access_layer.clone(),
692 manifest_ctx: Arc::new(ManifestContext::new(
694 manifest_manager,
695 RegionRoleState::Follower,
696 self.hook.clone(),
697 )),
698 file_purger,
699 provider: provider.clone(),
700 last_flush_millis: AtomicI64::new(now),
701 last_schedule_compaction_millis: AtomicI64::new(now),
702 time_provider: self.time_provider.clone(),
703 topic_latest_entry_id: AtomicU64::new(topic_latest_entry_id),
704 region_stats: RegionStats::new(),
705 stats: self.stats.clone(),
706 };
707
708 let region = Arc::new(region);
709
710 maybe_load_cache(®ion, config, &self.cache_manager);
711 maybe_preload_parquet_meta_cache(®ion, config, &self.cache_manager);
712
713 Ok(Some(region))
714 }
715}
716
717pub(crate) fn provider_from_wal_options<S: LogStore>(
718 region_id: RegionId,
719 wal_options: &WalOptions,
720) -> Result<Provider> {
721 match wal_options {
722 WalOptions::RaftEngine => {
723 ensure!(
724 TypeId::of::<RaftEngineLogStore>() == TypeId::of::<S>()
725 || TypeId::of::<NoopLogStore>() == TypeId::of::<S>(),
726 error::IncompatibleWalProviderChangeSnafu {
727 global: "`kafka`",
728 region: "`raft_engine`",
729 }
730 );
731 Ok(Provider::raft_engine_provider(region_id.as_u64()))
732 }
733 WalOptions::Kafka(options) => {
734 ensure!(
735 TypeId::of::<KafkaLogStore>() == TypeId::of::<S>()
736 || TypeId::of::<NoopLogStore>() == TypeId::of::<S>(),
737 error::IncompatibleWalProviderChangeSnafu {
738 global: "`raft_engine`",
739 region: "`kafka`",
740 }
741 );
742 Ok(Provider::kafka_provider(options.topic.clone()))
743 }
744 WalOptions::ObjectStore(options) => {
745 ensure!(
746 TypeId::of::<RaftEngineLogStore>() != TypeId::of::<S>(),
747 error::IncompatibleWalProviderChangeSnafu {
748 global: "`raft_engine`",
749 region: "`object_store`",
750 }
751 );
752 ensure!(
753 TypeId::of::<KafkaLogStore>() != TypeId::of::<S>(),
754 error::IncompatibleWalProviderChangeSnafu {
755 global: "`kafka`",
756 region: "`object_store`",
757 }
758 );
759 Ok(Provider::object_store_provider(
760 region_id,
761 options.prefix.clone(),
762 ))
763 }
764 WalOptions::Noop => Ok(Provider::noop_provider()),
765 }
766}
767
768#[cfg(not(feature = "test-shared-fs-region-migration"))]
769fn supports_open_region_object_storage_requirement(object_store: &ObjectStore) -> bool {
770 is_object_storage(object_store)
771}
772
773#[cfg(feature = "test-shared-fs-region-migration")]
774fn supports_open_region_object_storage_requirement(object_store: &ObjectStore) -> bool {
775 is_object_storage(object_store)
780 || object_store.info().scheme() == object_store::services::FS_SCHEME
781}
782
783pub(crate) fn version_builder_from_manifest(
785 manifest: &RegionManifest,
786 metadata: RegionMetadataRef,
787 file_purger: FilePurgerRef,
788 mutable: TimePartitionsRef,
789 region_options: RegionOptions,
790) -> VersionBuilder {
791 VersionBuilder::new(metadata, mutable)
792 .add_files(file_purger, manifest.files.values().cloned())
793 .flushed_entry_id(manifest.flushed_entry_id)
794 .flushed_sequence(manifest.flushed_sequence)
795 .truncated_entry_id(manifest.truncated_entry_id)
796 .compaction_time_window(manifest.compaction_time_window)
797 .options(region_options)
798}
799
800pub(crate) fn sanitize_region_options(manifest: &RegionManifest, options: &mut RegionOptions) {
802 match options.sst_format {
807 Some(format) if format != manifest.sst_format => {
808 common_telemetry::warn!(
809 "Overriding SST format from {:?} (manifest) to {:?} (options) for region {}",
810 manifest.sst_format,
811 format,
812 manifest.metadata.region_id,
813 );
814 }
815 Some(_) => {}
816 None => {
817 options.sst_format = Some(manifest.sst_format);
818 }
819 }
820 if let Some(manifest_append_mode) = manifest.append_mode
821 && options.append_mode != manifest_append_mode
822 {
823 common_telemetry::warn!(
824 "Overriding append_mode from {} to {} for region {}",
825 options.append_mode,
826 manifest_append_mode,
827 manifest.metadata.region_id,
828 );
829 options.append_mode = manifest_append_mode;
830 }
831 if options.append_mode && options.merge_mode.take().is_some() {
832 common_telemetry::warn!(
833 "Ignoring merge_mode because append_mode is enabled for region {}",
834 manifest.metadata.region_id,
835 );
836 }
837}
838
839pub(crate) fn sanitize_open_request_options(options: &mut HashMap<String, String>) {
841 let append_mode_enabled = options
842 .get("append_mode")
843 .is_some_and(|v| matches!(v.trim().to_ascii_lowercase().as_str(), "true" | "1"));
844
845 if append_mode_enabled && options.remove("merge_mode").is_some() {
846 common_telemetry::warn!(
847 "Ignoring merge_mode in open request options because append_mode is enabled"
848 );
849 }
850}
851
852pub fn get_object_store(
854 name: &Option<String>,
855 object_store_manager: &ObjectStoreManagerRef,
856) -> Result<object_store::ObjectStore> {
857 if let Some(name) = name {
858 Ok(object_store_manager
859 .find(name)
860 .with_context(|| ObjectStoreNotFoundSnafu {
861 object_store: name.clone(),
862 })?
863 .clone())
864 } else {
865 Ok(object_store_manager.default_object_store().clone())
866 }
867}
868
869pub(crate) fn check_recovered_region(
871 recovered: &RegionMetadata,
872 region_id: RegionId,
873 column_metadatas: &[ColumnMetadata],
874 primary_key: &[ColumnId],
875) -> Result<()> {
876 if recovered.region_id != region_id {
877 error!(
878 "Recovered region {}, expect region {}",
879 recovered.region_id, region_id
880 );
881 return RegionCorruptedSnafu {
882 region_id,
883 reason: format!(
884 "recovered metadata has different region id {}",
885 recovered.region_id
886 ),
887 }
888 .fail();
889 }
890 if recovered.column_metadatas != column_metadatas {
891 error!(
892 "Unexpected schema in recovered region {}, recovered: {:?}, expect: {:?}",
893 recovered.region_id, recovered.column_metadatas, column_metadatas
894 );
895
896 return RegionCorruptedSnafu {
897 region_id,
898 reason: "recovered metadata has different schema",
899 }
900 .fail();
901 }
902 if recovered.primary_key != primary_key {
903 error!(
904 "Unexpected primary key in recovered region {}, recovered: {:?}, expect: {:?}",
905 recovered.region_id, recovered.primary_key, primary_key
906 );
907
908 return RegionCorruptedSnafu {
909 region_id,
910 reason: "recovered metadata has different primary key",
911 }
912 .fail();
913 }
914
915 Ok(())
916}
917
918pub(crate) async fn replay_memtable<F>(
920 provider: &Provider,
921 mut wal_entry_reader: Box<dyn WalEntryReader>,
922 region_id: RegionId,
923 flushed_entry_id: EntryId,
924 version_control: &VersionControlRef,
925 allow_stale_entries: bool,
926 on_region_opened: F,
927) -> Result<EntryId>
928where
929 F: FnOnce(RegionId, EntryId, &Provider) -> BoxFuture<Result<()>> + Send,
930{
931 let now = Instant::now();
932 let mut rows_replayed = 0;
933 let mut last_entry_id = flushed_entry_id;
936 let replay_from_entry_id = flushed_entry_id + 1;
937 let region_metadata = version_control.current().version.metadata.clone();
938
939 let mut wal_stream = wal_entry_reader.read(provider, replay_from_entry_id)?;
940 while let Some(res) = wal_stream.next().await {
941 let (entry_id, entry) = res?;
942 if entry_id <= flushed_entry_id {
943 warn!(
944 "Stale WAL entries read during replay, region id: {}, flushed entry id: {}, entry id read: {}",
945 region_id, flushed_entry_id, entry_id
946 );
947 ensure!(
948 allow_stale_entries,
949 StaleLogEntrySnafu {
950 region_id,
951 flushed_entry_id,
952 unexpected_entry_id: entry_id,
953 }
954 );
955 }
956 last_entry_id = last_entry_id.max(entry_id);
957
958 let mut region_write_ctx = RegionWriteCtx::new(
959 region_id,
960 version_control,
961 provider.clone(),
962 None,
964 );
965 for mutation in entry.mutations {
966 rows_replayed += mutation
967 .rows
968 .as_ref()
969 .map(|rows| rows.rows.len())
970 .unwrap_or(0);
971 region_write_ctx.push_mutation(
972 mutation.op_type,
973 mutation.rows,
974 mutation.write_hint,
975 OptionOutputTx::none(),
976 Some(mutation.sequence),
978 false,
979 );
980 }
981
982 for bulk_entry in entry.bulk_entries {
983 let mut part = BulkPart::try_from(bulk_entry)?;
984 if region_metadata.primary_key_encoding != PrimaryKeyEncoding::Sparse {
991 part.fill_missing_columns(®ion_metadata)?;
992 }
993 rows_replayed += part.num_rows();
994 let bulk_sequence_from_wal = part.sequence;
996 ensure!(
997 region_write_ctx.push_bulk(
998 OptionOutputTx::none(),
999 part,
1000 Some(bulk_sequence_from_wal),
1001 false
1002 ),
1003 RegionCorruptedSnafu {
1004 region_id,
1005 reason: "unable to replay memtable with bulk entries",
1006 }
1007 );
1008 }
1009
1010 region_write_ctx.set_next_entry_id(last_entry_id + 1);
1012 region_write_ctx.write_memtable().await;
1013 region_write_ctx.write_bulk().await;
1014 region_write_ctx.publish_sequence_and_entry_id();
1017 }
1018
1019 (on_region_opened)(region_id, flushed_entry_id, provider).await?;
1022
1023 let series_count = version_control.current().series_count();
1024 info!(
1025 "Replay WAL for region: {}, provider: {:?}, rows recovered: {}, replay from entry id: {}, last entry id: {}, total timeseries replayed: {}, elapsed: {:?}",
1026 region_id,
1027 provider,
1028 rows_replayed,
1029 replay_from_entry_id,
1030 last_entry_id,
1031 series_count,
1032 now.elapsed()
1033 );
1034 Ok(last_entry_id)
1035}
1036
1037pub(crate) struct RegionLoadCacheTask {
1039 region: MitoRegionRef,
1040}
1041
1042impl RegionLoadCacheTask {
1043 pub(crate) fn new(region: MitoRegionRef) -> Self {
1044 Self { region }
1045 }
1046
1047 pub(crate) async fn fill_cache(&self, file_cache: &FileCache) {
1049 let region_id = self.region.region_id;
1050 let table_dir = self.region.access_layer.table_dir();
1051 let path_type = self.region.access_layer.path_type();
1052 let object_store = self.region.access_layer.object_store();
1053 let version_control = &self.region.version_control;
1054
1055 let mut files_to_download = Vec::new();
1057 let mut files_already_cached = 0;
1058
1059 {
1060 let version = version_control.current().version;
1061 for level in version.ssts.levels() {
1062 for file_handle in level.files.values() {
1063 let file_meta = file_handle.meta_ref();
1064 if file_meta.exists_index() {
1065 let puffin_key = IndexKey::new(
1066 file_meta.region_id,
1067 file_meta.file_id,
1068 FileType::Puffin(file_meta.index_version),
1069 );
1070
1071 if !file_cache.contains_key(&puffin_key) {
1072 files_to_download.push((
1073 puffin_key,
1074 file_meta.index_file_size,
1075 file_meta.time_range.1, ));
1077 } else {
1078 files_already_cached += 1;
1079 }
1080 }
1081 }
1082 }
1083 }
1086
1087 files_to_download.sort_by_key(|b| std::cmp::Reverse(b.2));
1089
1090 let total_files = files_to_download.len() as i64;
1091
1092 info!(
1093 "Starting background index cache preload for region {}, total_files_to_download: {}, files_already_cached: {}",
1094 region_id, total_files, files_already_cached
1095 );
1096
1097 CACHE_FILL_PENDING_FILES.add(total_files);
1098
1099 let mut files_downloaded = 0;
1100 let mut files_skipped = 0;
1101
1102 for (puffin_key, file_size, max_timestamp) in files_to_download {
1103 let current_size = file_cache.puffin_cache_size();
1104 let capacity = file_cache.puffin_cache_capacity();
1105 let region_state = self.region.state();
1106 if !can_load_cache(region_state) {
1107 info!(
1108 "Stopping index cache by state: {:?}, region: {}, current_size: {}, capacity: {}",
1109 region_state, region_id, current_size, capacity
1110 );
1111 break;
1112 }
1113
1114 if current_size + file_size > capacity {
1116 info!(
1117 "Stopping index cache preload due to capacity limit, region: {}, file_id: {}, current_size: {}, file_size: {}, capacity: {}, file_timestamp: {:?}",
1118 region_id, puffin_key.file_id, current_size, file_size, capacity, max_timestamp
1119 );
1120 files_skipped = (total_files - files_downloaded) as usize;
1121 CACHE_FILL_PENDING_FILES.sub(total_files - files_downloaded);
1122 break;
1123 }
1124
1125 let index_version = if let FileType::Puffin(version) = puffin_key.file_type {
1126 version
1127 } else {
1128 unreachable!("`files_to_download` should only contains Puffin files");
1129 };
1130 let index_id = RegionIndexId::new(
1131 RegionFileId::new(puffin_key.region_id, puffin_key.file_id),
1132 index_version,
1133 );
1134
1135 let index_remote_path = location::index_file_path(table_dir, index_id, path_type);
1136
1137 match file_cache
1138 .download(puffin_key, &index_remote_path, object_store, file_size)
1139 .await
1140 {
1141 Ok(_) => {
1142 debug!(
1143 "Downloaded index file to write cache, region: {}, file_id: {}",
1144 region_id, puffin_key.file_id
1145 );
1146 files_downloaded += 1;
1147 CACHE_FILL_DOWNLOADED_FILES.inc_by(1);
1148 CACHE_FILL_PENDING_FILES.dec();
1149 }
1150 Err(e) => {
1151 warn!(
1152 e; "Failed to download index file to write cache, region: {}, file_id: {}",
1153 region_id, puffin_key.file_id
1154 );
1155 CACHE_FILL_PENDING_FILES.dec();
1156 }
1157 }
1158 }
1159
1160 info!(
1161 "Completed background cache fill task for region {}, total_files: {}, files_downloaded: {}, files_already_cached: {}, files_skipped: {}",
1162 region_id, total_files, files_downloaded, files_already_cached, files_skipped
1163 );
1164 }
1165}
1166
1167fn maybe_load_cache(
1169 region: &MitoRegionRef,
1170 config: &MitoConfig,
1171 cache_manager: &Option<CacheManagerRef>,
1172) {
1173 let Some(cache_manager) = cache_manager else {
1174 return;
1175 };
1176 let Some(write_cache) = cache_manager.write_cache() else {
1177 return;
1178 };
1179
1180 let preload_enabled = config.preload_index_cache;
1181 if !preload_enabled {
1182 return;
1183 }
1184
1185 let task = RegionLoadCacheTask::new(region.clone());
1186 write_cache.load_region_cache(task);
1187}
1188
1189#[allow(clippy::too_many_arguments)]
1200async fn preload_parquet_meta_cache_for_files(
1201 region_id: RegionId,
1202 cache_manager: CacheManagerRef,
1203 sst_meta_cache_capacity: u64,
1204 table_dir: String,
1205 path_type: PathType,
1206 object_store: object_store::ObjectStore,
1207 region_metadata: RegionMetadataRef,
1208 mut files: Vec<FileHandle>,
1209) -> usize {
1210 if !cache_manager.sst_meta_cache_enabled()
1211 || sst_meta_cache_capacity == 0
1212 || cache_manager.sst_meta_cache_is_full()
1213 {
1214 return 0;
1215 }
1216
1217 let allow_direct_load = object_store.info().scheme() == object_store::services::FS_SCHEME;
1218
1219 files.sort_by_key(|b| std::cmp::Reverse(b.meta_ref().time_range.1));
1221
1222 let mut loaded = 0usize;
1223 for file_handle in files {
1224 if cache_manager.sst_meta_cache_is_full() {
1226 break;
1227 }
1228
1229 let file_id = file_handle.file_id();
1230 let mut cache_metrics = MetadataCacheMetrics::default();
1231 if let Some(metadata) = cache_manager
1232 .get_compact_sst_meta_data(file_id, PageIndexPolicy::Optional)
1233 .await
1234 {
1235 if file_handle.raw_primary_key_range().is_none()
1236 && let Some(primary_key_range) = extract_primary_key_range(
1237 metadata.parquet_metadata().as_ref(),
1238 ®ion_metadata,
1239 )
1240 {
1241 file_handle.set_primary_key_range(primary_key_range);
1242 }
1243 continue;
1244 }
1245
1246 if let Some(write_cache) = cache_manager.write_cache() {
1247 let key = IndexKey::new(file_id.region_id(), file_id.file_id(), FileType::Parquet);
1248 if let Some(metadata) = write_cache
1249 .file_cache()
1250 .get_sst_meta_data(
1251 key,
1252 &mut cache_metrics,
1253 PageIndexPolicy::Optional,
1254 &common_runtime::global_runtime(),
1255 )
1256 .await
1257 {
1258 let decoded = metadata.decoded();
1259 if file_handle.raw_primary_key_range().is_none()
1260 && let Some(primary_key_range) = extract_primary_key_range(
1261 decoded.parquet_metadata().as_ref(),
1262 ®ion_metadata,
1263 )
1264 {
1265 file_handle.set_primary_key_range(primary_key_range);
1266 }
1267 match metadata {
1268 SstMetaPreparation::Prepared(metadata) => {
1269 cache_manager.put_prepared_sst_meta(file_id, metadata, false);
1270 loaded += 1;
1271 }
1272 SstMetaPreparation::DecodedOnly { encoding_error, .. } => warn!(
1273 encoding_error;
1274 "Failed to encode file-cached SST metadata during preload, region: {}, file: {}",
1275 region_id,
1276 file_id.file_id()
1277 ),
1278 }
1279 continue;
1280 }
1281 }
1282
1283 if !allow_direct_load {
1284 continue;
1285 }
1286
1287 let file_size = file_handle.meta_ref().file_size;
1288 let file_path = file_handle.file_path(&table_dir, path_type);
1289 let mut loader = MetadataLoader::new(object_store.clone(), &file_path, file_size);
1290 loader.with_page_index_policy(PageIndexPolicy::Optional);
1291 match loader.load(&mut cache_metrics).await {
1292 Ok(metadata) => {
1293 if let Some(primary_key_range) =
1294 extract_primary_key_range(&metadata, ®ion_metadata)
1295 {
1296 file_handle.set_primary_key_range(primary_key_range);
1297 }
1298 match prepare_sst_meta(
1299 &file_path,
1300 metadata,
1301 None,
1304 PageIndexPolicy::Optional,
1305 &common_runtime::global_runtime(),
1306 )
1307 .await
1308 {
1309 Ok(SstMetaPreparation::Prepared(metadata)) => {
1310 cache_manager.put_prepared_sst_meta(file_id, metadata, false);
1313 loaded += 1;
1314 }
1315 Ok(SstMetaPreparation::DecodedOnly { encoding_error, .. }) => warn!(
1316 encoding_error;
1317 "Failed to encode preloaded parquet metadata, region: {}, file: {}",
1318 region_id,
1319 file_path
1320 ),
1321 Err(err) => {
1322 warn!(
1323 err; "Failed to prepare preloaded parquet metadata, region: {}, file: {}",
1324 region_id, file_path
1325 );
1326 }
1327 }
1328 }
1329 Err(err) => {
1330 warn!(
1332 err; "Failed to preload parquet metadata from local store, region: {}, file: {}",
1333 region_id, file_path
1334 );
1335 }
1336 }
1337 }
1338
1339 loaded
1340}
1341
1342fn maybe_preload_parquet_meta_cache(
1343 region: &MitoRegionRef,
1344 config: &MitoConfig,
1345 cache_manager: &Option<CacheManagerRef>,
1346) {
1347 let Some(cache_manager) = cache_manager else {
1348 return;
1349 };
1350 if !cache_manager.sst_meta_cache_enabled() {
1351 return;
1352 }
1353
1354 if config.sst_meta_cache_size.as_bytes() == 0 {
1356 return;
1357 }
1358 if !config.preload_index_cache {
1359 return;
1360 }
1361
1362 let region = region.clone();
1363 let cache_manager = cache_manager.clone();
1364 let sst_meta_cache_capacity = config.sst_meta_cache_size.as_bytes();
1365
1366 tokio::spawn(async move {
1367 let _permit = PARQUET_META_PRELOAD_SEMAPHORE.acquire().await.unwrap();
1369
1370 let region_id = region.region_id;
1371 let table_dir = region.access_layer.table_dir().to_string();
1372 let path_type = region.access_layer.path_type();
1373 let object_store = region.access_layer.object_store().clone();
1374 let region_metadata = region.version_control.current().version.metadata.clone();
1375
1376 let mut files = Vec::new();
1378 {
1379 let version = region.version_control.current().version;
1380 for level in version.ssts.levels() {
1381 for file_handle in level.files.values() {
1382 files.push(file_handle.clone());
1383 }
1384 }
1385 }
1386 let preloading_start = Instant::now();
1387 let loaded = preload_parquet_meta_cache_for_files(
1388 region_id,
1389 cache_manager,
1390 sst_meta_cache_capacity,
1391 table_dir,
1392 path_type,
1393 object_store,
1394 region_metadata,
1395 files,
1396 )
1397 .await;
1398 let preloading_cost = preloading_start.elapsed();
1399
1400 if loaded > 0 {
1401 info!(
1402 "Preloaded parquet metadata for region {}, loaded_files: {}, elapsed_ms: {}",
1403 region_id,
1404 loaded,
1405 preloading_cost.as_millis()
1406 );
1407 }
1408 });
1409}
1410
1411fn can_load_cache(state: RegionRoleState) -> bool {
1412 match state {
1413 RegionRoleState::Leader(RegionLeaderState::Writable)
1414 | RegionRoleState::Leader(RegionLeaderState::Staging)
1415 | RegionRoleState::Leader(RegionLeaderState::Altering)
1416 | RegionRoleState::Leader(RegionLeaderState::EnteringStaging)
1417 | RegionRoleState::Leader(RegionLeaderState::Editing)
1418 | RegionRoleState::Follower => true,
1419 RegionRoleState::Leader(RegionLeaderState::Downgrading)
1421 | RegionRoleState::Leader(RegionLeaderState::Dropping)
1422 | RegionRoleState::Leader(RegionLeaderState::Truncating) => false,
1423 }
1424}
1425
1426#[cfg(test)]
1427mod tests {
1428 use std::collections::HashMap;
1429 use std::sync::Arc;
1430
1431 use arrow_schema::extension::ExtensionType;
1432 use common_base::readable_size::ReadableSize;
1433 use common_error::ext::WhateverResult;
1434 use common_test_util::temp_dir::create_temp_dir;
1435 use common_time::Timestamp;
1436 use common_wal::options::{KafkaWalOptions, ObjectStoreWalOptions, WalOptions};
1437 use datatypes::arrow::array::{ArrayRef, BinaryArray, Int64Array};
1438 use datatypes::arrow::record_batch::RecordBatch;
1439 use datatypes::extension::json::{Json2ExtensionType, JsonMetadata};
1440 use datatypes::json::JsonSettings;
1441 use datatypes::prelude::ConcreteDataType;
1442 use datatypes::schema::ColumnSchema;
1443 use datatypes::types::json_type::{JsonNativeType, JsonObjectType};
1444 use log_store::kafka::log_store::KafkaLogStore;
1445 use log_store::noop::log_store::NoopLogStore;
1446 use log_store::raft_engine::log_store::RaftEngineLogStore;
1447 use object_store::ObjectStore;
1448 use object_store::services::{Fs, Memory, S3};
1449 use parquet::arrow::ArrowWriter;
1450 use parquet::file::metadata::{KeyValue, PageIndexPolicy};
1451 use parquet::file::properties::WriterProperties;
1452 use store_api::logstore::provider::Provider;
1453 use store_api::metadata::RegionMetadataBuilder;
1454 use store_api::region_request::PathType;
1455 use store_api::storage::{FileId, RegionId};
1456
1457 use super::{
1458 initial_pruned_entry_id, maybe_upgrade_json2_layout, preload_parquet_meta_cache_for_files,
1459 provider_from_wal_options, sanitize_region_options,
1460 supports_open_region_object_storage_requirement,
1461 };
1462 use crate::cache::CacheManager;
1463 use crate::cache::file_cache::{FileType, IndexKey};
1464 use crate::error;
1465 use crate::manifest::action::{RegionManifest, RemovedFilesRecord};
1466 use crate::region::options::RegionOptions;
1467 use crate::sst::FormatType;
1468 use crate::sst::file::{FileHandle, FileMeta};
1469 use crate::sst::file_purger::NoopFilePurger;
1470 use crate::sst::parquet::PARQUET_METADATA_KEY;
1471 use crate::test_util::TestEnv;
1472 use crate::test_util::sst_util::sst_region_metadata;
1473
1474 fn build_test_manifest(sst_format: FormatType) -> RegionManifest {
1475 RegionManifest {
1476 metadata: Arc::new(sst_region_metadata()),
1477 files: HashMap::new(),
1478 removed_files: RemovedFilesRecord::default(),
1479 flushed_entry_id: 0,
1480 flushed_sequence: 0,
1481 committed_sequence: None,
1482 manifest_version: 0,
1483 truncated_entry_id: None,
1484 compaction_time_window: None,
1485 sst_format,
1486 append_mode: None,
1487 }
1488 }
1489
1490 fn build_fs_object_store() -> ObjectStore {
1491 ObjectStore::new(Fs::default().root("/tmp")).unwrap()
1492 }
1493
1494 #[test]
1495 fn test_provider_from_object_store_wal_options() {
1496 let region_id = RegionId::new(1, 2);
1497 let wal_options = WalOptions::ObjectStore(ObjectStoreWalOptions::new("wal".to_string()));
1498
1499 let provider = provider_from_wal_options::<NoopLogStore>(region_id, &wal_options).unwrap();
1500 assert_eq!(
1501 Provider::object_store_provider(region_id, "wal".to_string()),
1502 provider
1503 );
1504
1505 let err =
1506 provider_from_wal_options::<RaftEngineLogStore>(region_id, &wal_options).unwrap_err();
1507 assert!(matches!(
1508 err,
1509 error::Error::IncompatibleWalProviderChange { .. }
1510 ));
1511 let err = provider_from_wal_options::<KafkaLogStore>(region_id, &wal_options).unwrap_err();
1512 assert!(matches!(
1513 err,
1514 error::Error::IncompatibleWalProviderChange { .. }
1515 ));
1516 }
1517
1518 #[test]
1519 fn test_initial_pruned_entry_id() {
1520 assert_eq!(0, initial_pruned_entry_id(&WalOptions::RaftEngine));
1521 assert_eq!(0, initial_pruned_entry_id(&WalOptions::Noop));
1522 assert_eq!(
1523 0,
1524 initial_pruned_entry_id(&WalOptions::ObjectStore(ObjectStoreWalOptions::new(
1525 "wal".to_string()
1526 )))
1527 );
1528 assert_eq!(
1529 0,
1530 initial_pruned_entry_id(&WalOptions::Kafka(KafkaWalOptions::new(
1531 "test_topic".to_string()
1532 )))
1533 );
1534 assert_eq!(
1535 42,
1536 initial_pruned_entry_id(&WalOptions::Kafka(KafkaWalOptions {
1537 topic: "test_topic".to_string(),
1538 initial_pruned_entry_id: Some(42),
1539 }))
1540 );
1541 }
1542
1543 #[test]
1544 fn test_upgrade_json2_layout() -> WhateverResult<()> {
1545 let settings = JsonSettings::try_new(vec![], Some(3))?;
1546 let extension = Json2ExtensionType::new(Arc::new(JsonMetadata::new_v1(settings.clone())));
1547 let mut column = ColumnSchema::new(
1548 "field_0",
1549 ConcreteDataType::json2(JsonNativeType::Object(JsonObjectType::new())),
1550 true,
1551 );
1552 column.with_extension_type(&extension);
1553
1554 let mut metadata = sst_region_metadata();
1555 metadata.column_metadatas[2].column_schema = column;
1556 let builder = RegionMetadataBuilder::from_existing(metadata);
1557 let metadata = Arc::new(builder.build()?);
1558
1559 let upgraded = maybe_upgrade_json2_layout(metadata)?;
1560 let column = &upgraded.column_metadatas[2].column_schema;
1561 let extension = column.extension_type::<Json2ExtensionType>()?.unwrap();
1562 assert!(extension.metadata().is_version_2());
1563 assert_eq!(&settings, extension.metadata().json_settings());
1564
1565 let arrow_schema = upgraded.schema.arrow_schema();
1566 let field = arrow_schema.field_with_name("field_0").unwrap();
1567 let extension = field.try_extension_type::<Json2ExtensionType>().unwrap();
1568 assert!(extension.metadata().is_version_2());
1569
1570 let unchanged = maybe_upgrade_json2_layout(upgraded.clone())?;
1571 assert!(Arc::ptr_eq(&upgraded, &unchanged));
1572 Ok(())
1573 }
1574
1575 #[test]
1576 #[cfg(not(feature = "test-shared-fs-region-migration"))]
1577 fn test_open_requirement_rejects_fs_object_store() {
1578 let object_store = build_fs_object_store();
1579
1580 assert!(!supports_open_region_object_storage_requirement(
1581 &object_store
1582 ));
1583 }
1584
1585 #[test]
1586 #[cfg(feature = "test-shared-fs-region-migration")]
1587 fn test_open_requirement_accepts_shared_fs_object_store_for_tests() {
1588 let object_store = build_fs_object_store();
1589
1590 assert!(supports_open_region_object_storage_requirement(
1591 &object_store
1592 ));
1593 }
1594
1595 #[test]
1596 fn test_open_requirement_accepts_s3_object_store() {
1597 let object_store = ObjectStore::new(
1598 S3::default()
1599 .bucket("test-bucket")
1600 .region("us-east-1")
1601 .disable_ec2_metadata(),
1602 )
1603 .unwrap();
1604
1605 assert!(supports_open_region_object_storage_requirement(
1606 &object_store
1607 ));
1608 }
1609
1610 #[test]
1611 fn test_sanitize_region_options_options_format_wins() {
1612 let manifest = build_test_manifest(FormatType::PrimaryKey);
1615 let mut options = RegionOptions {
1616 sst_format: Some(FormatType::Flat),
1617 ..Default::default()
1618 };
1619 sanitize_region_options(&manifest, &mut options);
1620 assert_eq!(options.sst_format, Some(FormatType::Flat));
1621 }
1622
1623 #[test]
1624 fn test_sanitize_region_options_fills_from_manifest_when_unset() {
1625 let manifest = build_test_manifest(FormatType::Flat);
1628 let mut options = RegionOptions {
1629 sst_format: None,
1630 ..Default::default()
1631 };
1632 sanitize_region_options(&manifest, &mut options);
1633 assert_eq!(options.sst_format, Some(FormatType::Flat));
1634 }
1635
1636 fn sst_parquet_bytes(batch: &RecordBatch) -> Vec<u8> {
1637 let key_value_meta = KeyValue::new(
1638 PARQUET_METADATA_KEY.to_string(),
1639 sst_region_metadata().to_json().unwrap(),
1640 );
1641 let props = WriterProperties::builder()
1642 .set_key_value_metadata(Some(vec![key_value_meta]))
1643 .build();
1644
1645 let mut parquet_bytes = Vec::new();
1646 let mut writer =
1647 ArrowWriter::try_new(&mut parquet_bytes, batch.schema(), Some(props)).unwrap();
1648 writer.write(batch).unwrap();
1649 writer.close().unwrap();
1650
1651 parquet_bytes
1652 }
1653
1654 #[tokio::test]
1655 async fn test_preload_parquet_meta_cache_uses_file_cache() {
1656 let env = TestEnv::new().await;
1657
1658 let local_store = ObjectStore::new(Memory::default()).unwrap();
1659 let write_cache = env
1660 .create_write_cache(local_store, ReadableSize::mb(1024))
1661 .await;
1662 let cache_manager = Arc::new(
1663 CacheManager::builder()
1664 .sst_meta_cache_size(1024 * 1024)
1665 .write_cache(Some(write_cache.clone()))
1666 .build(),
1667 );
1668
1669 let region_id = RegionId::new(1, 1);
1670 let file_id = FileId::random();
1671
1672 let col = Arc::new(Int64Array::from_iter_values([1, 2, 3])) as ArrayRef;
1673 let keys =
1674 ["a", "b", "c"].map(|tag| crate::test_util::sst_util::new_primary_key(&[tag, ""]));
1675 let primary_key = Arc::new(BinaryArray::from_iter_values(&keys)) as ArrayRef;
1676 let batch = RecordBatch::try_from_iter([
1677 ("col", col),
1678 (
1679 store_api::storage::consts::PRIMARY_KEY_COLUMN_NAME,
1680 primary_key,
1681 ),
1682 ])
1683 .unwrap();
1684 let parquet_bytes = sst_parquet_bytes(&batch);
1685 let file_size = parquet_bytes.len() as u64;
1686
1687 let file_meta = FileMeta {
1688 region_id,
1689 file_id,
1690 time_range: (Timestamp::new_millisecond(0), Timestamp::new_millisecond(1)),
1691 level: 0,
1692 file_size,
1693 max_row_group_uncompressed_size: 0,
1694 available_indexes: Default::default(),
1695 indexes: vec![],
1696 index_file_size: 0,
1697 index_version: 0,
1698 num_rows: 3,
1699 num_row_groups: 1,
1700 sequence: None,
1701 partition_expr: None,
1702 num_series: 0,
1703 ..Default::default()
1704 };
1705 let file_handle = FileHandle::new(file_meta, Arc::new(NoopFilePurger));
1706
1707 let table_dir = "test_table";
1708 let path_type = PathType::Bare;
1709 let remote_path = file_handle.file_path(table_dir, path_type);
1710
1711 let source_store = ObjectStore::new(Memory::default()).unwrap();
1712 source_store
1713 .write(&remote_path, parquet_bytes)
1714 .await
1715 .unwrap();
1716
1717 let index_key = IndexKey::new(region_id, file_id, FileType::Parquet);
1719 write_cache
1720 .file_cache()
1721 .download(index_key, &remote_path, &source_store, file_size)
1722 .await
1723 .unwrap();
1724
1725 let region_file_id = file_handle.file_id();
1726 assert!(
1727 cache_manager
1728 .get_parquet_meta_data_from_mem_cache(region_file_id)
1729 .is_none()
1730 );
1731
1732 let loaded = preload_parquet_meta_cache_for_files(
1733 region_id,
1734 cache_manager.clone(),
1735 1024 * 1024,
1736 table_dir.to_string(),
1737 path_type,
1738 source_store.clone(),
1739 Arc::new(sst_region_metadata()),
1740 vec![file_handle.clone()],
1741 )
1742 .await;
1743
1744 assert_eq!(loaded, 1);
1746 assert!(
1747 cache_manager
1748 .get_compact_sst_meta_data(
1749 region_file_id,
1750 parquet::file::metadata::PageIndexPolicy::Optional,
1751 )
1752 .await
1753 .is_some()
1754 );
1755 assert!(file_handle.raw_primary_key_range().is_some());
1756 }
1757
1758 #[tokio::test]
1759 async fn test_preload_parquet_meta_cache_skips_files_not_in_file_cache() {
1760 let cache_manager = Arc::new(
1761 CacheManager::builder()
1762 .sst_meta_cache_size(1024 * 1024)
1763 .build(),
1764 );
1765
1766 let region_id = RegionId::new(1, 1);
1767 let file_id = FileId::random();
1768
1769 let file_meta = FileMeta {
1771 region_id,
1772 file_id,
1773 time_range: (Timestamp::new_millisecond(0), Timestamp::new_millisecond(1)),
1774 level: 0,
1775 file_size: 0,
1776 max_row_group_uncompressed_size: 0,
1777 available_indexes: Default::default(),
1778 indexes: vec![],
1779 index_file_size: 0,
1780 index_version: 0,
1781 num_rows: 3,
1782 num_row_groups: 1,
1783 sequence: None,
1784 partition_expr: None,
1785 num_series: 0,
1786 ..Default::default()
1787 };
1788 let file_handle = FileHandle::new(file_meta, Arc::new(NoopFilePurger));
1789
1790 let table_dir = "test_table";
1791 let path_type = PathType::Bare;
1792 let remote_path = file_handle.file_path(table_dir, path_type);
1793
1794 let object_store = ObjectStore::new(Memory::default()).unwrap();
1796 object_store
1797 .write(&remote_path, b"noop".as_slice())
1798 .await
1799 .unwrap();
1800
1801 let region_file_id = file_handle.file_id();
1802 assert!(
1803 cache_manager
1804 .get_parquet_meta_data_from_mem_cache(region_file_id)
1805 .is_none()
1806 );
1807
1808 let loaded = preload_parquet_meta_cache_for_files(
1809 region_id,
1810 cache_manager.clone(),
1811 1024 * 1024,
1812 table_dir.to_string(),
1813 path_type,
1814 object_store,
1815 Arc::new(sst_region_metadata()),
1816 vec![file_handle],
1817 )
1818 .await;
1819
1820 assert_eq!(loaded, 0);
1821 assert!(
1822 cache_manager
1823 .get_parquet_meta_data_from_mem_cache(region_file_id)
1824 .is_none()
1825 );
1826 }
1827
1828 #[tokio::test]
1829 async fn test_preload_parquet_meta_cache_loads_from_local_fs() {
1830 let cache_manager = Arc::new(
1831 CacheManager::builder()
1832 .sst_meta_cache_size(1024 * 1024)
1833 .build(),
1834 );
1835
1836 let region_id = RegionId::new(1, 1);
1837 let file_id = FileId::random();
1838
1839 let col = Arc::new(Int64Array::from_iter_values([1, 2, 3])) as ArrayRef;
1840 let batch = RecordBatch::try_from_iter([("col", col)]).unwrap();
1841 let parquet_bytes = sst_parquet_bytes(&batch);
1842
1843 let file_meta = FileMeta {
1846 region_id,
1847 file_id,
1848 time_range: (Timestamp::new_millisecond(0), Timestamp::new_millisecond(1)),
1849 level: 0,
1850 file_size: 0,
1851 max_row_group_uncompressed_size: 0,
1852 available_indexes: Default::default(),
1853 indexes: vec![],
1854 index_file_size: 0,
1855 index_version: 0,
1856 num_rows: 3,
1857 num_row_groups: 1,
1858 sequence: None,
1859 partition_expr: None,
1860 num_series: 0,
1861 ..Default::default()
1862 };
1863 let file_handle = FileHandle::new(file_meta, Arc::new(NoopFilePurger));
1864
1865 let table_dir = "test_table";
1866 let path_type = PathType::Bare;
1867 let file_path = file_handle.file_path(table_dir, path_type);
1868
1869 let root = create_temp_dir("parquet-meta-preload");
1870 let object_store =
1871 ObjectStore::new(Fs::default().root(root.path().to_str().unwrap())).unwrap();
1872 object_store.write(&file_path, parquet_bytes).await.unwrap();
1873
1874 let region_file_id = file_handle.file_id();
1875 assert!(
1876 cache_manager
1877 .get_parquet_meta_data_from_mem_cache(region_file_id)
1878 .is_none()
1879 );
1880
1881 let loaded = preload_parquet_meta_cache_for_files(
1882 region_id,
1883 cache_manager.clone(),
1884 1024 * 1024,
1885 table_dir.to_string(),
1886 path_type,
1887 object_store,
1888 Arc::new(sst_region_metadata()),
1889 vec![file_handle],
1890 )
1891 .await;
1892
1893 assert_eq!(loaded, 1);
1894 assert!(
1895 cache_manager
1896 .get_compact_sst_meta_data(region_file_id, PageIndexPolicy::Optional)
1897 .await
1898 .is_some()
1899 );
1900 }
1901}