1#[cfg(test)]
18mod alter_test;
19#[cfg(test)]
20mod append_mode_test;
21#[cfg(test)]
22mod basic_test;
23#[cfg(test)]
24mod batch_catchup_test;
25#[cfg(test)]
26mod batch_open_test;
27#[cfg(test)]
28mod bump_committed_sequence_test;
29#[cfg(test)]
30mod catchup_test;
31#[cfg(test)]
32mod close_test;
33#[cfg(test)]
34pub(crate) mod compaction_test;
35#[cfg(test)]
36mod create_test;
37#[cfg(test)]
38mod drop_test;
39#[cfg(test)]
40mod edit_region_test;
41#[cfg(test)]
42mod file_ref_test;
43#[cfg(test)]
44mod filter_deleted_test;
45#[cfg(test)]
46mod flush_test;
47#[cfg(test)]
48mod index_build_test;
49#[cfg(any(test, feature = "test"))]
50pub mod listener;
51#[cfg(test)]
52mod merge_mode_test;
53#[cfg(test)]
54mod open_test;
55#[cfg(test)]
56mod parallel_test;
57#[cfg(test)]
58mod projection_test;
59#[cfg(test)]
60mod prune_test;
61pub mod region_hook;
62#[cfg(test)]
63mod row_selector_test;
64#[cfg(test)]
65mod scan_corrupt;
66#[cfg(test)]
67mod scan_test;
68#[cfg(test)]
69mod set_role_state_test;
70#[cfg(test)]
71mod skip_wal_test;
72#[cfg(test)]
73mod staging_test;
74#[cfg(test)]
75mod sync_test;
76#[cfg(test)]
77mod truncate_test;
78
79#[cfg(test)]
80mod copy_region_from_test;
81#[cfg(test)]
82mod remap_manifests_test;
83
84#[cfg(test)]
85mod apply_staging_manifest_test;
86#[cfg(test)]
87mod partition_filter_test;
88mod puffin_index;
89
90use std::any::Any;
91use std::collections::{HashMap, HashSet};
92use std::sync::Arc;
93use std::time::Instant;
94
95use api::region::RegionResponse;
96use async_trait::async_trait;
97use common_base::Plugins;
98use common_error::ext::BoxedError;
99use common_meta::error::UnexpectedSnafu;
100use common_meta::key::SchemaMetadataManagerRef;
101use common_recordbatch::{QueryMemoryTracker, SendableRecordBatchStream};
102use common_stat::get_total_memory_bytes;
103use common_telemetry::{info, tracing, warn};
104use common_wal::options::WalOptions;
105use futures::future::{join_all, try_join_all};
106use futures::stream::{self, Stream, StreamExt};
107use object_store::manager::ObjectStoreManagerRef;
108use region_hook::RegionHookRef;
109use snafu::{OptionExt, ResultExt, ensure};
110use store_api::ManifestVersion;
111use store_api::codec::PrimaryKeyEncoding;
112use store_api::logstore::LogStore;
113use store_api::logstore::provider::{KafkaProvider, Provider};
114use store_api::metadata::{ColumnMetadata, RegionMetadataRef};
115use store_api::metric_engine_consts::{
116 MANIFEST_INFO_EXTENSION_KEY, TABLE_COLUMN_METADATA_EXTENSION_KEY,
117};
118use store_api::region_engine::{
119 BatchResponses, MitoCopyRegionFromRequest, MitoCopyRegionFromResponse, RegionEngine,
120 RegionManifestInfo, RegionRole, RegionScannerRef, RegionStatistic, RemapManifestsRequest,
121 RemapManifestsResponse, SetRegionRoleStateResponse, SettableRegionRoleState,
122 SyncRegionFromRequest, SyncRegionFromResponse,
123};
124use store_api::region_info::RegionInfoEntry;
125use store_api::region_request::{
126 AffectedRows, RegionCatchupRequest, RegionOpenRequest, RegionRequest,
127};
128use store_api::sst_entry::{ManifestSstEntry, PuffinIndexMetaEntry, StorageSstEntry};
129use store_api::storage::{FileId, FileRefsManifest, RegionId, ScanRequest, SequenceNumber};
130use tokio::sync::{Semaphore, oneshot};
131
132use crate::access_layer::RegionFilePathFactory;
133use crate::cache::{CacheManagerRef, CacheStrategy};
134use crate::config::MitoConfig;
135use crate::engine::puffin_index::{IndexEntryContext, collect_index_entries_from_puffin};
136use crate::error::{
137 IncrementalQueryStaleSnafu, InvalidRequestSnafu, JoinSnafu, MitoManifestInfoSnafu, RecvSnafu,
138 RegionNotFoundSnafu, Result, SerdeJsonSnafu, SerializeColumnMetadataSnafu,
139 SnapshotFenceStaleSnafu,
140};
141#[cfg(feature = "enterprise")]
142use crate::extension::BoxedExtensionRangeProviderFactory;
143use crate::gc::GcLimiterRef;
144use crate::manifest::action::RegionEdit;
145use crate::memtable::MemtableStats;
146use crate::metrics::{
147 HANDLE_REQUEST_ELAPSED, SCAN_MEMORY_EXHAUSTED_TOTAL, SCAN_MEMORY_USAGE_BYTES,
148 SCAN_REQUESTS_REJECTED_TOTAL,
149};
150use crate::read::scan_region::{ScanRegion, Scanner};
151use crate::read::stream::ScanBatchStream;
152use crate::region::MitoRegionRef;
153use crate::region::opener::PartitionExprFetcherRef;
154use crate::region::options::parse_wal_options;
155use crate::request::{RegionEditRequest, WorkerRequest};
156use crate::sst::file::{FileMeta, RegionFileId, RegionIndexId};
157use crate::sst::file_ref::FileReferenceManagerRef;
158use crate::sst::index::intermediate::IntermediateManager;
159use crate::sst::index::puffin_manager::PuffinManagerFactory;
160use crate::wal::entry_distributor::{
161 DEFAULT_ENTRY_RECEIVER_BUFFER_SIZE, build_wal_entry_distributor_and_receivers,
162};
163use crate::wal::raw_entry_reader::{LogStoreRawEntryReader, RawEntryReader};
164use crate::worker::WorkerGroup;
165
166pub const MITO_ENGINE_NAME: &str = "mito";
167
168pub struct MitoEngineBuilder<'a, S: LogStore> {
169 data_home: &'a str,
170 config: MitoConfig,
171 log_store: Arc<S>,
172 object_store_manager: ObjectStoreManagerRef,
173 schema_metadata_manager: SchemaMetadataManagerRef,
174 file_ref_manager: FileReferenceManagerRef,
175 partition_expr_fetcher: PartitionExprFetcherRef,
176 plugins: Plugins,
177 #[cfg(feature = "enterprise")]
178 extension_range_provider_factory: Option<BoxedExtensionRangeProviderFactory>,
179}
180
181impl<'a, S: LogStore> MitoEngineBuilder<'a, S> {
182 #[allow(clippy::too_many_arguments)]
183 pub fn new(
184 data_home: &'a str,
185 config: MitoConfig,
186 log_store: Arc<S>,
187 object_store_manager: ObjectStoreManagerRef,
188 schema_metadata_manager: SchemaMetadataManagerRef,
189 file_ref_manager: FileReferenceManagerRef,
190 partition_expr_fetcher: PartitionExprFetcherRef,
191 plugins: Plugins,
192 ) -> Self {
193 Self {
194 data_home,
195 config,
196 log_store,
197 object_store_manager,
198 schema_metadata_manager,
199 file_ref_manager,
200 plugins,
201 partition_expr_fetcher,
202 #[cfg(feature = "enterprise")]
203 extension_range_provider_factory: None,
204 }
205 }
206
207 #[cfg(feature = "enterprise")]
208 #[must_use]
209 pub fn with_extension_range_provider_factory(
210 self,
211 extension_range_provider_factory: Option<BoxedExtensionRangeProviderFactory>,
212 ) -> Self {
213 Self {
214 extension_range_provider_factory,
215 ..self
216 }
217 }
218
219 pub async fn try_build(mut self) -> Result<MitoEngine> {
220 self.config.sanitize(self.data_home)?;
221
222 let config = Arc::new(self.config);
223 let region_hook = self.plugins.get::<RegionHookRef>();
226 let workers = WorkerGroup::start(
227 config.clone(),
228 self.log_store.clone(),
229 self.object_store_manager,
230 self.schema_metadata_manager,
231 self.file_ref_manager,
232 self.partition_expr_fetcher.clone(),
233 self.plugins,
234 )
235 .await?;
236 let wal_raw_entry_reader = Arc::new(LogStoreRawEntryReader::new(self.log_store));
237 let total_memory = get_total_memory_bytes().max(0) as u64;
238 let scan_memory_limit = config.scan_memory_limit.resolve(total_memory) as usize;
239 let scan_memory_tracker =
240 QueryMemoryTracker::builder(scan_memory_limit, config.scan_memory_on_exhausted)
241 .on_update(|usage| {
242 SCAN_MEMORY_USAGE_BYTES.set(usage as i64);
243 })
244 .on_exhausted(|| {
245 SCAN_MEMORY_EXHAUSTED_TOTAL.inc();
246 })
247 .on_reject(|| {
248 SCAN_REQUESTS_REJECTED_TOTAL.inc();
249 })
250 .build();
251
252 let inner = EngineInner {
253 workers,
254 config,
255 wal_raw_entry_reader,
256 scan_memory_tracker,
257 region_hook,
258 #[cfg(feature = "enterprise")]
259 extension_range_provider_factory: None,
260 };
261
262 #[cfg(feature = "enterprise")]
263 let inner =
264 inner.with_extension_range_provider_factory(self.extension_range_provider_factory);
265
266 Ok(MitoEngine {
267 inner: Arc::new(inner),
268 })
269 }
270}
271
272#[derive(Clone)]
274pub struct MitoEngine {
275 inner: Arc<EngineInner>,
276}
277
278impl MitoEngine {
279 #[allow(clippy::too_many_arguments)]
281 pub async fn new<S: LogStore>(
282 data_home: &str,
283 config: MitoConfig,
284 log_store: Arc<S>,
285 object_store_manager: ObjectStoreManagerRef,
286 schema_metadata_manager: SchemaMetadataManagerRef,
287 file_ref_manager: FileReferenceManagerRef,
288 partition_expr_fetcher: PartitionExprFetcherRef,
289 plugins: Plugins,
290 ) -> Result<MitoEngine> {
291 let builder = MitoEngineBuilder::new(
292 data_home,
293 config,
294 log_store,
295 object_store_manager,
296 schema_metadata_manager,
297 file_ref_manager,
298 partition_expr_fetcher,
299 plugins,
300 );
301 builder.try_build().await
302 }
303
304 pub fn mito_config(&self) -> &MitoConfig {
305 &self.inner.config
306 }
307
308 pub fn cache_manager(&self) -> CacheManagerRef {
309 self.inner.workers.cache_manager()
310 }
311
312 pub fn file_ref_manager(&self) -> FileReferenceManagerRef {
313 self.inner.workers.file_ref_manager()
314 }
315
316 pub fn gc_limiter(&self) -> GcLimiterRef {
317 self.inner.workers.gc_limiter()
318 }
319
320 pub fn object_store_manager(&self) -> &ObjectStoreManagerRef {
321 self.inner.workers.object_store_manager()
322 }
323
324 pub fn puffin_manager_factory(&self) -> &PuffinManagerFactory {
325 self.inner.workers.puffin_manager_factory()
326 }
327
328 pub fn intermediate_manager(&self) -> &IntermediateManager {
329 self.inner.workers.intermediate_manager()
330 }
331
332 pub fn schema_metadata_manager(&self) -> &SchemaMetadataManagerRef {
333 self.inner.workers.schema_metadata_manager()
334 }
335
336 pub fn region_hook(&self) -> Option<RegionHookRef> {
339 self.inner.region_hook.clone()
340 }
341
342 pub async fn get_snapshot_of_file_refs(
344 &self,
345 file_handle_regions: impl IntoIterator<Item = RegionId>,
346 related_regions: HashMap<RegionId, HashSet<RegionId>>,
347 ) -> Result<FileRefsManifest> {
348 let file_ref_mgr = self.file_ref_manager();
349
350 let file_handle_regions = file_handle_regions.into_iter().collect::<Vec<_>>();
351 let query_regions: Vec<MitoRegionRef> = file_handle_regions
354 .into_iter()
355 .filter_map(|region_id| self.find_region(region_id))
356 .collect();
357
358 let dst_region_to_src_regions: Vec<(MitoRegionRef, HashSet<RegionId>)> = {
359 let dst2src = related_regions
360 .into_iter()
361 .flat_map(|(src, dsts)| dsts.into_iter().map(move |dst| (dst, src)))
362 .fold(
363 HashMap::<RegionId, HashSet<RegionId>>::new(),
364 |mut acc, (k, v)| {
365 let entry = acc.entry(k).or_default();
366 entry.insert(v);
367 acc
368 },
369 );
370 let mut dst_region_to_src_regions = Vec::with_capacity(dst2src.len());
371 for (dst_region, srcs) in dst2src {
372 let Some(region) = self.find_region(dst_region) else {
373 return RegionNotFoundSnafu {
374 region_id: dst_region,
375 }
376 .fail();
377 };
378 dst_region_to_src_regions.push((region, srcs));
379 }
380 dst_region_to_src_regions
381 };
382
383 file_ref_mgr
384 .get_snapshot_of_file_refs(query_regions, dst_region_to_src_regions)
385 .await
386 }
387
388 pub fn is_region_exists(&self, region_id: RegionId) -> bool {
390 self.inner.workers.is_region_exists(region_id)
391 }
392
393 pub fn is_region_opening(&self, region_id: RegionId) -> bool {
395 self.inner.workers.is_region_opening(region_id)
396 }
397
398 pub fn is_region_catching_up(&self, region_id: RegionId) -> bool {
400 self.inner.workers.is_region_catching_up(region_id)
401 }
402
403 pub fn get_region_statistic(&self, region_id: RegionId) -> Option<RegionStatistic> {
405 self.find_region(region_id)
406 .map(|region| region.region_statistic())
407 }
408
409 pub fn get_primary_key_encoding(&self, region_id: RegionId) -> Option<PrimaryKeyEncoding> {
411 self.find_region(region_id)
412 .map(|r| r.primary_key_encoding())
413 }
414
415 #[tracing::instrument(skip_all)]
420 pub async fn scan_to_stream(
421 &self,
422 region_id: RegionId,
423 request: ScanRequest,
424 ) -> Result<SendableRecordBatchStream, BoxedError> {
425 self.scanner(region_id, request)
426 .await
427 .map_err(BoxedError::new)?
428 .scan()
429 .await
430 }
431
432 pub async fn scan_batch(
434 &self,
435 region_id: RegionId,
436 request: ScanRequest,
437 filter_deleted: bool,
438 ) -> Result<ScanBatchStream> {
439 let mut scan_region = self.scan_region(region_id, request)?;
440 scan_region.set_filter_deleted(filter_deleted);
441 scan_region.scanner().await?.scan_batch()
442 }
443
444 pub(crate) async fn scanner(
446 &self,
447 region_id: RegionId,
448 request: ScanRequest,
449 ) -> Result<Scanner> {
450 self.scan_region(region_id, request)?.scanner().await
451 }
452
453 #[tracing::instrument(skip_all, fields(region_id = %region_id))]
455 fn scan_region(&self, region_id: RegionId, request: ScanRequest) -> Result<ScanRegion> {
456 self.inner.scan_region(region_id, request)
457 }
458
459 pub async fn edit_region(&self, region_id: RegionId, edit: RegionEdit) -> Result<()> {
464 let _timer = HANDLE_REQUEST_ELAPSED
465 .with_label_values(&["edit_region"])
466 .start_timer();
467
468 ensure!(
469 is_valid_region_edit(&edit),
470 InvalidRequestSnafu {
471 region_id,
472 reason: "invalid region edit"
473 }
474 );
475
476 let (tx, rx) = oneshot::channel();
477 let request = WorkerRequest::EditRegion(RegionEditRequest::new(region_id, edit, true, tx));
478 self.inner
479 .workers
480 .submit_to_worker(region_id, request)
481 .await?;
482 rx.await.context(RecvSnafu)?
483 }
484
485 pub async fn copy_region_from(
489 &self,
490 region_id: RegionId,
491 request: MitoCopyRegionFromRequest,
492 ) -> Result<MitoCopyRegionFromResponse> {
493 self.inner.copy_region_from(region_id, request).await
494 }
495
496 #[cfg(test)]
497 pub(crate) fn get_region(&self, id: RegionId) -> Option<crate::region::MitoRegionRef> {
498 self.find_region(id)
499 }
500
501 pub fn find_region(&self, region_id: RegionId) -> Option<MitoRegionRef> {
502 self.inner.workers.get_region(region_id)
503 }
504
505 pub fn regions(&self) -> Vec<MitoRegionRef> {
507 self.inner.workers.all_regions().collect()
508 }
509
510 fn encode_manifest_info_to_extensions(
511 region_id: &RegionId,
512 manifest_info: RegionManifestInfo,
513 extensions: &mut HashMap<String, Vec<u8>>,
514 ) -> Result<()> {
515 let region_manifest_info = vec![(*region_id, manifest_info)];
516
517 extensions.insert(
518 MANIFEST_INFO_EXTENSION_KEY.to_string(),
519 RegionManifestInfo::encode_list(®ion_manifest_info).context(SerdeJsonSnafu)?,
520 );
521 info!(
522 "Added manifest info: {:?} to extensions, region_id: {:?}",
523 region_manifest_info, region_id
524 );
525 Ok(())
526 }
527
528 fn encode_column_metadatas_to_extensions(
529 region_id: &RegionId,
530 column_metadatas: Vec<ColumnMetadata>,
531 extensions: &mut HashMap<String, Vec<u8>>,
532 ) -> Result<()> {
533 extensions.insert(
534 TABLE_COLUMN_METADATA_EXTENSION_KEY.to_string(),
535 ColumnMetadata::encode_list(&column_metadatas).context(SerializeColumnMetadataSnafu)?,
536 );
537 info!(
538 "Added column metadatas: {:?} to extensions, region_id: {:?}",
539 column_metadatas, region_id
540 );
541 Ok(())
542 }
543
544 pub fn find_memtable_and_sst_stats(
547 &self,
548 region_id: RegionId,
549 ) -> Result<(Vec<MemtableStats>, Vec<FileMeta>)> {
550 let region = self
551 .find_region(region_id)
552 .context(RegionNotFoundSnafu { region_id })?;
553
554 let version = region.version();
555 let memtable_stats = version
556 .memtables
557 .list_memtables()
558 .iter()
559 .map(|x| x.stats())
560 .collect::<Vec<_>>();
561
562 let sst_stats = version
563 .ssts
564 .levels()
565 .iter()
566 .flat_map(|level| level.files().map(|x| x.meta_ref()))
567 .cloned()
568 .collect::<Vec<_>>();
569 Ok((memtable_stats, sst_stats))
570 }
571
572 pub async fn all_ssts_from_manifest(&self) -> Vec<ManifestSstEntry> {
574 let node_id = self.inner.workers.file_ref_manager().node_id();
575 let regions = self.inner.workers.all_regions();
576
577 let mut results = Vec::new();
578 for region in regions {
579 let mut entries = region.manifest_sst_entries().await;
580 for e in &mut entries {
581 e.node_id = node_id;
582 }
583 results.extend(entries);
584 }
585
586 results
587 }
588
589 pub async fn all_index_metas(&self) -> Vec<PuffinIndexMetaEntry> {
591 let node_id = self.inner.workers.file_ref_manager().node_id();
592 let cache_manager = self.inner.workers.cache_manager();
593 let puffin_metadata_cache = cache_manager.puffin_metadata_cache().cloned();
594 let bloom_filter_cache = cache_manager.bloom_filter_index_cache().cloned();
595 let inverted_index_cache = cache_manager.inverted_index_cache().cloned();
596
597 let mut results = Vec::new();
598
599 for region in self.inner.workers.all_regions() {
600 let manifest_entries = region.manifest_sst_entries().await;
601 let access_layer = region.access_layer.clone();
602 let table_dir = access_layer.table_dir().to_string();
603 let path_type = access_layer.path_type();
604 let object_store = access_layer.object_store().clone();
605 let puffin_factory = access_layer.puffin_manager_factory().clone();
606 let path_factory = RegionFilePathFactory::new(table_dir, path_type);
607
608 let entry_futures = manifest_entries.into_iter().map(|entry| {
609 let object_store = object_store.clone();
610 let path_factory = path_factory.clone();
611 let puffin_factory = puffin_factory.clone();
612 let puffin_metadata_cache = puffin_metadata_cache.clone();
613 let bloom_filter_cache = bloom_filter_cache.clone();
614 let inverted_index_cache = inverted_index_cache.clone();
615
616 async move {
617 let Some(index_file_path) = entry.index_file_path.as_ref() else {
618 return Vec::new();
619 };
620
621 let index_version = entry.index_version;
622 let file_id = match FileId::parse_str(&entry.file_id) {
623 Ok(file_id) => file_id,
624 Err(err) => {
625 warn!(
626 err;
627 "Failed to parse puffin index file id, table_dir: {}, file_id: {}",
628 entry.table_dir,
629 entry.file_id
630 );
631 return Vec::new();
632 }
633 };
634 let region_index_id = RegionIndexId::new(
637 RegionFileId::new(entry.origin_region_id, file_id),
638 index_version,
639 );
640 let context = IndexEntryContext {
641 table_dir: &entry.table_dir,
642 index_file_path: index_file_path.as_str(),
643 region_id: entry.region_id,
644 table_id: entry.table_id,
645 region_number: entry.region_number,
646 region_group: entry.region_group,
647 region_sequence: entry.region_sequence,
648 file_id: &entry.file_id,
649 index_file_size: entry.index_file_size,
650 node_id,
651 };
652
653 let manager = puffin_factory
654 .build(object_store, path_factory)
655 .with_puffin_metadata_cache(puffin_metadata_cache);
656
657 collect_index_entries_from_puffin(
658 manager,
659 region_index_id,
660 context,
661 bloom_filter_cache,
662 inverted_index_cache,
663 )
664 .await
665 }
666 });
667
668 let mut meta_stream = stream::iter(entry_futures).buffer_unordered(8); while let Some(mut metas) = meta_stream.next().await {
670 results.append(&mut metas);
671 }
672 }
673
674 results
675 }
676
677 pub async fn all_region_infos(&self) -> Vec<RegionInfoEntry> {
679 let node_id = self.inner.workers.file_ref_manager().node_id();
680 self.inner
681 .workers
682 .all_regions()
683 .map(|region| region.region_info_entry(node_id))
684 .collect()
685 }
686
687 pub fn all_ssts_from_storage(&self) -> impl Stream<Item = Result<StorageSstEntry>> {
689 let node_id = self.inner.workers.file_ref_manager().node_id();
690 let regions = self.inner.workers.all_regions();
691
692 let mut layers_distinct_table_dirs = HashMap::new();
693 for region in regions {
694 let table_dir = region.access_layer.table_dir();
695 if !layers_distinct_table_dirs.contains_key(table_dir) {
696 layers_distinct_table_dirs
697 .insert(table_dir.to_string(), region.access_layer.clone());
698 }
699 }
700
701 stream::iter(layers_distinct_table_dirs)
702 .map(|(_, access_layer)| access_layer.storage_sst_entries())
703 .flatten()
704 .map(move |entry| {
705 entry.map(move |mut entry| {
706 entry.node_id = node_id;
707 entry
708 })
709 })
710 }
711}
712
713fn is_valid_region_edit(edit: &RegionEdit) -> bool {
717 (!edit.files_to_add.is_empty() || !edit.files_to_remove.is_empty())
718 && matches!(
719 edit,
720 RegionEdit {
721 files_to_add: _,
722 files_to_remove: _,
723 timestamp_ms: _,
724 compaction_time_window: None,
725 flushed_entry_id: None,
726 flushed_sequence: None,
727 ..
728 }
729 )
730}
731
732struct EngineInner {
734 workers: WorkerGroup,
736 config: Arc<MitoConfig>,
738 wal_raw_entry_reader: Arc<dyn RawEntryReader>,
740 scan_memory_tracker: QueryMemoryTracker,
742 region_hook: Option<RegionHookRef>,
745 #[cfg(feature = "enterprise")]
746 extension_range_provider_factory: Option<BoxedExtensionRangeProviderFactory>,
747}
748
749type TopicGroupedRegionOpenRequests = HashMap<String, Vec<(RegionId, RegionOpenRequest)>>;
750
751fn prepare_batch_open_requests(
753 requests: Vec<(RegionId, RegionOpenRequest)>,
754) -> Result<(
755 TopicGroupedRegionOpenRequests,
756 Vec<(RegionId, RegionOpenRequest)>,
757)> {
758 let mut topic_to_regions: HashMap<String, Vec<(RegionId, RegionOpenRequest)>> = HashMap::new();
759 let mut remaining_regions: Vec<(RegionId, RegionOpenRequest)> = Vec::new();
760 for (region_id, request) in requests {
761 match parse_wal_options(&request.options).context(SerdeJsonSnafu)? {
762 WalOptions::Kafka(options) => {
763 topic_to_regions
764 .entry(options.topic)
765 .or_default()
766 .push((region_id, request));
767 }
768 WalOptions::RaftEngine | WalOptions::Noop => {
769 remaining_regions.push((region_id, request));
770 }
771 }
772 }
773
774 Ok((topic_to_regions, remaining_regions))
775}
776
777impl EngineInner {
778 #[cfg(feature = "enterprise")]
779 #[must_use]
780 fn with_extension_range_provider_factory(
781 self,
782 extension_range_provider_factory: Option<BoxedExtensionRangeProviderFactory>,
783 ) -> Self {
784 Self {
785 extension_range_provider_factory,
786 ..self
787 }
788 }
789
790 async fn stop(&self) -> Result<()> {
792 self.workers.stop().await
793 }
794
795 fn find_region(&self, region_id: RegionId) -> Result<MitoRegionRef> {
796 self.workers
797 .get_region(region_id)
798 .context(RegionNotFoundSnafu { region_id })
799 }
800
801 fn get_metadata(&self, region_id: RegionId) -> Result<RegionMetadataRef> {
805 let region = self.find_region(region_id)?;
807 Ok(region.metadata())
808 }
809
810 async fn open_topic_regions(
811 &self,
812 topic: String,
813 region_requests: Vec<(RegionId, RegionOpenRequest)>,
814 ) -> Result<Vec<(RegionId, Result<AffectedRows>)>> {
815 let now = Instant::now();
816 let region_ids = region_requests
817 .iter()
818 .map(|(region_id, _)| *region_id)
819 .collect::<Vec<_>>();
820 let provider = Provider::kafka_provider(topic);
821 let (distributor, entry_receivers) = build_wal_entry_distributor_and_receivers(
822 provider.clone(),
823 self.wal_raw_entry_reader.clone(),
824 ®ion_ids,
825 DEFAULT_ENTRY_RECEIVER_BUFFER_SIZE,
826 );
827
828 let mut responses = Vec::with_capacity(region_requests.len());
829 for ((region_id, request), entry_receiver) in
830 region_requests.into_iter().zip(entry_receivers)
831 {
832 let (request, receiver) =
833 WorkerRequest::new_open_region_request(region_id, request, Some(entry_receiver));
834 self.workers.submit_to_worker(region_id, request).await?;
835 responses.push(async move { receiver.await.context(RecvSnafu)? });
836 }
837
838 let distribution =
840 common_runtime::spawn_global(async move { distributor.distribute().await });
841 let responses = join_all(responses).await;
843 distribution.await.context(JoinSnafu)??;
844
845 let num_failure = responses.iter().filter(|r| r.is_err()).count();
846 info!(
847 "Opened {} regions for topic '{}', failures: {}, elapsed: {:?}",
848 region_ids.len() - num_failure,
849 provider.as_kafka_provider().unwrap(),
851 num_failure,
852 now.elapsed(),
853 );
854 Ok(region_ids.into_iter().zip(responses).collect())
855 }
856
857 async fn handle_batch_open_requests(
858 &self,
859 parallelism: usize,
860 requests: Vec<(RegionId, RegionOpenRequest)>,
861 ) -> Result<Vec<(RegionId, Result<AffectedRows>)>> {
862 let semaphore = Arc::new(Semaphore::new(parallelism));
863 let (topic_to_region_requests, remaining_region_requests) =
864 prepare_batch_open_requests(requests)?;
865 let mut responses =
866 Vec::with_capacity(topic_to_region_requests.len() + remaining_region_requests.len());
867
868 if !topic_to_region_requests.is_empty() {
869 let mut tasks = Vec::with_capacity(topic_to_region_requests.len());
870 for (topic, region_requests) in topic_to_region_requests {
871 let semaphore_moved = semaphore.clone();
872 tasks.push(async move {
873 let _permit = semaphore_moved.acquire().await.unwrap();
875 self.open_topic_regions(topic, region_requests).await
876 })
877 }
878 let r = try_join_all(tasks).await?;
879 responses.extend(r.into_iter().flatten());
880 }
881
882 if !remaining_region_requests.is_empty() {
883 let mut tasks = Vec::with_capacity(remaining_region_requests.len());
884 let mut region_ids = Vec::with_capacity(remaining_region_requests.len());
885 for (region_id, request) in remaining_region_requests {
886 let semaphore_moved = semaphore.clone();
887 region_ids.push(region_id);
888 tasks.push(async move {
889 let _permit = semaphore_moved.acquire().await.unwrap();
891 let (request, receiver) =
892 WorkerRequest::new_open_region_request(region_id, request, None);
893
894 self.workers.submit_to_worker(region_id, request).await?;
895
896 receiver.await.context(RecvSnafu)?
897 })
898 }
899
900 let results = join_all(tasks).await;
901 responses.extend(region_ids.into_iter().zip(results));
902 }
903
904 Ok(responses)
905 }
906
907 async fn catchup_topic_regions(
908 &self,
909 provider: Provider,
910 region_requests: Vec<(RegionId, RegionCatchupRequest)>,
911 ) -> Result<Vec<(RegionId, Result<AffectedRows>)>> {
912 let now = Instant::now();
913 let region_ids = region_requests
914 .iter()
915 .map(|(region_id, _)| *region_id)
916 .collect::<Vec<_>>();
917 let (distributor, entry_receivers) = build_wal_entry_distributor_and_receivers(
918 provider.clone(),
919 self.wal_raw_entry_reader.clone(),
920 ®ion_ids,
921 DEFAULT_ENTRY_RECEIVER_BUFFER_SIZE,
922 );
923
924 let mut responses = Vec::with_capacity(region_requests.len());
925 for ((region_id, request), entry_receiver) in
926 region_requests.into_iter().zip(entry_receivers)
927 {
928 let (request, receiver) =
929 WorkerRequest::new_catchup_region_request(region_id, request, Some(entry_receiver));
930 self.workers.submit_to_worker(region_id, request).await?;
931 responses.push(async move { receiver.await.context(RecvSnafu)? });
932 }
933
934 let distribution =
936 common_runtime::spawn_global(async move { distributor.distribute().await });
937 let responses = join_all(responses).await;
939 distribution.await.context(JoinSnafu)??;
940
941 let num_failure = responses.iter().filter(|r| r.is_err()).count();
942 info!(
943 "Caught up {} regions for topic '{}', failures: {}, elapsed: {:?}",
944 region_ids.len() - num_failure,
945 provider.as_kafka_provider().unwrap(),
947 num_failure,
948 now.elapsed(),
949 );
950
951 Ok(region_ids.into_iter().zip(responses).collect())
952 }
953
954 async fn handle_batch_catchup_requests(
955 &self,
956 parallelism: usize,
957 requests: Vec<(RegionId, RegionCatchupRequest)>,
958 ) -> Result<Vec<(RegionId, Result<AffectedRows>)>> {
959 let mut responses = Vec::with_capacity(requests.len());
960 let mut topic_regions: HashMap<Arc<KafkaProvider>, Vec<_>> = HashMap::new();
961 let mut remaining_region_requests = vec![];
962
963 for (region_id, request) in requests {
964 match self.workers.get_region(region_id) {
965 Some(region) => match region.provider.as_kafka_provider() {
966 Some(provider) => {
967 topic_regions
968 .entry(provider.clone())
969 .or_default()
970 .push((region_id, request));
971 }
972 None => {
973 remaining_region_requests.push((region_id, request));
974 }
975 },
976 None => responses.push((region_id, RegionNotFoundSnafu { region_id }.fail())),
977 }
978 }
979
980 let semaphore = Arc::new(Semaphore::new(parallelism));
981
982 if !topic_regions.is_empty() {
983 let mut tasks = Vec::with_capacity(topic_regions.len());
984 for (provider, region_requests) in topic_regions {
985 let semaphore_moved = semaphore.clone();
986 tasks.push(async move {
987 let _permit = semaphore_moved.acquire().await.unwrap();
989 self.catchup_topic_regions(Provider::Kafka(provider), region_requests)
990 .await
991 })
992 }
993
994 let r = try_join_all(tasks).await?;
995 responses.extend(r.into_iter().flatten());
996 }
997
998 if !remaining_region_requests.is_empty() {
999 let mut tasks = Vec::with_capacity(remaining_region_requests.len());
1000 let mut region_ids = Vec::with_capacity(remaining_region_requests.len());
1001 for (region_id, request) in remaining_region_requests {
1002 let semaphore_moved = semaphore.clone();
1003 region_ids.push(region_id);
1004 tasks.push(async move {
1005 let _permit = semaphore_moved.acquire().await.unwrap();
1007 let (request, receiver) =
1008 WorkerRequest::new_catchup_region_request(region_id, request, None);
1009
1010 self.workers.submit_to_worker(region_id, request).await?;
1011
1012 receiver.await.context(RecvSnafu)?
1013 })
1014 }
1015
1016 let results = join_all(tasks).await;
1017 responses.extend(region_ids.into_iter().zip(results));
1018 }
1019
1020 Ok(responses)
1021 }
1022
1023 async fn handle_request(
1025 &self,
1026 region_id: RegionId,
1027 request: RegionRequest,
1028 ) -> Result<AffectedRows> {
1029 let region_metadata = self.get_metadata(region_id).ok();
1030 let (request, receiver) =
1031 WorkerRequest::try_from_region_request(region_id, request, region_metadata)?;
1032 self.workers.submit_to_worker(region_id, request).await?;
1033
1034 receiver.await.context(RecvSnafu)?
1035 }
1036
1037 fn get_committed_sequence(&self, region_id: RegionId) -> Result<SequenceNumber> {
1039 self.find_region(region_id)
1041 .map(|r| r.find_committed_sequence())
1042 }
1043
1044 #[tracing::instrument(skip_all, fields(region_id = %region_id))]
1046 fn scan_region(&self, region_id: RegionId, mut request: ScanRequest) -> Result<ScanRegion> {
1047 let query_start = Instant::now();
1048 let region = self.find_region(region_id)?;
1050 let version_data = region.version_control.current();
1051 let version = version_data.version;
1052
1053 if request.snapshot_on_scan && request.memtable_max_sequence.is_none() {
1054 request.memtable_max_sequence = Some(version_data.committed_sequence);
1055 }
1056
1057 if let Some(given_seq) = request.memtable_min_sequence {
1058 let min_readable_seq = version.flushed_sequence;
1059 ensure!(
1060 given_seq >= min_readable_seq,
1061 IncrementalQueryStaleSnafu {
1062 region_id,
1063 given_seq,
1064 min_readable_seq,
1065 }
1066 );
1067 }
1068
1069 if let Some(given_seq) = request.memtable_max_sequence
1070 && !request.skip_sst_files
1071 {
1072 let min_enforceable_seq = version.flushed_sequence;
1079 ensure!(
1080 given_seq >= min_enforceable_seq,
1081 SnapshotFenceStaleSnafu {
1082 region_id,
1083 given_seq,
1084 min_enforceable_seq,
1085 }
1086 );
1087 }
1088
1089 let cache_manager = self.workers.cache_manager();
1091
1092 let scan_region = ScanRegion::new(
1093 version,
1094 region.access_layer.clone(),
1095 request,
1096 CacheStrategy::EnableAll(cache_manager),
1097 )
1098 .with_query_stat_counters(region.region_stats.query_stat_counters())
1099 .with_max_concurrent_scan_files(self.config.max_concurrent_scan_files)
1100 .with_ignore_inverted_index(self.config.inverted_index.apply_on_query.disabled())
1101 .with_ignore_fulltext_index(self.config.fulltext_index.apply_on_query.disabled())
1102 .with_ignore_bloom_filter(self.config.bloom_filter_index.apply_on_query.disabled())
1103 .with_start_time(query_start);
1104
1105 #[cfg(feature = "enterprise")]
1106 let scan_region = self.maybe_fill_extension_range_provider(scan_region, region);
1107
1108 Ok(scan_region)
1109 }
1110
1111 #[cfg(feature = "enterprise")]
1112 fn maybe_fill_extension_range_provider(
1113 &self,
1114 mut scan_region: ScanRegion,
1115 region: MitoRegionRef,
1116 ) -> ScanRegion {
1117 if region.is_follower()
1118 && let Some(factory) = self.extension_range_provider_factory.as_ref()
1119 {
1120 scan_region
1121 .set_extension_range_provider(factory.create_extension_range_provider(region));
1122 }
1123 scan_region
1124 }
1125
1126 fn set_region_role(&self, region_id: RegionId, role: RegionRole) -> Result<()> {
1128 let region = self.find_region(region_id)?;
1129 region.set_role(role);
1130 Ok(())
1131 }
1132
1133 async fn set_region_role_state_gracefully(
1135 &self,
1136 region_id: RegionId,
1137 region_role_state: SettableRegionRoleState,
1138 ) -> Result<SetRegionRoleStateResponse> {
1139 let (request, receiver) =
1142 WorkerRequest::new_set_readonly_gracefully(region_id, region_role_state);
1143 self.workers.submit_to_worker(region_id, request).await?;
1144
1145 receiver.await.context(RecvSnafu)
1146 }
1147
1148 async fn sync_region(
1149 &self,
1150 region_id: RegionId,
1151 manifest_info: RegionManifestInfo,
1152 ) -> Result<(ManifestVersion, bool)> {
1153 ensure!(manifest_info.is_mito(), MitoManifestInfoSnafu);
1154 let manifest_version = manifest_info.data_manifest_version();
1155 let (request, receiver) =
1156 WorkerRequest::new_sync_region_request(region_id, manifest_version);
1157 self.workers.submit_to_worker(region_id, request).await?;
1158
1159 receiver.await.context(RecvSnafu)?
1160 }
1161
1162 async fn remap_manifests(
1163 &self,
1164 request: RemapManifestsRequest,
1165 ) -> Result<RemapManifestsResponse> {
1166 let region_id = request.region_id;
1167 let (request, receiver) = WorkerRequest::try_from_remap_manifests_request(request)?;
1168 self.workers.submit_to_worker(region_id, request).await?;
1169 let manifest_paths = receiver.await.context(RecvSnafu)??;
1170 Ok(RemapManifestsResponse { manifest_paths })
1171 }
1172
1173 async fn copy_region_from(
1174 &self,
1175 region_id: RegionId,
1176 request: MitoCopyRegionFromRequest,
1177 ) -> Result<MitoCopyRegionFromResponse> {
1178 let (request, receiver) =
1179 WorkerRequest::try_from_copy_region_from_request(region_id, request)?;
1180 self.workers.submit_to_worker(region_id, request).await?;
1181 let response = receiver.await.context(RecvSnafu)??;
1182 Ok(response)
1183 }
1184
1185 fn role(&self, region_id: RegionId) -> Option<RegionRole> {
1186 self.workers
1187 .get_region(region_id)
1188 .map(|region| region.region_role())
1189 }
1190}
1191
1192fn map_batch_responses(responses: Vec<(RegionId, Result<AffectedRows>)>) -> BatchResponses {
1193 responses
1194 .into_iter()
1195 .map(|(region_id, response)| {
1196 (
1197 region_id,
1198 response.map(RegionResponse::new).map_err(BoxedError::new),
1199 )
1200 })
1201 .collect()
1202}
1203
1204#[async_trait]
1205impl RegionEngine for MitoEngine {
1206 fn name(&self) -> &str {
1207 MITO_ENGINE_NAME
1208 }
1209
1210 #[tracing::instrument(skip_all)]
1211 async fn handle_batch_open_requests(
1212 &self,
1213 parallelism: usize,
1214 requests: Vec<(RegionId, RegionOpenRequest)>,
1215 ) -> Result<BatchResponses, BoxedError> {
1216 self.inner
1218 .handle_batch_open_requests(parallelism, requests)
1219 .await
1220 .map(map_batch_responses)
1221 .map_err(BoxedError::new)
1222 }
1223
1224 #[tracing::instrument(skip_all)]
1225 async fn handle_batch_catchup_requests(
1226 &self,
1227 parallelism: usize,
1228 requests: Vec<(RegionId, RegionCatchupRequest)>,
1229 ) -> Result<BatchResponses, BoxedError> {
1230 self.inner
1231 .handle_batch_catchup_requests(parallelism, requests)
1232 .await
1233 .map(map_batch_responses)
1234 .map_err(BoxedError::new)
1235 }
1236
1237 #[tracing::instrument(skip_all)]
1238 async fn handle_request(
1239 &self,
1240 region_id: RegionId,
1241 request: RegionRequest,
1242 ) -> Result<RegionResponse, BoxedError> {
1243 let _timer = HANDLE_REQUEST_ELAPSED
1244 .with_label_values(&[request.request_type()])
1245 .start_timer();
1246
1247 let is_alter = matches!(request, RegionRequest::Alter(_));
1248 let is_create = matches!(request, RegionRequest::Create(_));
1249 let mut response = self
1250 .inner
1251 .handle_request(region_id, request)
1252 .await
1253 .map(RegionResponse::new)
1254 .map_err(BoxedError::new)?;
1255
1256 if is_alter {
1257 self.handle_alter_response(region_id, &mut response)
1258 .map_err(BoxedError::new)?;
1259 } else if is_create {
1260 self.handle_create_response(region_id, &mut response)
1261 .map_err(BoxedError::new)?;
1262 }
1263
1264 Ok(response)
1265 }
1266
1267 #[tracing::instrument(skip_all)]
1268 async fn handle_query(
1269 &self,
1270 region_id: RegionId,
1271 request: ScanRequest,
1272 ) -> Result<RegionScannerRef, BoxedError> {
1273 self.scan_region(region_id, request)
1274 .map_err(BoxedError::new)?
1275 .region_scanner()
1276 .await
1277 .map_err(BoxedError::new)
1278 }
1279
1280 fn query_memory_tracker(&self) -> Option<QueryMemoryTracker> {
1281 Some(self.inner.scan_memory_tracker.clone())
1282 }
1283
1284 async fn get_committed_sequence(
1285 &self,
1286 region_id: RegionId,
1287 ) -> Result<SequenceNumber, BoxedError> {
1288 self.inner
1289 .get_committed_sequence(region_id)
1290 .map_err(BoxedError::new)
1291 }
1292
1293 async fn get_metadata(
1295 &self,
1296 region_id: RegionId,
1297 ) -> std::result::Result<RegionMetadataRef, BoxedError> {
1298 self.inner.get_metadata(region_id).map_err(BoxedError::new)
1299 }
1300
1301 async fn stop(&self) -> std::result::Result<(), BoxedError> {
1307 self.inner.stop().await.map_err(BoxedError::new)
1308 }
1309
1310 fn region_statistic(&self, region_id: RegionId) -> Option<RegionStatistic> {
1311 self.get_region_statistic(region_id)
1312 }
1313
1314 fn set_region_role(&self, region_id: RegionId, role: RegionRole) -> Result<(), BoxedError> {
1315 self.inner
1316 .set_region_role(region_id, role)
1317 .map_err(BoxedError::new)
1318 }
1319
1320 async fn set_region_role_state_gracefully(
1321 &self,
1322 region_id: RegionId,
1323 region_role_state: SettableRegionRoleState,
1324 ) -> Result<SetRegionRoleStateResponse, BoxedError> {
1325 let _timer = HANDLE_REQUEST_ELAPSED
1326 .with_label_values(&["set_region_role_state_gracefully"])
1327 .start_timer();
1328
1329 self.inner
1330 .set_region_role_state_gracefully(region_id, region_role_state)
1331 .await
1332 .map_err(BoxedError::new)
1333 }
1334
1335 async fn sync_region(
1336 &self,
1337 region_id: RegionId,
1338 request: SyncRegionFromRequest,
1339 ) -> Result<SyncRegionFromResponse, BoxedError> {
1340 let manifest_info = request
1341 .into_region_manifest_info()
1342 .context(UnexpectedSnafu {
1343 err_msg: "Expected a manifest info request",
1344 })
1345 .map_err(BoxedError::new)?;
1346 let (_, synced) = self
1347 .inner
1348 .sync_region(region_id, manifest_info)
1349 .await
1350 .map_err(BoxedError::new)?;
1351
1352 Ok(SyncRegionFromResponse::Mito { synced })
1353 }
1354
1355 async fn remap_manifests(
1356 &self,
1357 request: RemapManifestsRequest,
1358 ) -> Result<RemapManifestsResponse, BoxedError> {
1359 self.inner
1360 .remap_manifests(request)
1361 .await
1362 .map_err(BoxedError::new)
1363 }
1364
1365 fn role(&self, region_id: RegionId) -> Option<RegionRole> {
1366 self.inner.role(region_id)
1367 }
1368
1369 fn as_any(&self) -> &dyn Any {
1370 self
1371 }
1372}
1373
1374impl MitoEngine {
1375 fn handle_alter_response(
1376 &self,
1377 region_id: RegionId,
1378 response: &mut RegionResponse,
1379 ) -> Result<()> {
1380 if let Some(statistic) = self.region_statistic(region_id) {
1381 Self::encode_manifest_info_to_extensions(
1382 ®ion_id,
1383 statistic.manifest,
1384 &mut response.extensions,
1385 )?;
1386 }
1387 let column_metadatas = self
1388 .inner
1389 .find_region(region_id)
1390 .ok()
1391 .map(|r| r.metadata().column_metadatas.clone());
1392 if let Some(column_metadatas) = column_metadatas {
1393 Self::encode_column_metadatas_to_extensions(
1394 ®ion_id,
1395 column_metadatas,
1396 &mut response.extensions,
1397 )?;
1398 }
1399 Ok(())
1400 }
1401
1402 fn handle_create_response(
1403 &self,
1404 region_id: RegionId,
1405 response: &mut RegionResponse,
1406 ) -> Result<()> {
1407 let column_metadatas = self
1408 .inner
1409 .find_region(region_id)
1410 .ok()
1411 .map(|r| r.metadata().column_metadatas.clone());
1412 if let Some(column_metadatas) = column_metadatas {
1413 Self::encode_column_metadatas_to_extensions(
1414 ®ion_id,
1415 column_metadatas,
1416 &mut response.extensions,
1417 )?;
1418 }
1419 Ok(())
1420 }
1421}
1422
1423#[cfg(any(test, feature = "test"))]
1425#[allow(clippy::too_many_arguments)]
1426impl MitoEngine {
1427 pub async fn new_for_test<S: LogStore>(
1429 data_home: &str,
1430 mut config: MitoConfig,
1431 log_store: Arc<S>,
1432 object_store_manager: ObjectStoreManagerRef,
1433 write_buffer_manager: Option<crate::flush::WriteBufferManagerRef>,
1434 listener: Option<crate::engine::listener::EventListenerRef>,
1435 time_provider: crate::time_provider::TimeProviderRef,
1436 schema_metadata_manager: SchemaMetadataManagerRef,
1437 file_ref_manager: FileReferenceManagerRef,
1438 partition_expr_fetcher: PartitionExprFetcherRef,
1439 ) -> Result<MitoEngine> {
1440 config.sanitize(data_home)?;
1441
1442 let config = Arc::new(config);
1443 let wal_raw_entry_reader = Arc::new(LogStoreRawEntryReader::new(log_store.clone()));
1444 let total_memory = get_total_memory_bytes().max(0) as u64;
1445 let scan_memory_limit = config.scan_memory_limit.resolve(total_memory) as usize;
1446 let scan_memory_tracker =
1447 QueryMemoryTracker::builder(scan_memory_limit, config.scan_memory_on_exhausted)
1448 .on_update(|usage| {
1449 SCAN_MEMORY_USAGE_BYTES.set(usage as i64);
1450 })
1451 .on_exhausted(|| {
1452 SCAN_MEMORY_EXHAUSTED_TOTAL.inc();
1453 })
1454 .on_reject(|| {
1455 SCAN_REQUESTS_REJECTED_TOTAL.inc();
1456 })
1457 .build();
1458 Ok(MitoEngine {
1459 inner: Arc::new(EngineInner {
1460 workers: WorkerGroup::start_for_test(
1461 config.clone(),
1462 log_store,
1463 object_store_manager,
1464 write_buffer_manager,
1465 listener,
1466 schema_metadata_manager,
1467 file_ref_manager,
1468 time_provider,
1469 partition_expr_fetcher,
1470 )
1471 .await?,
1472 config,
1473 wal_raw_entry_reader,
1474 scan_memory_tracker,
1475 region_hook: None,
1476 #[cfg(feature = "enterprise")]
1477 extension_range_provider_factory: None,
1478 }),
1479 })
1480 }
1481
1482 pub fn purge_scheduler(&self) -> &crate::schedule::scheduler::SchedulerRef {
1484 self.inner.workers.purge_scheduler()
1485 }
1486}
1487
1488#[cfg(test)]
1489mod tests {
1490 use std::time::Duration;
1491
1492 use super::*;
1493 use crate::sst::file::FileMeta;
1494
1495 #[test]
1496 fn test_is_valid_region_edit() {
1497 let edit = RegionEdit {
1499 files_to_add: vec![FileMeta::default()],
1500 files_to_remove: vec![],
1501 timestamp_ms: None,
1502 compaction_time_window: None,
1503 flushed_entry_id: None,
1504 flushed_sequence: None,
1505 committed_sequence: None,
1506 };
1507 assert!(is_valid_region_edit(&edit));
1508
1509 let edit = RegionEdit {
1511 files_to_add: vec![],
1512 files_to_remove: vec![],
1513 timestamp_ms: None,
1514 compaction_time_window: None,
1515 flushed_entry_id: None,
1516 flushed_sequence: None,
1517 committed_sequence: None,
1518 };
1519 assert!(!is_valid_region_edit(&edit));
1520
1521 let edit = RegionEdit {
1523 files_to_add: vec![],
1524 files_to_remove: vec![FileMeta::default()],
1525 timestamp_ms: None,
1526 compaction_time_window: None,
1527 flushed_entry_id: None,
1528 flushed_sequence: None,
1529 committed_sequence: None,
1530 };
1531 assert!(is_valid_region_edit(&edit));
1532
1533 let edit = RegionEdit {
1535 files_to_add: vec![FileMeta::default()],
1536 files_to_remove: vec![FileMeta::default()],
1537 timestamp_ms: None,
1538 compaction_time_window: None,
1539 flushed_entry_id: None,
1540 flushed_sequence: None,
1541 committed_sequence: None,
1542 };
1543 assert!(is_valid_region_edit(&edit));
1544
1545 let edit = RegionEdit {
1547 files_to_add: vec![FileMeta::default()],
1548 files_to_remove: vec![],
1549 timestamp_ms: None,
1550 compaction_time_window: Some(Duration::from_secs(1)),
1551 flushed_entry_id: None,
1552 flushed_sequence: None,
1553 committed_sequence: None,
1554 };
1555 assert!(!is_valid_region_edit(&edit));
1556 let edit = RegionEdit {
1557 files_to_add: vec![FileMeta::default()],
1558 files_to_remove: vec![],
1559 timestamp_ms: None,
1560 compaction_time_window: None,
1561 flushed_entry_id: Some(1),
1562 flushed_sequence: None,
1563 committed_sequence: None,
1564 };
1565 assert!(!is_valid_region_edit(&edit));
1566 let edit = RegionEdit {
1567 files_to_add: vec![FileMeta::default()],
1568 files_to_remove: vec![],
1569 timestamp_ms: None,
1570 compaction_time_window: None,
1571 flushed_entry_id: None,
1572 flushed_sequence: Some(1),
1573 committed_sequence: None,
1574 };
1575 assert!(!is_valid_region_edit(&edit));
1576 }
1577}