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