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::{debug, info, tracing, warn};
104use common_wal::options::WalOptions;
105use datafusion::execution::memory_pool::{GreedyMemoryPool, MemoryPool, UnboundedMemoryPool};
106use futures::future::{join_all, try_join_all};
107use futures::stream::{self, Stream, StreamExt};
108use object_store::manager::ObjectStoreManagerRef;
109use region_hook::RegionHookRef;
110use snafu::{OptionExt, ResultExt, ensure};
111use store_api::ManifestVersion;
112use store_api::codec::PrimaryKeyEncoding;
113use store_api::logstore::LogStore;
114use store_api::logstore::provider::{KafkaProvider, Provider};
115use store_api::metadata::{ColumnMetadata, RegionMetadataRef};
116use store_api::metric_engine_consts::{
117 MANIFEST_INFO_EXTENSION_KEY, TABLE_COLUMN_METADATA_EXTENSION_KEY,
118};
119use store_api::region_engine::{
120 BatchResponses, MitoCopyRegionFromRequest, MitoCopyRegionFromResponse, RegionEngine,
121 RegionManifestInfo, RegionRole, RegionScannerRef, RegionStatistic, RemapManifestsRequest,
122 RemapManifestsResponse, SetRegionRoleStateResponse, SettableRegionRoleState,
123 SyncRegionFromRequest, SyncRegionFromResponse,
124};
125use store_api::region_info::RegionInfoEntry;
126use store_api::region_request::{
127 AffectedRows, RegionCatchupRequest, RegionOpenRequest, RegionRequest,
128};
129use store_api::sst_entry::{ManifestSstEntry, PuffinIndexMetaEntry, StorageSstEntry};
130use store_api::storage::{FileId, FileRefsManifest, RegionId, ScanRequest, SequenceNumber};
131use tokio::sync::{Semaphore, oneshot};
132
133use crate::access_layer::RegionFilePathFactory;
134use crate::cache::{CacheManagerRef, CacheStrategy};
135use crate::config::MitoConfig;
136use crate::engine::puffin_index::{IndexEntryContext, collect_index_entries_from_puffin};
137use crate::error::{
138 IncrementalQueryStaleSnafu, InvalidRequestSnafu, JoinSnafu, MitoManifestInfoSnafu, RecvSnafu,
139 RegionNotFoundSnafu, Result, SequenceRangeUnsupportedSnafu, SerdeJsonSnafu,
140 SerializeColumnMetadataSnafu, SnapshotFenceStaleSnafu,
141};
142#[cfg(feature = "enterprise")]
143use crate::extension::BoxedExtensionRangeProviderFactory;
144use crate::gc::GcLimiterRef;
145use crate::manifest::action::RegionEdit;
146use crate::memtable::MemtableStats;
147use crate::metrics::{
148 HANDLE_REQUEST_ELAPSED, SCAN_MEMORY_EXHAUSTED_TOTAL, SCAN_MEMORY_USAGE_BYTES,
149 SCAN_REQUESTS_REJECTED_TOTAL,
150};
151use crate::read::scan_region::{ScanRegion, Scanner, exact_sequence_range};
152use crate::read::stream::ScanBatchStream;
153use crate::region::MitoRegionRef;
154use crate::region::opener::PartitionExprFetcherRef;
155use crate::region::options::parse_wal_options;
156use crate::request::{RegionEditRequest, WorkerRequest};
157use crate::sst::file::{FileMeta, RegionFileId, RegionIndexId};
158use crate::sst::file_ref::FileReferenceManagerRef;
159use crate::sst::index::intermediate::IntermediateManager;
160use crate::sst::index::puffin_manager::PuffinManagerFactory;
161use crate::wal::entry_distributor::{
162 DEFAULT_ENTRY_RECEIVER_BUFFER_SIZE, build_wal_entry_distributor_and_receivers,
163};
164use crate::wal::raw_entry_reader::{LogStoreRawEntryReader, RawEntryReader};
165use crate::worker::WorkerGroup;
166
167pub const MITO_ENGINE_NAME: &str = "mito";
168
169pub struct MitoEngineBuilder<'a, S: LogStore> {
170 data_home: &'a str,
171 config: MitoConfig,
172 log_store: Arc<S>,
173 object_store_manager: ObjectStoreManagerRef,
174 schema_metadata_manager: SchemaMetadataManagerRef,
175 file_ref_manager: FileReferenceManagerRef,
176 partition_expr_fetcher: PartitionExprFetcherRef,
177 plugins: Plugins,
178 #[cfg(feature = "enterprise")]
179 extension_range_provider_factory: Option<BoxedExtensionRangeProviderFactory>,
180}
181
182impl<'a, S: LogStore> MitoEngineBuilder<'a, S> {
183 #[allow(clippy::too_many_arguments)]
184 pub fn new(
185 data_home: &'a str,
186 config: MitoConfig,
187 log_store: Arc<S>,
188 object_store_manager: ObjectStoreManagerRef,
189 schema_metadata_manager: SchemaMetadataManagerRef,
190 file_ref_manager: FileReferenceManagerRef,
191 partition_expr_fetcher: PartitionExprFetcherRef,
192 plugins: Plugins,
193 ) -> Self {
194 Self {
195 data_home,
196 config,
197 log_store,
198 object_store_manager,
199 schema_metadata_manager,
200 file_ref_manager,
201 plugins,
202 partition_expr_fetcher,
203 #[cfg(feature = "enterprise")]
204 extension_range_provider_factory: None,
205 }
206 }
207
208 #[cfg(feature = "enterprise")]
209 #[must_use]
210 pub fn with_extension_range_provider_factory(
211 self,
212 extension_range_provider_factory: Option<BoxedExtensionRangeProviderFactory>,
213 ) -> Self {
214 Self {
215 extension_range_provider_factory,
216 ..self
217 }
218 }
219
220 pub async fn try_build(mut self) -> Result<MitoEngine> {
221 self.config.sanitize(self.data_home)?;
222
223 let config = Arc::new(self.config);
224 let region_hook = self.plugins.get::<RegionHookRef>();
227 let workers = WorkerGroup::start(
228 config.clone(),
229 self.log_store.clone(),
230 self.object_store_manager,
231 self.schema_metadata_manager,
232 self.file_ref_manager,
233 self.partition_expr_fetcher.clone(),
234 self.plugins,
235 )
236 .await?;
237 let wal_raw_entry_reader = Arc::new(LogStoreRawEntryReader::new(self.log_store));
238 let total_memory = get_total_memory_bytes().max(0) as u64;
239 let scan_memory_limit = config.scan_memory_limit.resolve(total_memory) as usize;
240 let scan_memory_pool = new_scan_memory_pool(scan_memory_limit);
241 let scan_memory_tracker =
242 QueryMemoryTracker::builder(scan_memory_limit, config.scan_memory_on_exhausted)
243 .on_update(|usage| {
244 SCAN_MEMORY_USAGE_BYTES.set(usage as i64);
245 })
246 .on_exhausted(|| {
247 SCAN_MEMORY_EXHAUSTED_TOTAL.inc();
248 })
249 .on_reject(|| {
250 SCAN_REQUESTS_REJECTED_TOTAL.inc();
251 })
252 .build();
253
254 let inner = EngineInner {
255 workers,
256 config,
257 wal_raw_entry_reader,
258 scan_memory_tracker,
259 scan_memory_pool,
260 region_hook,
261 #[cfg(feature = "enterprise")]
262 extension_range_provider_factory: None,
263 };
264
265 #[cfg(feature = "enterprise")]
266 let inner =
267 inner.with_extension_range_provider_factory(self.extension_range_provider_factory);
268
269 Ok(MitoEngine {
270 inner: Arc::new(inner),
271 })
272 }
273}
274
275#[derive(Clone)]
277pub struct MitoEngine {
278 inner: Arc<EngineInner>,
279}
280
281impl MitoEngine {
282 #[allow(clippy::too_many_arguments)]
284 pub async fn new<S: LogStore>(
285 data_home: &str,
286 config: MitoConfig,
287 log_store: Arc<S>,
288 object_store_manager: ObjectStoreManagerRef,
289 schema_metadata_manager: SchemaMetadataManagerRef,
290 file_ref_manager: FileReferenceManagerRef,
291 partition_expr_fetcher: PartitionExprFetcherRef,
292 plugins: Plugins,
293 ) -> Result<MitoEngine> {
294 let builder = MitoEngineBuilder::new(
295 data_home,
296 config,
297 log_store,
298 object_store_manager,
299 schema_metadata_manager,
300 file_ref_manager,
301 partition_expr_fetcher,
302 plugins,
303 );
304 builder.try_build().await
305 }
306
307 pub fn mito_config(&self) -> &MitoConfig {
308 &self.inner.config
309 }
310
311 pub fn cache_manager(&self) -> CacheManagerRef {
312 self.inner.workers.cache_manager()
313 }
314
315 pub fn file_ref_manager(&self) -> FileReferenceManagerRef {
316 self.inner.workers.file_ref_manager()
317 }
318
319 pub fn gc_limiter(&self) -> GcLimiterRef {
320 self.inner.workers.gc_limiter()
321 }
322
323 pub fn object_store_manager(&self) -> &ObjectStoreManagerRef {
324 self.inner.workers.object_store_manager()
325 }
326
327 pub fn puffin_manager_factory(&self) -> &PuffinManagerFactory {
328 self.inner.workers.puffin_manager_factory()
329 }
330
331 pub fn intermediate_manager(&self) -> &IntermediateManager {
332 self.inner.workers.intermediate_manager()
333 }
334
335 pub fn schema_metadata_manager(&self) -> &SchemaMetadataManagerRef {
336 self.inner.workers.schema_metadata_manager()
337 }
338
339 pub fn region_hook(&self) -> Option<RegionHookRef> {
342 self.inner.region_hook.clone()
343 }
344
345 pub async fn get_snapshot_of_file_refs(
347 &self,
348 file_handle_regions: impl IntoIterator<Item = RegionId>,
349 related_regions: HashMap<RegionId, HashSet<RegionId>>,
350 ) -> Result<FileRefsManifest> {
351 let file_ref_mgr = self.file_ref_manager();
352
353 let file_handle_regions = file_handle_regions.into_iter().collect::<Vec<_>>();
354 let query_regions: Vec<MitoRegionRef> = file_handle_regions
357 .into_iter()
358 .filter_map(|region_id| self.find_region(region_id))
359 .collect();
360
361 let dst_region_to_src_regions: Vec<(MitoRegionRef, HashSet<RegionId>)> = {
362 let dst2src = related_regions
363 .into_iter()
364 .flat_map(|(src, dsts)| dsts.into_iter().map(move |dst| (dst, src)))
365 .fold(
366 HashMap::<RegionId, HashSet<RegionId>>::new(),
367 |mut acc, (k, v)| {
368 let entry = acc.entry(k).or_default();
369 entry.insert(v);
370 acc
371 },
372 );
373 let mut dst_region_to_src_regions = Vec::with_capacity(dst2src.len());
374 for (dst_region, srcs) in dst2src {
375 let Some(region) = self.find_region(dst_region) else {
376 return RegionNotFoundSnafu {
377 region_id: dst_region,
378 }
379 .fail();
380 };
381 dst_region_to_src_regions.push((region, srcs));
382 }
383 dst_region_to_src_regions
384 };
385
386 file_ref_mgr
387 .get_snapshot_of_file_refs(query_regions, dst_region_to_src_regions)
388 .await
389 }
390
391 pub fn is_region_exists(&self, region_id: RegionId) -> bool {
393 self.inner.workers.is_region_exists(region_id)
394 }
395
396 pub fn is_region_opening(&self, region_id: RegionId) -> bool {
398 self.inner.workers.is_region_opening(region_id)
399 }
400
401 pub fn is_region_catching_up(&self, region_id: RegionId) -> bool {
403 self.inner.workers.is_region_catching_up(region_id)
404 }
405
406 pub fn get_region_statistic(&self, region_id: RegionId) -> Option<RegionStatistic> {
408 self.find_region(region_id)
409 .map(|region| region.region_statistic())
410 }
411
412 pub fn get_primary_key_encoding(&self, region_id: RegionId) -> Option<PrimaryKeyEncoding> {
414 self.find_region(region_id)
415 .map(|r| r.primary_key_encoding())
416 }
417
418 #[tracing::instrument(skip_all)]
423 pub async fn scan_to_stream(
424 &self,
425 region_id: RegionId,
426 request: ScanRequest,
427 ) -> Result<SendableRecordBatchStream, BoxedError> {
428 self.scanner(region_id, request)
429 .await
430 .map_err(BoxedError::new)?
431 .scan()
432 .await
433 }
434
435 pub async fn scan_batch(
437 &self,
438 region_id: RegionId,
439 request: ScanRequest,
440 filter_deleted: bool,
441 ) -> Result<ScanBatchStream> {
442 let mut scan_region = self.scan_region(region_id, request)?;
443 scan_region.set_filter_deleted(filter_deleted);
444 scan_region.scanner().await?.scan_batch()
445 }
446
447 pub(crate) async fn scanner(
449 &self,
450 region_id: RegionId,
451 request: ScanRequest,
452 ) -> Result<Scanner> {
453 self.scan_region(region_id, request)?.scanner().await
454 }
455
456 #[tracing::instrument(skip_all, fields(region_id = %region_id))]
458 fn scan_region(&self, region_id: RegionId, request: ScanRequest) -> Result<ScanRegion> {
459 self.inner.scan_region(region_id, request)
460 }
461
462 pub async fn edit_region(&self, region_id: RegionId, edit: RegionEdit) -> Result<()> {
467 let _timer = HANDLE_REQUEST_ELAPSED
468 .with_label_values(&["edit_region"])
469 .start_timer();
470
471 ensure!(
472 is_valid_region_edit(&edit),
473 InvalidRequestSnafu {
474 region_id,
475 reason: "invalid region edit"
476 }
477 );
478
479 let (tx, rx) = oneshot::channel();
480 let request = WorkerRequest::EditRegion(RegionEditRequest::new(region_id, edit, true, tx));
481 self.inner
482 .workers
483 .submit_to_worker(region_id, request)
484 .await?;
485 rx.await.context(RecvSnafu)?
486 }
487
488 pub async fn copy_region_from(
492 &self,
493 region_id: RegionId,
494 request: MitoCopyRegionFromRequest,
495 ) -> Result<MitoCopyRegionFromResponse> {
496 self.inner.copy_region_from(region_id, request).await
497 }
498
499 #[cfg(test)]
500 pub(crate) fn get_region(&self, id: RegionId) -> Option<crate::region::MitoRegionRef> {
501 self.find_region(id)
502 }
503
504 pub fn find_region(&self, region_id: RegionId) -> Option<MitoRegionRef> {
505 self.inner.workers.get_region(region_id)
506 }
507
508 pub fn regions(&self) -> Vec<MitoRegionRef> {
510 self.inner.workers.all_regions().collect()
511 }
512
513 fn encode_manifest_info_to_extensions(
514 region_id: &RegionId,
515 manifest_info: RegionManifestInfo,
516 extensions: &mut HashMap<String, Vec<u8>>,
517 ) -> Result<()> {
518 let region_manifest_info = vec![(*region_id, manifest_info)];
519
520 extensions.insert(
521 MANIFEST_INFO_EXTENSION_KEY.to_string(),
522 RegionManifestInfo::encode_list(®ion_manifest_info).context(SerdeJsonSnafu)?,
523 );
524 info!(
525 "Added manifest info: {:?} to extensions, region_id: {:?}",
526 region_manifest_info, region_id
527 );
528 Ok(())
529 }
530
531 fn encode_column_metadatas_to_extensions(
532 region_id: &RegionId,
533 column_metadatas: Vec<ColumnMetadata>,
534 extensions: &mut HashMap<String, Vec<u8>>,
535 ) -> Result<()> {
536 extensions.insert(
537 TABLE_COLUMN_METADATA_EXTENSION_KEY.to_string(),
538 ColumnMetadata::encode_list(&column_metadatas).context(SerializeColumnMetadataSnafu)?,
539 );
540 info!(
541 "Added column metadatas: {:?} to extensions, region_id: {:?}",
542 column_metadatas, region_id
543 );
544 Ok(())
545 }
546
547 pub fn find_memtable_and_sst_stats(
550 &self,
551 region_id: RegionId,
552 ) -> Result<(Vec<MemtableStats>, Vec<FileMeta>)> {
553 let region = self
554 .find_region(region_id)
555 .context(RegionNotFoundSnafu { region_id })?;
556
557 let version = region.version();
558 let memtable_stats = version
559 .memtables
560 .list_memtables()
561 .iter()
562 .map(|x| x.stats())
563 .collect::<Vec<_>>();
564
565 let sst_stats = version
566 .ssts
567 .levels()
568 .iter()
569 .flat_map(|level| level.files().map(|x| x.meta_ref()))
570 .cloned()
571 .collect::<Vec<_>>();
572 Ok((memtable_stats, sst_stats))
573 }
574
575 pub async fn all_ssts_from_manifest(&self) -> Vec<ManifestSstEntry> {
577 let node_id = self.inner.workers.file_ref_manager().node_id();
578 let regions = self.inner.workers.all_regions();
579
580 let mut results = Vec::new();
581 for region in regions {
582 let mut entries = region.manifest_sst_entries().await;
583 for e in &mut entries {
584 e.node_id = node_id;
585 }
586 results.extend(entries);
587 }
588
589 results
590 }
591
592 pub async fn all_index_metas(&self) -> Vec<PuffinIndexMetaEntry> {
594 let node_id = self.inner.workers.file_ref_manager().node_id();
595 let cache_manager = self.inner.workers.cache_manager();
596 let puffin_metadata_cache = cache_manager.puffin_metadata_cache().cloned();
597 let bloom_filter_cache = cache_manager.bloom_filter_index_cache().cloned();
598 let inverted_index_cache = cache_manager.inverted_index_cache().cloned();
599
600 let mut results = Vec::new();
601
602 for region in self.inner.workers.all_regions() {
603 let manifest_entries = region.manifest_sst_entries().await;
604 let access_layer = region.access_layer.clone();
605 let table_dir = access_layer.table_dir().to_string();
606 let path_type = access_layer.path_type();
607 let object_store = access_layer.object_store().clone();
608 let puffin_factory = access_layer.puffin_manager_factory().clone();
609 let path_factory = RegionFilePathFactory::new(table_dir, path_type);
610
611 let entry_futures = manifest_entries.into_iter().map(|entry| {
612 let object_store = object_store.clone();
613 let path_factory = path_factory.clone();
614 let puffin_factory = puffin_factory.clone();
615 let puffin_metadata_cache = puffin_metadata_cache.clone();
616 let bloom_filter_cache = bloom_filter_cache.clone();
617 let inverted_index_cache = inverted_index_cache.clone();
618
619 async move {
620 let Some(index_file_path) = entry.index_file_path.as_ref() else {
621 return Vec::new();
622 };
623
624 let index_version = entry.index_version;
625 let file_id = match FileId::parse_str(&entry.file_id) {
626 Ok(file_id) => file_id,
627 Err(err) => {
628 warn!(
629 err;
630 "Failed to parse puffin index file id, table_dir: {}, file_id: {}",
631 entry.table_dir,
632 entry.file_id
633 );
634 return Vec::new();
635 }
636 };
637 let region_index_id = RegionIndexId::new(
640 RegionFileId::new(entry.origin_region_id, file_id),
641 index_version,
642 );
643 let context = IndexEntryContext {
644 table_dir: &entry.table_dir,
645 index_file_path: index_file_path.as_str(),
646 region_id: entry.region_id,
647 table_id: entry.table_id,
648 region_number: entry.region_number,
649 region_group: entry.region_group,
650 region_sequence: entry.region_sequence,
651 file_id: &entry.file_id,
652 index_file_size: entry.index_file_size,
653 node_id,
654 };
655
656 let manager = puffin_factory
657 .build(object_store, path_factory)
658 .with_puffin_metadata_cache(puffin_metadata_cache);
659
660 collect_index_entries_from_puffin(
661 manager,
662 region_index_id,
663 context,
664 bloom_filter_cache,
665 inverted_index_cache,
666 )
667 .await
668 }
669 });
670
671 let mut meta_stream = stream::iter(entry_futures).buffer_unordered(8); while let Some(mut metas) = meta_stream.next().await {
673 results.append(&mut metas);
674 }
675 }
676
677 results
678 }
679
680 pub async fn all_region_infos(&self) -> Vec<RegionInfoEntry> {
682 let node_id = self.inner.workers.file_ref_manager().node_id();
683 self.inner
684 .workers
685 .all_regions()
686 .map(|region| region.region_info_entry(node_id))
687 .collect()
688 }
689
690 pub fn all_ssts_from_storage(&self) -> impl Stream<Item = Result<StorageSstEntry>> {
692 let node_id = self.inner.workers.file_ref_manager().node_id();
693 let regions = self.inner.workers.all_regions();
694
695 let mut layers_distinct_table_dirs = HashMap::new();
696 for region in regions {
697 let table_dir = region.access_layer.table_dir();
698 if !layers_distinct_table_dirs.contains_key(table_dir) {
699 layers_distinct_table_dirs
700 .insert(table_dir.to_string(), region.access_layer.clone());
701 }
702 }
703
704 stream::iter(layers_distinct_table_dirs)
705 .map(|(_, access_layer)| access_layer.storage_sst_entries())
706 .flatten()
707 .map(move |entry| {
708 entry.map(move |mut entry| {
709 entry.node_id = node_id;
710 entry
711 })
712 })
713 }
714}
715
716fn is_valid_region_edit(edit: &RegionEdit) -> bool {
720 (!edit.files_to_add.is_empty() || !edit.files_to_remove.is_empty())
721 && matches!(
722 edit,
723 RegionEdit {
724 files_to_add: _,
725 files_to_remove: _,
726 timestamp_ms: _,
727 compaction_time_window: None,
728 flushed_entry_id: None,
729 flushed_sequence: None,
730 ..
731 }
732 )
733}
734
735struct EngineInner {
737 workers: WorkerGroup,
739 config: Arc<MitoConfig>,
741 wal_raw_entry_reader: Arc<dyn RawEntryReader>,
743 scan_memory_tracker: QueryMemoryTracker,
745 scan_memory_pool: Arc<dyn MemoryPool>,
747 region_hook: Option<RegionHookRef>,
750 #[cfg(feature = "enterprise")]
751 extension_range_provider_factory: Option<BoxedExtensionRangeProviderFactory>,
752}
753
754type TopicGroupedRegionOpenRequests = HashMap<String, Vec<(RegionId, RegionOpenRequest)>>;
755
756fn prepare_batch_open_requests(
758 requests: Vec<(RegionId, RegionOpenRequest)>,
759) -> Result<(
760 TopicGroupedRegionOpenRequests,
761 Vec<(RegionId, RegionOpenRequest)>,
762)> {
763 let mut topic_to_regions: HashMap<String, Vec<(RegionId, RegionOpenRequest)>> = HashMap::new();
764 let mut remaining_regions: Vec<(RegionId, RegionOpenRequest)> = Vec::new();
765 for (region_id, request) in requests {
766 match parse_wal_options(&request.options).context(SerdeJsonSnafu)? {
767 WalOptions::Kafka(options) => {
768 topic_to_regions
769 .entry(options.topic)
770 .or_default()
771 .push((region_id, request));
772 }
773 WalOptions::RaftEngine | WalOptions::Noop => {
774 remaining_regions.push((region_id, request));
775 }
776 }
777 }
778
779 Ok((topic_to_regions, remaining_regions))
780}
781
782impl EngineInner {
783 #[cfg(feature = "enterprise")]
784 #[must_use]
785 fn with_extension_range_provider_factory(
786 self,
787 extension_range_provider_factory: Option<BoxedExtensionRangeProviderFactory>,
788 ) -> Self {
789 Self {
790 extension_range_provider_factory,
791 ..self
792 }
793 }
794
795 async fn stop(&self) -> Result<()> {
797 self.workers.stop().await
798 }
799
800 fn find_region(&self, region_id: RegionId) -> Result<MitoRegionRef> {
801 self.workers
802 .get_region(region_id)
803 .context(RegionNotFoundSnafu { region_id })
804 }
805
806 fn get_metadata(&self, region_id: RegionId) -> Result<RegionMetadataRef> {
810 let region = self.find_region(region_id)?;
812 Ok(region.metadata())
813 }
814
815 async fn open_topic_regions(
816 &self,
817 topic: String,
818 region_requests: Vec<(RegionId, RegionOpenRequest)>,
819 ) -> Result<Vec<(RegionId, Result<AffectedRows>)>> {
820 let now = Instant::now();
821 let region_ids = region_requests
822 .iter()
823 .map(|(region_id, _)| *region_id)
824 .collect::<Vec<_>>();
825 let provider = Provider::kafka_provider(topic);
826 let (distributor, entry_receivers) = build_wal_entry_distributor_and_receivers(
827 provider.clone(),
828 self.wal_raw_entry_reader.clone(),
829 ®ion_ids,
830 DEFAULT_ENTRY_RECEIVER_BUFFER_SIZE,
831 );
832
833 let mut responses = Vec::with_capacity(region_requests.len());
834 for ((region_id, request), entry_receiver) in
835 region_requests.into_iter().zip(entry_receivers)
836 {
837 let (request, receiver) =
838 WorkerRequest::new_open_region_request(region_id, request, Some(entry_receiver));
839 self.workers.submit_to_worker(region_id, request).await?;
840 responses.push(async move { receiver.await.context(RecvSnafu)? });
841 }
842
843 let distribution =
845 common_runtime::spawn_global(async move { distributor.distribute().await });
846 let responses = join_all(responses).await;
848 distribution.await.context(JoinSnafu)??;
849
850 let num_failure = responses.iter().filter(|r| r.is_err()).count();
851 info!(
852 "Opened {} regions for topic '{}', failures: {}, elapsed: {:?}",
853 region_ids.len() - num_failure,
854 provider.as_kafka_provider().unwrap(),
856 num_failure,
857 now.elapsed(),
858 );
859 Ok(region_ids.into_iter().zip(responses).collect())
860 }
861
862 async fn handle_batch_open_requests(
863 &self,
864 parallelism: usize,
865 requests: Vec<(RegionId, RegionOpenRequest)>,
866 ) -> Result<Vec<(RegionId, Result<AffectedRows>)>> {
867 let semaphore = Arc::new(Semaphore::new(parallelism));
868 let (topic_to_region_requests, remaining_region_requests) =
869 prepare_batch_open_requests(requests)?;
870 let mut responses =
871 Vec::with_capacity(topic_to_region_requests.len() + remaining_region_requests.len());
872
873 if !topic_to_region_requests.is_empty() {
874 let mut tasks = Vec::with_capacity(topic_to_region_requests.len());
875 for (topic, region_requests) in topic_to_region_requests {
876 let semaphore_moved = semaphore.clone();
877 tasks.push(async move {
878 let _permit = semaphore_moved.acquire().await.unwrap();
880 self.open_topic_regions(topic, region_requests).await
881 })
882 }
883 let r = try_join_all(tasks).await?;
884 responses.extend(r.into_iter().flatten());
885 }
886
887 if !remaining_region_requests.is_empty() {
888 let mut tasks = Vec::with_capacity(remaining_region_requests.len());
889 let mut region_ids = Vec::with_capacity(remaining_region_requests.len());
890 for (region_id, request) in remaining_region_requests {
891 let semaphore_moved = semaphore.clone();
892 region_ids.push(region_id);
893 tasks.push(async move {
894 let _permit = semaphore_moved.acquire().await.unwrap();
896 let (request, receiver) =
897 WorkerRequest::new_open_region_request(region_id, request, None);
898
899 self.workers.submit_to_worker(region_id, request).await?;
900
901 receiver.await.context(RecvSnafu)?
902 })
903 }
904
905 let results = join_all(tasks).await;
906 responses.extend(region_ids.into_iter().zip(results));
907 }
908
909 Ok(responses)
910 }
911
912 async fn catchup_topic_regions(
913 &self,
914 provider: Provider,
915 region_requests: Vec<(RegionId, RegionCatchupRequest)>,
916 ) -> Result<Vec<(RegionId, Result<AffectedRows>)>> {
917 let now = Instant::now();
918 let region_ids = region_requests
919 .iter()
920 .map(|(region_id, _)| *region_id)
921 .collect::<Vec<_>>();
922 let (distributor, entry_receivers) = build_wal_entry_distributor_and_receivers(
923 provider.clone(),
924 self.wal_raw_entry_reader.clone(),
925 ®ion_ids,
926 DEFAULT_ENTRY_RECEIVER_BUFFER_SIZE,
927 );
928
929 let mut responses = Vec::with_capacity(region_requests.len());
930 for ((region_id, request), entry_receiver) in
931 region_requests.into_iter().zip(entry_receivers)
932 {
933 let (request, receiver) =
934 WorkerRequest::new_catchup_region_request(region_id, request, Some(entry_receiver));
935 self.workers.submit_to_worker(region_id, request).await?;
936 responses.push(async move { receiver.await.context(RecvSnafu)? });
937 }
938
939 let distribution =
941 common_runtime::spawn_global(async move { distributor.distribute().await });
942 let responses = join_all(responses).await;
944 distribution.await.context(JoinSnafu)??;
945
946 let num_failure = responses.iter().filter(|r| r.is_err()).count();
947 info!(
948 "Caught up {} regions for topic '{}', failures: {}, elapsed: {:?}",
949 region_ids.len() - num_failure,
950 provider.as_kafka_provider().unwrap(),
952 num_failure,
953 now.elapsed(),
954 );
955
956 Ok(region_ids.into_iter().zip(responses).collect())
957 }
958
959 async fn handle_batch_catchup_requests(
960 &self,
961 parallelism: usize,
962 requests: Vec<(RegionId, RegionCatchupRequest)>,
963 ) -> Result<Vec<(RegionId, Result<AffectedRows>)>> {
964 let mut responses = Vec::with_capacity(requests.len());
965 let mut topic_regions: HashMap<Arc<KafkaProvider>, Vec<_>> = HashMap::new();
966 let mut remaining_region_requests = vec![];
967
968 for (region_id, request) in requests {
969 match self.workers.get_region(region_id) {
970 Some(region) => match region.provider.as_kafka_provider() {
971 Some(provider) => {
972 topic_regions
973 .entry(provider.clone())
974 .or_default()
975 .push((region_id, request));
976 }
977 None => {
978 remaining_region_requests.push((region_id, request));
979 }
980 },
981 None => responses.push((region_id, RegionNotFoundSnafu { region_id }.fail())),
982 }
983 }
984
985 let semaphore = Arc::new(Semaphore::new(parallelism));
986
987 if !topic_regions.is_empty() {
988 let mut tasks = Vec::with_capacity(topic_regions.len());
989 for (provider, region_requests) in topic_regions {
990 let semaphore_moved = semaphore.clone();
991 tasks.push(async move {
992 let _permit = semaphore_moved.acquire().await.unwrap();
994 self.catchup_topic_regions(Provider::Kafka(provider), region_requests)
995 .await
996 })
997 }
998
999 let r = try_join_all(tasks).await?;
1000 responses.extend(r.into_iter().flatten());
1001 }
1002
1003 if !remaining_region_requests.is_empty() {
1004 let mut tasks = Vec::with_capacity(remaining_region_requests.len());
1005 let mut region_ids = Vec::with_capacity(remaining_region_requests.len());
1006 for (region_id, request) in remaining_region_requests {
1007 let semaphore_moved = semaphore.clone();
1008 region_ids.push(region_id);
1009 tasks.push(async move {
1010 let _permit = semaphore_moved.acquire().await.unwrap();
1012 let (request, receiver) =
1013 WorkerRequest::new_catchup_region_request(region_id, request, None);
1014
1015 self.workers.submit_to_worker(region_id, request).await?;
1016
1017 receiver.await.context(RecvSnafu)?
1018 })
1019 }
1020
1021 let results = join_all(tasks).await;
1022 responses.extend(region_ids.into_iter().zip(results));
1023 }
1024
1025 Ok(responses)
1026 }
1027
1028 async fn handle_request(
1030 &self,
1031 region_id: RegionId,
1032 request: RegionRequest,
1033 ) -> Result<AffectedRows> {
1034 let region_metadata = self.get_metadata(region_id).ok();
1035 let (request, receiver) =
1036 WorkerRequest::try_from_region_request(region_id, request, region_metadata)?;
1037 self.workers.submit_to_worker(region_id, request).await?;
1038
1039 receiver.await.context(RecvSnafu)?
1040 }
1041
1042 fn get_committed_sequence(&self, region_id: RegionId) -> Result<SequenceNumber> {
1044 self.find_region(region_id)
1046 .map(|r| r.find_committed_sequence())
1047 }
1048
1049 #[tracing::instrument(skip_all, fields(region_id = %region_id))]
1051 fn scan_region(&self, region_id: RegionId, mut request: ScanRequest) -> Result<ScanRegion> {
1052 let query_start = Instant::now();
1053 let region = self.find_region(region_id)?;
1055 let version_data = region.version_control.current();
1056 let version = version_data.version;
1057
1058 if request.snapshot_on_scan && request.memtable_max_sequence.is_none() {
1059 request.memtable_max_sequence = Some(version_data.committed_sequence);
1060 }
1061
1062 let exact_selection = if request.exact_sequence_range {
1066 match exact_sequence_range(&request, &version) {
1067 Ok((files, sequence_range)) => Some((files, sequence_range)),
1068 Err(err) => {
1069 debug!(
1070 "Scan region {} exact sequence range denied: min={:?}, max={:?}, denial_reason=foreign_file_missing_barrier",
1071 region_id, request.memtable_min_sequence, request.memtable_max_sequence,
1072 );
1073 return Err(err);
1074 }
1075 }
1076 } else {
1077 None
1078 };
1079 let exact_sequence_range = exact_selection
1080 .as_ref()
1081 .and_then(|(_, sequence_range)| *sequence_range);
1082
1083 #[cfg(feature = "enterprise")]
1090 let extension_provider_blocks_exact = region.is_follower()
1091 && self.extension_range_provider_factory.is_some()
1092 && exact_sequence_range.is_some();
1093 #[cfg(not(feature = "enterprise"))]
1094 let extension_provider_blocks_exact = false;
1095 let exact_sequence_range = if extension_provider_blocks_exact {
1096 None
1097 } else {
1098 exact_sequence_range
1099 };
1100 let exact_selection = exact_selection.map(|(files, _)| (files, exact_sequence_range));
1101 let exact_denial_reason = if !request.exact_sequence_range {
1102 "not_requested"
1103 } else if request.skip_sst_files {
1104 "sst_files_skipped"
1105 } else if request.memtable_min_sequence.is_none() {
1106 "missing_lower_bound"
1107 } else if request.memtable_max_sequence.is_none() {
1108 "missing_upper_bound"
1109 } else if !version.options.preserve_row_sequence {
1110 "preserve_row_sequence_disabled"
1111 } else if extension_provider_blocks_exact {
1112 "extension_provider"
1113 } else if exact_sequence_range.is_none() {
1114 "selected_file_barrier_not_admitted"
1115 } else {
1116 "none"
1117 };
1118
1119 debug!(
1120 "Scan region {} exact sequence range: requested={}, min={:?}, max={:?}, committed={}, flushed={}, available={}, selected_files={}, denial_reason={}",
1121 region_id,
1122 request.exact_sequence_range,
1123 request.memtable_min_sequence,
1124 request.memtable_max_sequence,
1125 version_data.committed_sequence,
1126 version.flushed_sequence,
1127 exact_sequence_range.is_some(),
1128 exact_selection.as_ref().map_or(0, |(files, _)| files.len()),
1129 exact_denial_reason,
1130 );
1131 validate_sequence_fences(
1132 &request,
1133 &version,
1134 region_id,
1135 exact_sequence_range.is_some(),
1136 extension_provider_blocks_exact,
1137 )?;
1138
1139 let cache_manager = self.workers.cache_manager();
1141
1142 let scan_region = ScanRegion::new(
1143 version,
1144 region.access_layer.clone(),
1145 request,
1146 CacheStrategy::EnableAll(cache_manager),
1147 )
1148 .with_query_stat_counters(region.region_stats.query_stat_counters())
1149 .with_max_concurrent_scan_files(self.config.max_concurrent_scan_files)
1150 .with_scan_memory_pool(self.scan_memory_pool.clone())
1151 .with_experimental_series_scan_v2(self.config.experimental_series_scan_v2)
1152 .with_ignore_inverted_index(self.config.inverted_index.apply_on_query.disabled())
1153 .with_ignore_fulltext_index(self.config.fulltext_index.apply_on_query.disabled())
1154 .with_ignore_bloom_filter(self.config.bloom_filter_index.apply_on_query.disabled())
1155 .with_start_time(query_start);
1156 let scan_region = if let Some(selection) = exact_selection {
1157 scan_region.with_exact_selection(selection)
1158 } else {
1159 scan_region
1160 };
1161
1162 #[cfg(feature = "enterprise")]
1163 let scan_region = self.maybe_fill_extension_range_provider(scan_region, region);
1164
1165 Ok(scan_region)
1166 }
1167
1168 #[cfg(feature = "enterprise")]
1169 fn maybe_fill_extension_range_provider(
1170 &self,
1171 mut scan_region: ScanRegion,
1172 region: MitoRegionRef,
1173 ) -> ScanRegion {
1174 if region.is_follower()
1175 && let Some(factory) = self.extension_range_provider_factory.as_ref()
1176 {
1177 scan_region
1178 .set_extension_range_provider(factory.create_extension_range_provider(region));
1179 }
1180 scan_region
1181 }
1182
1183 fn set_region_role(&self, region_id: RegionId, role: RegionRole) -> Result<()> {
1185 let region = self.find_region(region_id)?;
1186 region.set_role(role);
1187 Ok(())
1188 }
1189
1190 async fn set_region_role_state_gracefully(
1192 &self,
1193 region_id: RegionId,
1194 region_role_state: SettableRegionRoleState,
1195 ) -> Result<SetRegionRoleStateResponse> {
1196 let (request, receiver) =
1199 WorkerRequest::new_set_readonly_gracefully(region_id, region_role_state);
1200 self.workers.submit_to_worker(region_id, request).await?;
1201
1202 receiver.await.context(RecvSnafu)
1203 }
1204
1205 async fn sync_region(
1206 &self,
1207 region_id: RegionId,
1208 manifest_info: RegionManifestInfo,
1209 ) -> Result<(ManifestVersion, bool)> {
1210 ensure!(manifest_info.is_mito(), MitoManifestInfoSnafu);
1211 let manifest_version = manifest_info.data_manifest_version();
1212 let (request, receiver) =
1213 WorkerRequest::new_sync_region_request(region_id, manifest_version);
1214 self.workers.submit_to_worker(region_id, request).await?;
1215
1216 receiver.await.context(RecvSnafu)?
1217 }
1218
1219 async fn remap_manifests(
1220 &self,
1221 request: RemapManifestsRequest,
1222 ) -> Result<RemapManifestsResponse> {
1223 let region_id = request.region_id;
1224 let (request, receiver) = WorkerRequest::try_from_remap_manifests_request(request)?;
1225 self.workers.submit_to_worker(region_id, request).await?;
1226 let manifest_paths = receiver.await.context(RecvSnafu)??;
1227 Ok(RemapManifestsResponse { manifest_paths })
1228 }
1229
1230 async fn copy_region_from(
1231 &self,
1232 region_id: RegionId,
1233 request: MitoCopyRegionFromRequest,
1234 ) -> Result<MitoCopyRegionFromResponse> {
1235 let (request, receiver) =
1236 WorkerRequest::try_from_copy_region_from_request(region_id, request)?;
1237 self.workers.submit_to_worker(region_id, request).await?;
1238 let response = receiver.await.context(RecvSnafu)??;
1239 Ok(response)
1240 }
1241
1242 fn role(&self, region_id: RegionId) -> Option<RegionRole> {
1243 self.workers
1244 .get_region(region_id)
1245 .map(|region| region.region_role())
1246 }
1247}
1248
1249fn validate_sequence_fences(
1250 request: &ScanRequest,
1251 version: &crate::region::version::Version,
1252 region_id: RegionId,
1253 exact_sequence_range_available: bool,
1254 extension_provider_blocks_exact: bool,
1255) -> Result<()> {
1256 if request.exact_sequence_range {
1257 if !exact_sequence_range_available {
1263 if let Some(given_seq) = request.memtable_min_sequence {
1264 let min_readable_seq = version.flushed_sequence;
1265 ensure!(
1266 given_seq >= min_readable_seq,
1267 IncrementalQueryStaleSnafu {
1268 region_id,
1269 given_seq,
1270 min_readable_seq,
1271 }
1272 );
1273 }
1274
1275 if let Some(given_seq) = request.memtable_max_sequence {
1276 let min_enforceable_seq = version.flushed_sequence;
1277 ensure!(
1278 given_seq >= min_enforceable_seq,
1279 SnapshotFenceStaleSnafu {
1280 region_id,
1281 given_seq,
1282 min_enforceable_seq,
1283 }
1284 );
1285 }
1286
1287 return SequenceRangeUnsupportedSnafu {
1293 region_id,
1294 min_seq: request.memtable_min_sequence.unwrap_or_default(),
1295 max_seq: request.memtable_max_sequence.unwrap_or_default(),
1296 reason: sequence_range_unsupported_reason(version, extension_provider_blocks_exact),
1297 }
1298 .fail();
1299 }
1300 } else {
1301 if let Some(given_seq) = request.memtable_min_sequence {
1305 let min_readable_seq = version.flushed_sequence;
1306 ensure!(
1307 given_seq >= min_readable_seq,
1308 IncrementalQueryStaleSnafu {
1309 region_id,
1310 given_seq,
1311 min_readable_seq,
1312 }
1313 );
1314 }
1315
1316 if let Some(given_seq) = request.memtable_max_sequence
1317 && !request.skip_sst_files
1318 {
1319 let min_enforceable_seq = version.flushed_sequence;
1326 ensure!(
1327 given_seq >= min_enforceable_seq,
1328 SnapshotFenceStaleSnafu {
1329 region_id,
1330 given_seq,
1331 min_enforceable_seq,
1332 }
1333 );
1334 }
1335 }
1336
1337 Ok(())
1338}
1339
1340fn sequence_range_unsupported_reason(
1341 version: &crate::region::version::Version,
1342 extension_provider_blocks_exact: bool,
1343) -> String {
1344 if extension_provider_blocks_exact {
1345 "region is a follower with an extension range provider, whose streams cannot be filtered by sequence".to_string()
1346 } else if version.options.preserve_row_sequence {
1347 "region has files without preserved per-row sequences".to_string()
1348 } else {
1349 "region does not preserve per-row sequences (preserve_row_sequence is off)".to_string()
1350 }
1351}
1352
1353fn map_batch_responses(responses: Vec<(RegionId, Result<AffectedRows>)>) -> BatchResponses {
1354 responses
1355 .into_iter()
1356 .map(|(region_id, response)| {
1357 (
1358 region_id,
1359 response.map(RegionResponse::new).map_err(BoxedError::new),
1360 )
1361 })
1362 .collect()
1363}
1364
1365#[async_trait]
1366impl RegionEngine for MitoEngine {
1367 fn name(&self) -> &str {
1368 MITO_ENGINE_NAME
1369 }
1370
1371 #[tracing::instrument(skip_all)]
1372 async fn handle_batch_open_requests(
1373 &self,
1374 parallelism: usize,
1375 requests: Vec<(RegionId, RegionOpenRequest)>,
1376 ) -> Result<BatchResponses, BoxedError> {
1377 self.inner
1379 .handle_batch_open_requests(parallelism, requests)
1380 .await
1381 .map(map_batch_responses)
1382 .map_err(BoxedError::new)
1383 }
1384
1385 #[tracing::instrument(skip_all)]
1386 async fn handle_batch_catchup_requests(
1387 &self,
1388 parallelism: usize,
1389 requests: Vec<(RegionId, RegionCatchupRequest)>,
1390 ) -> Result<BatchResponses, BoxedError> {
1391 self.inner
1392 .handle_batch_catchup_requests(parallelism, requests)
1393 .await
1394 .map(map_batch_responses)
1395 .map_err(BoxedError::new)
1396 }
1397
1398 #[tracing::instrument(skip_all)]
1399 async fn handle_request(
1400 &self,
1401 region_id: RegionId,
1402 request: RegionRequest,
1403 ) -> Result<RegionResponse, BoxedError> {
1404 let _timer = HANDLE_REQUEST_ELAPSED
1405 .with_label_values(&[request.request_type()])
1406 .start_timer();
1407
1408 let is_alter = matches!(request, RegionRequest::Alter(_));
1409 let is_create = matches!(request, RegionRequest::Create(_));
1410 let mut response = self
1411 .inner
1412 .handle_request(region_id, request)
1413 .await
1414 .map(RegionResponse::new)
1415 .map_err(BoxedError::new)?;
1416
1417 if is_alter {
1418 self.handle_alter_response(region_id, &mut response)
1419 .map_err(BoxedError::new)?;
1420 } else if is_create {
1421 self.handle_create_response(region_id, &mut response)
1422 .map_err(BoxedError::new)?;
1423 }
1424
1425 Ok(response)
1426 }
1427
1428 #[tracing::instrument(skip_all)]
1429 async fn handle_query(
1430 &self,
1431 region_id: RegionId,
1432 request: ScanRequest,
1433 ) -> Result<RegionScannerRef, BoxedError> {
1434 self.scan_region(region_id, request)
1435 .map_err(BoxedError::new)?
1436 .region_scanner()
1437 .await
1438 .map_err(BoxedError::new)
1439 }
1440
1441 fn query_memory_tracker(&self) -> Option<QueryMemoryTracker> {
1442 Some(self.inner.scan_memory_tracker.clone())
1443 }
1444
1445 async fn get_committed_sequence(
1446 &self,
1447 region_id: RegionId,
1448 ) -> Result<SequenceNumber, BoxedError> {
1449 self.inner
1450 .get_committed_sequence(region_id)
1451 .map_err(BoxedError::new)
1452 }
1453
1454 async fn get_metadata(
1456 &self,
1457 region_id: RegionId,
1458 ) -> std::result::Result<RegionMetadataRef, BoxedError> {
1459 self.inner.get_metadata(region_id).map_err(BoxedError::new)
1460 }
1461
1462 async fn stop(&self) -> std::result::Result<(), BoxedError> {
1468 self.inner.stop().await.map_err(BoxedError::new)
1469 }
1470
1471 fn region_statistic(&self, region_id: RegionId) -> Option<RegionStatistic> {
1472 self.get_region_statistic(region_id)
1473 }
1474
1475 fn set_region_role(&self, region_id: RegionId, role: RegionRole) -> Result<(), BoxedError> {
1476 self.inner
1477 .set_region_role(region_id, role)
1478 .map_err(BoxedError::new)
1479 }
1480
1481 async fn set_region_role_state_gracefully(
1482 &self,
1483 region_id: RegionId,
1484 region_role_state: SettableRegionRoleState,
1485 ) -> Result<SetRegionRoleStateResponse, BoxedError> {
1486 let _timer = HANDLE_REQUEST_ELAPSED
1487 .with_label_values(&["set_region_role_state_gracefully"])
1488 .start_timer();
1489
1490 self.inner
1491 .set_region_role_state_gracefully(region_id, region_role_state)
1492 .await
1493 .map_err(BoxedError::new)
1494 }
1495
1496 async fn sync_region(
1497 &self,
1498 region_id: RegionId,
1499 request: SyncRegionFromRequest,
1500 ) -> Result<SyncRegionFromResponse, BoxedError> {
1501 let manifest_info = request
1502 .into_region_manifest_info()
1503 .context(UnexpectedSnafu {
1504 err_msg: "Expected a manifest info request",
1505 })
1506 .map_err(BoxedError::new)?;
1507 let (_, synced) = self
1508 .inner
1509 .sync_region(region_id, manifest_info)
1510 .await
1511 .map_err(BoxedError::new)?;
1512
1513 Ok(SyncRegionFromResponse::Mito { synced })
1514 }
1515
1516 async fn remap_manifests(
1517 &self,
1518 request: RemapManifestsRequest,
1519 ) -> Result<RemapManifestsResponse, BoxedError> {
1520 self.inner
1521 .remap_manifests(request)
1522 .await
1523 .map_err(BoxedError::new)
1524 }
1525
1526 fn role(&self, region_id: RegionId) -> Option<RegionRole> {
1527 self.inner.role(region_id)
1528 }
1529
1530 fn as_any(&self) -> &dyn Any {
1531 self
1532 }
1533}
1534
1535impl MitoEngine {
1536 fn handle_alter_response(
1537 &self,
1538 region_id: RegionId,
1539 response: &mut RegionResponse,
1540 ) -> Result<()> {
1541 if let Some(statistic) = self.region_statistic(region_id) {
1542 Self::encode_manifest_info_to_extensions(
1543 ®ion_id,
1544 statistic.manifest,
1545 &mut response.extensions,
1546 )?;
1547 }
1548 let column_metadatas = self
1549 .inner
1550 .find_region(region_id)
1551 .ok()
1552 .map(|r| r.metadata().column_metadatas.clone());
1553 if let Some(column_metadatas) = column_metadatas {
1554 Self::encode_column_metadatas_to_extensions(
1555 ®ion_id,
1556 column_metadatas,
1557 &mut response.extensions,
1558 )?;
1559 }
1560 Ok(())
1561 }
1562
1563 fn handle_create_response(
1564 &self,
1565 region_id: RegionId,
1566 response: &mut RegionResponse,
1567 ) -> Result<()> {
1568 let column_metadatas = self
1569 .inner
1570 .find_region(region_id)
1571 .ok()
1572 .map(|r| r.metadata().column_metadatas.clone());
1573 if let Some(column_metadatas) = column_metadatas {
1574 Self::encode_column_metadatas_to_extensions(
1575 ®ion_id,
1576 column_metadatas,
1577 &mut response.extensions,
1578 )?;
1579 }
1580 Ok(())
1581 }
1582}
1583
1584#[cfg(any(test, feature = "test"))]
1586#[allow(clippy::too_many_arguments)]
1587impl MitoEngine {
1588 pub async fn new_for_test<S: LogStore>(
1590 data_home: &str,
1591 mut config: MitoConfig,
1592 log_store: Arc<S>,
1593 object_store_manager: ObjectStoreManagerRef,
1594 write_buffer_manager: Option<crate::flush::WriteBufferManagerRef>,
1595 listener: Option<crate::engine::listener::EventListenerRef>,
1596 time_provider: crate::time_provider::TimeProviderRef,
1597 schema_metadata_manager: SchemaMetadataManagerRef,
1598 file_ref_manager: FileReferenceManagerRef,
1599 partition_expr_fetcher: PartitionExprFetcherRef,
1600 ) -> Result<MitoEngine> {
1601 config.sanitize(data_home)?;
1602
1603 let config = Arc::new(config);
1604 let wal_raw_entry_reader = Arc::new(LogStoreRawEntryReader::new(log_store.clone()));
1605 let total_memory = get_total_memory_bytes().max(0) as u64;
1606 let scan_memory_limit = config.scan_memory_limit.resolve(total_memory) as usize;
1607 let scan_memory_pool = new_scan_memory_pool(scan_memory_limit);
1608 let scan_memory_tracker =
1609 QueryMemoryTracker::builder(scan_memory_limit, config.scan_memory_on_exhausted)
1610 .on_update(|usage| {
1611 SCAN_MEMORY_USAGE_BYTES.set(usage as i64);
1612 })
1613 .on_exhausted(|| {
1614 SCAN_MEMORY_EXHAUSTED_TOTAL.inc();
1615 })
1616 .on_reject(|| {
1617 SCAN_REQUESTS_REJECTED_TOTAL.inc();
1618 })
1619 .build();
1620 Ok(MitoEngine {
1621 inner: Arc::new(EngineInner {
1622 workers: WorkerGroup::start_for_test(
1623 config.clone(),
1624 log_store,
1625 object_store_manager,
1626 write_buffer_manager,
1627 listener,
1628 schema_metadata_manager,
1629 file_ref_manager,
1630 time_provider,
1631 partition_expr_fetcher,
1632 )
1633 .await?,
1634 config,
1635 wal_raw_entry_reader,
1636 scan_memory_tracker,
1637 scan_memory_pool,
1638 region_hook: None,
1639 #[cfg(feature = "enterprise")]
1640 extension_range_provider_factory: None,
1641 }),
1642 })
1643 }
1644
1645 pub fn purge_scheduler(&self) -> &crate::schedule::scheduler::SchedulerRef {
1647 self.inner.workers.purge_scheduler()
1648 }
1649}
1650
1651fn new_scan_memory_pool(scan_memory_limit: usize) -> Arc<dyn MemoryPool> {
1652 if scan_memory_limit == 0 {
1653 Arc::new(UnboundedMemoryPool::default())
1654 } else {
1655 Arc::new(GreedyMemoryPool::new(scan_memory_limit))
1656 }
1657}
1658
1659#[cfg(test)]
1660mod tests {
1661 use std::time::Duration;
1662
1663 use datafusion::execution::memory_pool::MemoryConsumer;
1664
1665 use super::*;
1666 use crate::sst::file::FileMeta;
1667
1668 #[test]
1669 fn test_is_valid_region_edit() {
1670 let edit = RegionEdit {
1672 files_to_add: vec![FileMeta::default()],
1673 files_to_remove: vec![],
1674 timestamp_ms: None,
1675 compaction_time_window: None,
1676 flushed_entry_id: None,
1677 flushed_sequence: None,
1678 committed_sequence: None,
1679 };
1680 assert!(is_valid_region_edit(&edit));
1681
1682 let edit = RegionEdit {
1684 files_to_add: vec![],
1685 files_to_remove: vec![],
1686 timestamp_ms: None,
1687 compaction_time_window: None,
1688 flushed_entry_id: None,
1689 flushed_sequence: None,
1690 committed_sequence: None,
1691 };
1692 assert!(!is_valid_region_edit(&edit));
1693
1694 let edit = RegionEdit {
1696 files_to_add: vec![],
1697 files_to_remove: vec![FileMeta::default()],
1698 timestamp_ms: None,
1699 compaction_time_window: None,
1700 flushed_entry_id: None,
1701 flushed_sequence: None,
1702 committed_sequence: None,
1703 };
1704 assert!(is_valid_region_edit(&edit));
1705
1706 let edit = RegionEdit {
1708 files_to_add: vec![FileMeta::default()],
1709 files_to_remove: vec![FileMeta::default()],
1710 timestamp_ms: None,
1711 compaction_time_window: None,
1712 flushed_entry_id: None,
1713 flushed_sequence: None,
1714 committed_sequence: None,
1715 };
1716 assert!(is_valid_region_edit(&edit));
1717
1718 let edit = RegionEdit {
1720 files_to_add: vec![FileMeta::default()],
1721 files_to_remove: vec![],
1722 timestamp_ms: None,
1723 compaction_time_window: Some(Duration::from_secs(1)),
1724 flushed_entry_id: None,
1725 flushed_sequence: None,
1726 committed_sequence: None,
1727 };
1728 assert!(!is_valid_region_edit(&edit));
1729 let edit = RegionEdit {
1730 files_to_add: vec![FileMeta::default()],
1731 files_to_remove: vec![],
1732 timestamp_ms: None,
1733 compaction_time_window: None,
1734 flushed_entry_id: Some(1),
1735 flushed_sequence: None,
1736 committed_sequence: None,
1737 };
1738 assert!(!is_valid_region_edit(&edit));
1739 let edit = RegionEdit {
1740 files_to_add: vec![FileMeta::default()],
1741 files_to_remove: vec![],
1742 timestamp_ms: None,
1743 compaction_time_window: None,
1744 flushed_entry_id: None,
1745 flushed_sequence: Some(1),
1746 committed_sequence: None,
1747 };
1748 assert!(!is_valid_region_edit(&edit));
1749 }
1750
1751 #[test]
1752 fn test_scan_memory_pool_is_shared_across_consumers() {
1753 let pool = new_scan_memory_pool(100);
1754 let cloned_pool = pool.clone();
1755 let first = MemoryConsumer::new("first-scan").register(&pool);
1756 let second = MemoryConsumer::new("second-scan").register(&cloned_pool);
1757
1758 first.try_grow(60).unwrap();
1759 assert!(second.try_grow(50).is_err());
1760 second.try_grow(40).unwrap();
1761 assert_eq!(100, pool.reserved());
1762 }
1763}