1mod handle_alter;
18mod handle_apply_staging;
19mod handle_bulk_insert;
20mod handle_catchup;
21mod handle_close;
22mod handle_compaction;
23mod handle_copy_region;
24mod handle_create;
25mod handle_drop;
26mod handle_enter_staging;
27mod handle_flush;
28mod handle_manifest;
29mod handle_open;
30mod handle_rebuild_index;
31mod handle_remap;
32mod handle_truncate;
33mod handle_write;
34
35use std::collections::HashMap;
36use std::path::Path;
37use std::sync::Arc;
38use std::sync::atomic::{AtomicBool, Ordering};
39use std::time::Duration;
40
41use common_base::Plugins;
42use common_error::ext::BoxedError;
43use common_meta::key::SchemaMetadataManagerRef;
44use common_runtime::JoinHandle;
45use common_stat::get_total_memory_bytes;
46use common_telemetry::{error, info, warn};
47use futures::future::try_join_all;
48use object_store::manager::ObjectStoreManagerRef;
49use prometheus::{Histogram, IntGauge};
50use rand::{Rng, rng};
51use snafu::{ResultExt, ensure};
52use store_api::logstore::LogStore;
53use store_api::region_engine::{
54 SetRegionRoleStateResponse, SetRegionRoleStateSuccess, SettableRegionRoleState,
55};
56use store_api::storage::{FileId, RegionId};
57use tokio::sync::mpsc::{Receiver, Sender};
58use tokio::sync::{Mutex, Semaphore, mpsc, oneshot, watch};
59
60use crate::cache::write_cache::{WriteCache, WriteCacheRef};
61use crate::cache::{CacheManager, CacheManagerRef};
62use crate::compaction::CompactionScheduler;
63use crate::compaction::memory_manager::{CompactionMemoryManager, new_compaction_memory_manager};
64use crate::config::MitoConfig;
65use crate::error::{self, CreateDirSnafu, JoinSnafu, Result, WorkerStoppedSnafu};
66use crate::flush::{FlushScheduler, WriteBufferManagerImpl, WriteBufferManagerRef};
67use crate::gc::{GcLimiter, GcLimiterRef};
68use crate::memtable::MemtableBuilderProvider;
69use crate::metrics::{REGION_COUNT, REQUEST_WAIT_TIME, WRITE_STALLING};
70use crate::region::opener::PartitionExprFetcherRef;
71use crate::region::{
72 CatchupRegions, CatchupRegionsRef, MitoRegionRef, OpeningRegions, OpeningRegionsRef, RegionMap,
73 RegionMapRef,
74};
75use crate::request::{
76 BackgroundNotify, BulkInsertRequest, DdlRequest, OptionOutputTx, SenderBulkRequest,
77 SenderDdlRequest, SenderWriteRequest, WorkerRequest, WorkerRequestWithTime,
78};
79use crate::schedule::scheduler::{LocalScheduler, SchedulerRef};
80use crate::sst::file::RegionFileId;
81use crate::sst::file_ref::FileReferenceManagerRef;
82use crate::sst::index::IndexBuildScheduler;
83use crate::sst::index::intermediate::IntermediateManager;
84use crate::sst::index::puffin_manager::PuffinManagerFactory;
85use crate::time_provider::{StdTimeProvider, TimeProviderRef};
86use crate::wal::Wal;
87use crate::worker::handle_manifest::RegionEditQueues;
88
89pub(crate) type WorkerId = u32;
91
92pub(crate) const DROPPING_MARKER_FILE: &str = ".dropping";
93
94pub(crate) const CHECK_REGION_INTERVAL: Duration = Duration::from_secs(60);
96pub(crate) const MAX_INITIAL_CHECK_DELAY_SECS: u64 = 60 * 3;
98
99#[cfg_attr(doc, aquamarine::aquamarine)]
100pub(crate) struct WorkerGroup {
137 workers: Vec<RegionWorker>,
139 flush_job_pool: SchedulerRef,
141 compact_job_pool: SchedulerRef,
143 index_build_job_pool: SchedulerRef,
145 purge_scheduler: SchedulerRef,
147 cache_manager: CacheManagerRef,
149 file_ref_manager: FileReferenceManagerRef,
151 gc_limiter: GcLimiterRef,
153 object_store_manager: ObjectStoreManagerRef,
155 puffin_manager_factory: PuffinManagerFactory,
157 intermediate_manager: IntermediateManager,
159 schema_metadata_manager: SchemaMetadataManagerRef,
161}
162
163impl WorkerGroup {
164 pub(crate) async fn start<S: LogStore>(
168 config: Arc<MitoConfig>,
169 log_store: Arc<S>,
170 object_store_manager: ObjectStoreManagerRef,
171 schema_metadata_manager: SchemaMetadataManagerRef,
172 file_ref_manager: FileReferenceManagerRef,
173 partition_expr_fetcher: PartitionExprFetcherRef,
174 plugins: Plugins,
175 ) -> Result<WorkerGroup> {
176 let (flush_sender, flush_receiver) = watch::channel(());
177 let write_buffer_manager = Arc::new(
178 WriteBufferManagerImpl::new(config.global_write_buffer_size.as_bytes() as usize)
179 .with_notifier(flush_sender.clone()),
180 );
181 let puffin_manager_factory = PuffinManagerFactory::new(
182 &config.index.aux_path,
183 config.index.staging_size.as_bytes(),
184 Some(config.index.write_buffer_size.as_bytes() as _),
185 config.index.staging_ttl,
186 )
187 .await?;
188 let intermediate_manager = IntermediateManager::init_fs(&config.index.aux_path)
189 .await?
190 .with_buffer_size(Some(config.index.write_buffer_size.as_bytes() as _));
191 let index_build_job_pool =
192 Arc::new(LocalScheduler::new(config.max_background_index_builds));
193 let flush_job_pool = Arc::new(LocalScheduler::new(config.max_background_flushes));
194 let compact_job_pool = Arc::new(LocalScheduler::new(config.max_background_compactions));
195 let flush_semaphore = Arc::new(Semaphore::new(config.max_background_flushes));
196 let purge_scheduler = Arc::new(LocalScheduler::new(config.max_background_purges));
198 let write_cache = write_cache_from_config(
199 &config,
200 puffin_manager_factory.clone(),
201 intermediate_manager.clone(),
202 )
203 .await?;
204 let cache_manager = Arc::new(
205 CacheManager::builder()
206 .sst_meta_cache_size(config.sst_meta_cache_size.as_bytes())
207 .vector_cache_size(config.vector_cache_size.as_bytes())
208 .page_cache_size(config.page_cache_size.as_bytes())
209 .selector_result_cache_size(config.selector_result_cache_size.as_bytes())
210 .range_result_cache_size(config.range_result_cache_size.as_bytes())
211 .prefilter_result_cache_size(config.prefilter_result_cache_size.as_bytes())
212 .index_metadata_size(config.index.metadata_cache_size.as_bytes())
213 .index_content_size(config.index.content_cache_size.as_bytes())
214 .index_content_page_size(config.index.content_cache_page_size.as_bytes())
215 .index_result_cache_size(config.index.result_cache_size.as_bytes())
216 .puffin_metadata_size(config.index.metadata_cache_size.as_bytes())
217 .write_cache(write_cache)
218 .build(),
219 );
220 let time_provider = Arc::new(StdTimeProvider);
221 let total_memory = get_total_memory_bytes();
222 let total_memory = if total_memory > 0 {
223 total_memory as u64
224 } else {
225 0
226 };
227 let compaction_limit_bytes = config
228 .experimental_compaction_memory_limit
229 .resolve(total_memory);
230 let compaction_memory_manager =
231 Arc::new(new_compaction_memory_manager(compaction_limit_bytes));
232 let gc_limiter = Arc::new(GcLimiter::new(config.gc.max_concurrent_gc_job));
233
234 let workers = (0..config.num_workers)
235 .map(|id| {
236 WorkerStarter {
237 id: id as WorkerId,
238 config: config.clone(),
239 log_store: log_store.clone(),
240 object_store_manager: object_store_manager.clone(),
241 write_buffer_manager: write_buffer_manager.clone(),
242 index_build_job_pool: index_build_job_pool.clone(),
243 flush_job_pool: flush_job_pool.clone(),
244 compact_job_pool: compact_job_pool.clone(),
245 purge_scheduler: purge_scheduler.clone(),
246 listener: WorkerListener::default(),
247 cache_manager: cache_manager.clone(),
248 compaction_memory_manager: compaction_memory_manager.clone(),
249 puffin_manager_factory: puffin_manager_factory.clone(),
250 intermediate_manager: intermediate_manager.clone(),
251 time_provider: time_provider.clone(),
252 flush_sender: flush_sender.clone(),
253 flush_receiver: flush_receiver.clone(),
254 plugins: plugins.clone(),
255 schema_metadata_manager: schema_metadata_manager.clone(),
256 file_ref_manager: file_ref_manager.clone(),
257 partition_expr_fetcher: partition_expr_fetcher.clone(),
258 flush_semaphore: flush_semaphore.clone(),
259 }
260 .start()
261 })
262 .collect::<Result<Vec<_>>>()?;
263
264 Ok(WorkerGroup {
265 workers,
266 flush_job_pool,
267 compact_job_pool,
268 index_build_job_pool,
269 purge_scheduler,
270 cache_manager,
271 file_ref_manager,
272 gc_limiter,
273 object_store_manager,
274 puffin_manager_factory,
275 intermediate_manager,
276 schema_metadata_manager,
277 })
278 }
279
280 pub(crate) async fn stop(&self) -> Result<()> {
282 info!("Stop region worker group");
283
284 self.compact_job_pool.stop(true).await?;
287 self.flush_job_pool.stop(true).await?;
289 self.purge_scheduler.stop(true).await?;
291 self.index_build_job_pool.stop(true).await?;
293
294 try_join_all(self.workers.iter().map(|worker| worker.stop())).await?;
295
296 Ok(())
297 }
298
299 pub(crate) async fn submit_to_worker(
301 &self,
302 region_id: RegionId,
303 request: WorkerRequest,
304 ) -> Result<()> {
305 self.worker(region_id).submit_request(request).await
306 }
307
308 pub(crate) fn is_region_exists(&self, region_id: RegionId) -> bool {
310 self.worker(region_id).is_region_exists(region_id)
311 }
312
313 pub(crate) fn is_region_opening(&self, region_id: RegionId) -> bool {
315 self.worker(region_id).is_region_opening(region_id)
316 }
317
318 pub(crate) fn is_region_catching_up(&self, region_id: RegionId) -> bool {
320 self.worker(region_id).is_region_catching_up(region_id)
321 }
322
323 pub(crate) fn get_region(&self, region_id: RegionId) -> Option<MitoRegionRef> {
327 self.worker(region_id).get_region(region_id)
328 }
329
330 pub(crate) fn cache_manager(&self) -> CacheManagerRef {
332 self.cache_manager.clone()
333 }
334
335 pub(crate) fn file_ref_manager(&self) -> FileReferenceManagerRef {
336 self.file_ref_manager.clone()
337 }
338
339 pub(crate) fn gc_limiter(&self) -> GcLimiterRef {
340 self.gc_limiter.clone()
341 }
342
343 pub(crate) fn worker(&self, region_id: RegionId) -> &RegionWorker {
345 let index = region_id_to_index(region_id, self.workers.len());
346
347 &self.workers[index]
348 }
349
350 pub(crate) fn all_regions(&self) -> impl Iterator<Item = MitoRegionRef> + use<'_> {
351 self.workers
352 .iter()
353 .flat_map(|worker| worker.regions.list_regions())
354 }
355
356 pub(crate) fn object_store_manager(&self) -> &ObjectStoreManagerRef {
357 &self.object_store_manager
358 }
359
360 pub(crate) fn puffin_manager_factory(&self) -> &PuffinManagerFactory {
361 &self.puffin_manager_factory
362 }
363
364 pub(crate) fn intermediate_manager(&self) -> &IntermediateManager {
365 &self.intermediate_manager
366 }
367
368 pub(crate) fn schema_metadata_manager(&self) -> &SchemaMetadataManagerRef {
369 &self.schema_metadata_manager
370 }
371}
372
373#[cfg(any(test, feature = "test"))]
375impl WorkerGroup {
376 #[allow(clippy::too_many_arguments)]
380 pub(crate) async fn start_for_test<S: LogStore>(
381 config: Arc<MitoConfig>,
382 log_store: Arc<S>,
383 object_store_manager: ObjectStoreManagerRef,
384 write_buffer_manager: Option<WriteBufferManagerRef>,
385 listener: Option<crate::engine::listener::EventListenerRef>,
386 schema_metadata_manager: SchemaMetadataManagerRef,
387 file_ref_manager: FileReferenceManagerRef,
388 time_provider: TimeProviderRef,
389 partition_expr_fetcher: PartitionExprFetcherRef,
390 ) -> Result<WorkerGroup> {
391 let (flush_sender, flush_receiver) = watch::channel(());
392 let write_buffer_manager = write_buffer_manager.unwrap_or_else(|| {
393 Arc::new(
394 WriteBufferManagerImpl::new(config.global_write_buffer_size.as_bytes() as usize)
395 .with_notifier(flush_sender.clone()),
396 )
397 });
398 let index_build_job_pool =
399 Arc::new(LocalScheduler::new(config.max_background_index_builds));
400 let flush_job_pool = Arc::new(LocalScheduler::new(config.max_background_flushes));
401 let compact_job_pool = Arc::new(LocalScheduler::new(config.max_background_compactions));
402 let flush_semaphore = Arc::new(Semaphore::new(config.max_background_flushes));
403 let purge_scheduler = Arc::new(LocalScheduler::new(config.max_background_flushes));
404 let puffin_manager_factory = PuffinManagerFactory::new(
405 &config.index.aux_path,
406 config.index.staging_size.as_bytes(),
407 Some(config.index.write_buffer_size.as_bytes() as _),
408 config.index.staging_ttl,
409 )
410 .await?;
411 let intermediate_manager = IntermediateManager::init_fs(&config.index.aux_path)
412 .await?
413 .with_buffer_size(Some(config.index.write_buffer_size.as_bytes() as _));
414 let write_cache = write_cache_from_config(
415 &config,
416 puffin_manager_factory.clone(),
417 intermediate_manager.clone(),
418 )
419 .await?;
420 let cache_manager = Arc::new(
421 CacheManager::builder()
422 .sst_meta_cache_size(config.sst_meta_cache_size.as_bytes())
423 .vector_cache_size(config.vector_cache_size.as_bytes())
424 .page_cache_size(config.page_cache_size.as_bytes())
425 .selector_result_cache_size(config.selector_result_cache_size.as_bytes())
426 .range_result_cache_size(config.range_result_cache_size.as_bytes())
427 .prefilter_result_cache_size(config.prefilter_result_cache_size.as_bytes())
428 .write_cache(write_cache)
429 .build(),
430 );
431 let total_memory = get_total_memory_bytes();
432 let total_memory = if total_memory > 0 {
433 total_memory as u64
434 } else {
435 0
436 };
437 let compaction_limit_bytes = config
438 .experimental_compaction_memory_limit
439 .resolve(total_memory);
440 let compaction_memory_manager =
441 Arc::new(new_compaction_memory_manager(compaction_limit_bytes));
442 let gc_limiter = Arc::new(GcLimiter::new(config.gc.max_concurrent_gc_job));
443 let workers = (0..config.num_workers)
444 .map(|id| {
445 WorkerStarter {
446 id: id as WorkerId,
447 config: config.clone(),
448 log_store: log_store.clone(),
449 object_store_manager: object_store_manager.clone(),
450 write_buffer_manager: write_buffer_manager.clone(),
451 index_build_job_pool: index_build_job_pool.clone(),
452 flush_job_pool: flush_job_pool.clone(),
453 compact_job_pool: compact_job_pool.clone(),
454 purge_scheduler: purge_scheduler.clone(),
455 listener: WorkerListener::new(listener.clone()),
456 cache_manager: cache_manager.clone(),
457 compaction_memory_manager: compaction_memory_manager.clone(),
458 puffin_manager_factory: puffin_manager_factory.clone(),
459 intermediate_manager: intermediate_manager.clone(),
460 time_provider: time_provider.clone(),
461 flush_sender: flush_sender.clone(),
462 flush_receiver: flush_receiver.clone(),
463 plugins: Plugins::new(),
464 schema_metadata_manager: schema_metadata_manager.clone(),
465 file_ref_manager: file_ref_manager.clone(),
466 partition_expr_fetcher: partition_expr_fetcher.clone(),
467 flush_semaphore: flush_semaphore.clone(),
468 }
469 .start()
470 })
471 .collect::<Result<Vec<_>>>()?;
472
473 Ok(WorkerGroup {
474 workers,
475 flush_job_pool,
476 compact_job_pool,
477 index_build_job_pool,
478 purge_scheduler,
479 cache_manager,
480 file_ref_manager,
481 gc_limiter,
482 object_store_manager,
483 puffin_manager_factory,
484 intermediate_manager,
485 schema_metadata_manager,
486 })
487 }
488
489 pub(crate) fn purge_scheduler(&self) -> &SchedulerRef {
491 &self.purge_scheduler
492 }
493}
494
495fn region_id_to_index(id: RegionId, num_workers: usize) -> usize {
496 ((id.table_id() as usize % num_workers) + (id.region_number() as usize % num_workers))
497 % num_workers
498}
499
500pub async fn write_cache_from_config(
501 config: &MitoConfig,
502 puffin_manager_factory: PuffinManagerFactory,
503 intermediate_manager: IntermediateManager,
504) -> Result<Option<WriteCacheRef>> {
505 if !config.enable_write_cache {
506 return Ok(None);
507 }
508
509 tokio::fs::create_dir_all(Path::new(&config.write_cache_path))
510 .await
511 .context(CreateDirSnafu {
512 dir: &config.write_cache_path,
513 })?;
514
515 let cache = WriteCache::new_fs(
516 &config.write_cache_path,
517 config.write_cache_size,
518 config.write_cache_ttl,
519 Some(config.index_cache_percent),
520 config.enable_refill_cache_on_read,
521 puffin_manager_factory,
522 intermediate_manager,
523 config.manifest_cache_size,
524 )
525 .await?;
526 Ok(Some(Arc::new(cache)))
527}
528
529pub(crate) fn worker_init_check_delay() -> Duration {
531 let init_check_delay = rng().random_range(0..MAX_INITIAL_CHECK_DELAY_SECS);
532 Duration::from_secs(init_check_delay)
533}
534
535struct WorkerStarter<S> {
537 id: WorkerId,
538 config: Arc<MitoConfig>,
539 log_store: Arc<S>,
540 object_store_manager: ObjectStoreManagerRef,
541 write_buffer_manager: WriteBufferManagerRef,
542 compact_job_pool: SchedulerRef,
543 index_build_job_pool: SchedulerRef,
544 flush_job_pool: SchedulerRef,
545 purge_scheduler: SchedulerRef,
546 listener: WorkerListener,
547 cache_manager: CacheManagerRef,
548 compaction_memory_manager: Arc<CompactionMemoryManager>,
549 puffin_manager_factory: PuffinManagerFactory,
550 intermediate_manager: IntermediateManager,
551 time_provider: TimeProviderRef,
552 flush_sender: watch::Sender<()>,
554 flush_receiver: watch::Receiver<()>,
556 plugins: Plugins,
557 schema_metadata_manager: SchemaMetadataManagerRef,
558 file_ref_manager: FileReferenceManagerRef,
559 partition_expr_fetcher: PartitionExprFetcherRef,
560 flush_semaphore: Arc<Semaphore>,
561}
562
563impl<S: LogStore> WorkerStarter<S> {
564 fn start(self) -> Result<RegionWorker> {
566 let regions = Arc::new(RegionMap::default());
567 let opening_regions = Arc::new(OpeningRegions::default());
568 let catchup_regions = Arc::new(CatchupRegions::default());
569 let (sender, receiver) = mpsc::channel(self.config.worker_channel_size);
570
571 let running = Arc::new(AtomicBool::new(true));
572 let now = self.time_provider.current_time_millis();
573 let id_string = self.id.to_string();
574 let mut worker_thread = RegionWorkerLoop {
575 id: self.id,
576 config: self.config.clone(),
577 regions: regions.clone(),
578 catchup_regions: catchup_regions.clone(),
579 dropping_regions: Arc::new(RegionMap::default()),
580 opening_regions: opening_regions.clone(),
581 sender: sender.clone(),
582 receiver,
583 wal: Wal::new(self.log_store),
584 object_store_manager: self.object_store_manager.clone(),
585 running: running.clone(),
586 memtable_builder_provider: MemtableBuilderProvider::new(
587 Some(self.write_buffer_manager.clone()),
588 self.config.clone(),
589 ),
590 purge_scheduler: self.purge_scheduler.clone(),
591 write_buffer_manager: self.write_buffer_manager,
592 index_build_scheduler: IndexBuildScheduler::new(
593 self.index_build_job_pool,
594 self.config.max_background_index_builds,
595 ),
596 flush_scheduler: FlushScheduler::new(self.flush_job_pool),
597 compaction_scheduler: CompactionScheduler::new(
598 self.compact_job_pool,
599 sender.clone(),
600 self.cache_manager.clone(),
601 self.config.clone(),
602 self.listener.clone(),
603 self.plugins.clone(),
604 self.compaction_memory_manager.clone(),
605 self.config.experimental_compaction_on_exhausted,
606 ),
607 stalled_requests: StalledRequests::default(),
608 listener: self.listener,
609 cache_manager: self.cache_manager,
610 puffin_manager_factory: self.puffin_manager_factory,
611 intermediate_manager: self.intermediate_manager,
612 time_provider: self.time_provider,
613 last_periodical_check_millis: now,
614 flush_sender: self.flush_sender,
615 flush_receiver: self.flush_receiver,
616 stalling_count: WRITE_STALLING.with_label_values(&[&id_string]),
617 region_count: REGION_COUNT.with_label_values(&[&id_string]),
618 request_wait_time: REQUEST_WAIT_TIME.with_label_values(&[&id_string]),
619 region_edit_queues: RegionEditQueues::default(),
620 schema_metadata_manager: self.schema_metadata_manager,
621 file_ref_manager: self.file_ref_manager.clone(),
622 partition_expr_fetcher: self.partition_expr_fetcher,
623 flush_semaphore: self.flush_semaphore,
624 plugins: self.plugins,
625 };
626 let handle = common_runtime::spawn_global(async move {
627 worker_thread.run().await;
628 });
629
630 Ok(RegionWorker {
631 id: self.id,
632 regions,
633 opening_regions,
634 catchup_regions,
635 sender,
636 handle: Mutex::new(Some(handle)),
637 running,
638 })
639 }
640}
641
642pub(crate) struct RegionWorker {
644 id: WorkerId,
646 regions: RegionMapRef,
648 opening_regions: OpeningRegionsRef,
650 catchup_regions: CatchupRegionsRef,
652 sender: Sender<WorkerRequestWithTime>,
654 handle: Mutex<Option<JoinHandle<()>>>,
656 running: Arc<AtomicBool>,
658}
659
660impl RegionWorker {
661 async fn submit_request(&self, request: WorkerRequest) -> Result<()> {
663 ensure!(self.is_running(), WorkerStoppedSnafu { id: self.id });
664 let request_with_time = WorkerRequestWithTime::new(request);
665 if self.sender.send(request_with_time).await.is_err() {
666 warn!(
667 "Worker {} is already exited but the running flag is still true",
668 self.id
669 );
670 self.set_running(false);
672 return WorkerStoppedSnafu { id: self.id }.fail();
673 }
674
675 Ok(())
676 }
677
678 async fn stop(&self) -> Result<()> {
682 let handle = self.handle.lock().await.take();
683 if let Some(handle) = handle {
684 info!("Stop region worker {}", self.id);
685
686 self.set_running(false);
687 if self
688 .sender
689 .send(WorkerRequestWithTime::new(WorkerRequest::Stop))
690 .await
691 .is_err()
692 {
693 warn!("Worker {} is already exited before stop", self.id);
694 }
695
696 handle.await.context(JoinSnafu)?;
697 }
698
699 Ok(())
700 }
701
702 fn is_running(&self) -> bool {
704 self.running.load(Ordering::Relaxed)
705 }
706
707 fn set_running(&self, value: bool) {
709 self.running.store(value, Ordering::Relaxed)
710 }
711
712 fn is_region_exists(&self, region_id: RegionId) -> bool {
714 self.regions.is_region_exists(region_id)
715 }
716
717 fn is_region_opening(&self, region_id: RegionId) -> bool {
719 self.opening_regions.is_region_exists(region_id)
720 }
721
722 fn is_region_catching_up(&self, region_id: RegionId) -> bool {
724 self.catchup_regions.is_region_exists(region_id)
725 }
726
727 fn get_region(&self, region_id: RegionId) -> Option<MitoRegionRef> {
729 self.regions.get_region(region_id)
730 }
731
732 #[cfg(test)]
733 pub(crate) fn opening_regions(&self) -> &OpeningRegionsRef {
735 &self.opening_regions
736 }
737
738 #[cfg(test)]
739 pub(crate) fn catchup_regions(&self) -> &CatchupRegionsRef {
741 &self.catchup_regions
742 }
743}
744
745impl Drop for RegionWorker {
746 fn drop(&mut self) {
747 if self.is_running() {
748 self.set_running(false);
749 }
751 }
752}
753
754type RequestBuffer = Vec<WorkerRequest>;
755
756#[derive(Default)]
760pub(crate) struct StalledRequests {
761 pub(crate) requests:
768 HashMap<RegionId, (usize, Vec<SenderWriteRequest>, Vec<SenderBulkRequest>)>,
769 pub(crate) estimated_size: usize,
771}
772
773impl StalledRequests {
774 pub(crate) fn estimated_size(&self, region_id: &RegionId) -> usize {
776 self.requests
777 .get(region_id)
778 .map(|(size, _, _)| *size)
779 .unwrap_or_default()
780 }
781
782 pub(crate) fn append(
784 &mut self,
785 requests: &mut Vec<SenderWriteRequest>,
786 bulk_requests: &mut Vec<SenderBulkRequest>,
787 ) {
788 for req in requests.drain(..) {
789 self.push(req);
790 }
791 for req in bulk_requests.drain(..) {
792 self.push_bulk(req);
793 }
794 }
795
796 pub(crate) fn push(&mut self, req: SenderWriteRequest) {
798 let (size, requests, _) = self.requests.entry(req.request.region_id).or_default();
799 let req_size = req.request.estimated_size();
800 *size += req_size;
801 self.estimated_size += req_size;
802 requests.push(req);
803 }
804
805 pub(crate) fn push_bulk(&mut self, req: SenderBulkRequest) {
806 let region_id = req.region_id;
807 let (size, _, requests) = self.requests.entry(region_id).or_default();
808 let req_size = req.request.estimated_size();
809 *size += req_size;
810 self.estimated_size += req_size;
811 requests.push(req);
812 }
813
814 pub(crate) fn remove(
816 &mut self,
817 region_id: &RegionId,
818 ) -> (Vec<SenderWriteRequest>, Vec<SenderBulkRequest>) {
819 if let Some((size, write_reqs, bulk_reqs)) = self.requests.remove(region_id) {
820 self.estimated_size -= size;
821 (write_reqs, bulk_reqs)
822 } else {
823 (vec![], vec![])
824 }
825 }
826
827 pub(crate) fn stalled_count(&self) -> usize {
829 self.requests
830 .values()
831 .map(|(_, reqs, bulk_reqs)| reqs.len() + bulk_reqs.len())
832 .sum()
833 }
834}
835
836struct RegionWorkerLoop<S> {
838 id: WorkerId,
840 config: Arc<MitoConfig>,
842 regions: RegionMapRef,
844 dropping_regions: RegionMapRef,
846 opening_regions: OpeningRegionsRef,
848 catchup_regions: CatchupRegionsRef,
850 sender: Sender<WorkerRequestWithTime>,
852 receiver: Receiver<WorkerRequestWithTime>,
854 wal: Wal<S>,
856 object_store_manager: ObjectStoreManagerRef,
858 running: Arc<AtomicBool>,
860 memtable_builder_provider: MemtableBuilderProvider,
862 purge_scheduler: SchedulerRef,
864 write_buffer_manager: WriteBufferManagerRef,
866 index_build_scheduler: IndexBuildScheduler,
868 flush_scheduler: FlushScheduler,
870 compaction_scheduler: CompactionScheduler,
872 stalled_requests: StalledRequests,
874 listener: WorkerListener,
876 cache_manager: CacheManagerRef,
878 puffin_manager_factory: PuffinManagerFactory,
880 intermediate_manager: IntermediateManager,
882 time_provider: TimeProviderRef,
884 last_periodical_check_millis: i64,
886 flush_sender: watch::Sender<()>,
888 flush_receiver: watch::Receiver<()>,
890 stalling_count: IntGauge,
892 region_count: IntGauge,
894 request_wait_time: Histogram,
896 region_edit_queues: RegionEditQueues,
898 schema_metadata_manager: SchemaMetadataManagerRef,
900 file_ref_manager: FileReferenceManagerRef,
902 partition_expr_fetcher: PartitionExprFetcherRef,
904 flush_semaphore: Arc<Semaphore>,
906 plugins: Plugins,
908}
909
910impl<S: LogStore> RegionWorkerLoop<S> {
911 async fn run(&mut self) {
913 let init_check_delay = worker_init_check_delay();
914 info!(
915 "Start region worker thread {}, init_check_delay: {:?}",
916 self.id, init_check_delay
917 );
918 self.last_periodical_check_millis += init_check_delay.as_millis() as i64;
919
920 let mut write_req_buffer: Vec<SenderWriteRequest> =
922 Vec::with_capacity(self.config.worker_request_batch_size);
923 let mut bulk_req_buffer: Vec<SenderBulkRequest> =
924 Vec::with_capacity(self.config.worker_request_batch_size);
925 let mut ddl_req_buffer: Vec<SenderDdlRequest> =
926 Vec::with_capacity(self.config.worker_request_batch_size);
927 let mut general_req_buffer: Vec<WorkerRequest> =
928 RequestBuffer::with_capacity(self.config.worker_request_batch_size);
929
930 while self.running.load(Ordering::Relaxed) {
931 write_req_buffer.clear();
933 ddl_req_buffer.clear();
934 general_req_buffer.clear();
935 let mut bulk_insert_req_num = 0;
936
937 let max_wait_time = self.time_provider.wait_duration(CHECK_REGION_INTERVAL);
938 let sleep = tokio::time::sleep(max_wait_time);
939 tokio::pin!(sleep);
940
941 tokio::select! {
942 request_opt = self.receiver.recv() => {
943 match request_opt {
944 Some(request_with_time) => {
945 let wait_time = request_with_time.created_at.elapsed();
947 self.request_wait_time.observe(wait_time.as_secs_f64());
948
949 match request_with_time.request {
950 WorkerRequest::Write(sender_req) => write_req_buffer.push(sender_req),
951 WorkerRequest::Ddl(sender_req) => ddl_req_buffer.push(sender_req),
952 WorkerRequest::BulkInserts(bulk_insert) => {
953 bulk_insert_req_num += 1;
954 self.buffer_bulk_insert_request(bulk_insert, &mut bulk_req_buffer)
955 .await;
956 }
957 req => general_req_buffer.push(req),
958 }
959 },
960 None => break,
962 }
963 }
964 recv_res = self.flush_receiver.changed() => {
965 if recv_res.is_err() {
966 break;
968 } else {
969 self.maybe_flush_worker();
974 self.handle_stalled_requests().await;
976 continue;
977 }
978 }
979 _ = &mut sleep => {
980 self.handle_periodical_tasks();
982 continue;
983 }
984 }
985
986 if self.flush_receiver.has_changed().unwrap_or(false) {
987 self.handle_stalled_requests().await;
991 }
992
993 for _ in 1..self.config.worker_request_batch_size {
995 match self.receiver.try_recv() {
997 Ok(request_with_time) => {
998 let wait_time = request_with_time.created_at.elapsed();
1000 self.request_wait_time.observe(wait_time.as_secs_f64());
1001
1002 match request_with_time.request {
1003 WorkerRequest::Write(sender_req) => write_req_buffer.push(sender_req),
1004 WorkerRequest::Ddl(sender_req) => ddl_req_buffer.push(sender_req),
1005 WorkerRequest::BulkInserts(bulk_insert) => {
1006 bulk_insert_req_num += 1;
1007 self.buffer_bulk_insert_request(bulk_insert, &mut bulk_req_buffer)
1008 .await
1009 }
1010 req => general_req_buffer.push(req),
1011 }
1012 }
1013 Err(_) => break,
1015 }
1016 }
1017
1018 self.listener.on_recv_requests(
1019 write_req_buffer.len()
1020 + ddl_req_buffer.len()
1021 + general_req_buffer.len()
1022 + bulk_insert_req_num,
1023 );
1024
1025 self.handle_requests(
1026 &mut write_req_buffer,
1027 &mut ddl_req_buffer,
1028 &mut general_req_buffer,
1029 &mut bulk_req_buffer,
1030 )
1031 .await;
1032
1033 self.handle_periodical_tasks();
1034 }
1035
1036 self.clean().await;
1037
1038 info!("Exit region worker thread {}", self.id);
1039 }
1040
1041 async fn buffer_bulk_insert_request(
1042 &mut self,
1043 bulk_insert: BulkInsertRequest,
1044 bulk_requests: &mut Vec<SenderBulkRequest>,
1045 ) {
1046 let BulkInsertRequest {
1047 metadata,
1048 request,
1049 sender,
1050 } = bulk_insert;
1051
1052 if let Some(region_metadata) = metadata {
1053 self.handle_bulk_insert_batch(region_metadata, request, bulk_requests, sender)
1054 .await;
1055 } else {
1056 error!("Cannot find region metadata for {}", request.region_id);
1057 sender.send(
1058 error::RegionNotFoundSnafu {
1059 region_id: request.region_id,
1060 }
1061 .fail(),
1062 );
1063 }
1064 }
1065
1066 async fn handle_requests(
1070 &mut self,
1071 write_requests: &mut Vec<SenderWriteRequest>,
1072 ddl_requests: &mut Vec<SenderDdlRequest>,
1073 general_requests: &mut Vec<WorkerRequest>,
1074 bulk_requests: &mut Vec<SenderBulkRequest>,
1075 ) {
1076 for worker_req in general_requests.drain(..) {
1077 match worker_req {
1078 WorkerRequest::Write(_) | WorkerRequest::Ddl(_) => {
1079 continue;
1081 }
1082 WorkerRequest::BulkInserts(_) => unreachable!("bulk inserts are buffered"),
1083 WorkerRequest::Background { region_id, notify } => {
1084 if matches!(
1085 ¬ify,
1086 BackgroundNotify::RegionEdit(edit_result)
1087 if edit_result.update_region_state
1088 ) {
1089 self.handle_buffered_region_write_requests(
1098 ®ion_id,
1099 write_requests,
1100 bulk_requests,
1101 )
1102 .await;
1103 }
1104 self.handle_background_notify(region_id, notify).await;
1106 }
1107 WorkerRequest::SetRegionRoleStateGracefully {
1108 region_id,
1109 region_role_state,
1110 sender,
1111 } => {
1112 self.set_role_state_gracefully(region_id, region_role_state, sender)
1113 .await;
1114 }
1115 WorkerRequest::EditRegion(request) => {
1116 self.handle_region_edit(request);
1117 }
1118 WorkerRequest::Stop => {
1119 debug_assert!(!self.running.load(Ordering::Relaxed));
1120 }
1121 WorkerRequest::SyncRegion(req) => {
1122 self.handle_region_sync(req).await;
1123 }
1124 WorkerRequest::RemapManifests(req) => {
1125 self.handle_remap_manifests_request(req);
1126 }
1127 WorkerRequest::CopyRegionFrom(req) => {
1128 self.handle_copy_region_from_request(req);
1129 }
1130 }
1131 }
1132
1133 self.handle_write_requests(write_requests, bulk_requests, true)
1136 .await;
1137
1138 self.handle_ddl_requests(ddl_requests).await;
1139 }
1140
1141 async fn handle_ddl_requests(&mut self, ddl_requests: &mut Vec<SenderDdlRequest>) {
1143 if ddl_requests.is_empty() {
1144 return;
1145 }
1146
1147 for ddl in ddl_requests.drain(..) {
1148 let res = match ddl.request {
1149 DdlRequest::Create(req) => self.handle_create_request(ddl.region_id, req).await,
1150 DdlRequest::Drop(req) => {
1151 self.handle_drop_request(ddl.region_id, req, ddl.sender)
1152 .await;
1153 continue;
1154 }
1155 DdlRequest::Open((req, wal_entry_receiver)) => {
1156 self.handle_open_request(ddl.region_id, req, wal_entry_receiver, ddl.sender)
1157 .await;
1158 continue;
1159 }
1160 DdlRequest::OfflineCleanup(req) => {
1161 self.handle_offline_cleanup_request(ddl.region_id, req)
1162 .await
1163 }
1164 DdlRequest::Close(req) => {
1165 self.handle_close_request(ddl.region_id, req, ddl.sender)
1166 .await;
1167 continue;
1168 }
1169 DdlRequest::Alter(req) => {
1170 self.handle_alter_request(ddl.region_id, req, ddl.sender)
1171 .await;
1172 continue;
1173 }
1174 DdlRequest::Flush(req) => {
1175 self.handle_flush_request(ddl.region_id, req, ddl.sender);
1176 continue;
1177 }
1178 DdlRequest::Compact(req) => {
1179 self.handle_compaction_request(ddl.region_id, req, ddl.sender)
1180 .await;
1181 continue;
1182 }
1183 DdlRequest::BuildIndex(req) => {
1184 self.handle_build_index_request(ddl.region_id, req, ddl.sender)
1185 .await;
1186 continue;
1187 }
1188 DdlRequest::Truncate(req) => {
1189 self.handle_truncate_request(ddl.region_id, req, ddl.sender)
1190 .await;
1191 continue;
1192 }
1193 DdlRequest::Catchup((req, wal_entry_receiver)) => {
1194 self.handle_catchup_request(ddl.region_id, req, wal_entry_receiver, ddl.sender)
1195 .await;
1196 continue;
1197 }
1198 DdlRequest::EnterStaging(req) => {
1199 self.handle_enter_staging_request(
1200 ddl.region_id,
1201 req.partition_directive,
1202 ddl.sender,
1203 )
1204 .await;
1205 continue;
1206 }
1207 DdlRequest::ApplyStagingManifest(req) => {
1208 self.handle_apply_staging_manifest_request(ddl.region_id, req, ddl.sender)
1209 .await;
1210 continue;
1211 }
1212 };
1213
1214 ddl.sender.send(res);
1215 }
1216 }
1217
1218 fn handle_periodical_tasks(&mut self) {
1220 let interval = CHECK_REGION_INTERVAL.as_millis() as i64;
1221 if self
1222 .time_provider
1223 .elapsed_since(self.last_periodical_check_millis)
1224 < interval
1225 {
1226 return;
1227 }
1228
1229 self.last_periodical_check_millis = self.time_provider.current_time_millis();
1230
1231 if let Err(e) = self.flush_periodically() {
1232 error!(e; "Failed to flush regions periodically");
1233 }
1234 }
1235
1236 async fn handle_background_notify(&mut self, region_id: RegionId, notify: BackgroundNotify) {
1238 match notify {
1239 BackgroundNotify::CompactionPickFinished(req) => {
1240 self.handle_compaction_pick_finished(region_id, req).await
1241 }
1242 BackgroundNotify::FlushFinished(req) => {
1243 self.handle_flush_finished(region_id, req).await
1244 }
1245 BackgroundNotify::FlushFailed(req) => self.handle_flush_failed(region_id, req).await,
1246 BackgroundNotify::IndexBuildFinished(req) => {
1247 self.handle_index_build_finished(region_id, req).await
1248 }
1249 BackgroundNotify::IndexBuildStopped(req) => {
1250 self.handle_index_build_stopped(region_id, req).await
1251 }
1252 BackgroundNotify::IndexBuildFailed(req) => {
1253 self.handle_index_build_failed(region_id, req).await
1254 }
1255 BackgroundNotify::IndexBuildRetry(req) => {
1256 self.handle_rebuild_index(req, OptionOutputTx::new(None))
1257 .await
1258 }
1259 BackgroundNotify::CompactionFinished(req) => {
1260 self.handle_compaction_finished(region_id, req).await
1261 }
1262 BackgroundNotify::CompactionCancelled(req) => {
1263 self.handle_compaction_cancelled(region_id, req).await
1264 }
1265 BackgroundNotify::CompactionFailed(req) => self.handle_compaction_failure(req).await,
1266 BackgroundNotify::Truncate(req) => self.handle_truncate_result(req).await,
1267 BackgroundNotify::DiscardUnflushed(req) => {
1268 self.handle_discard_unflushed_result(req).await
1269 }
1270 BackgroundNotify::RegionChange(req) => {
1271 self.handle_manifest_region_change_result(req).await
1272 }
1273 BackgroundNotify::EnterStaging(req) => self.handle_enter_staging_result(req).await,
1274 BackgroundNotify::RegionEdit(req) => self.handle_region_edit_result(req).await,
1275 BackgroundNotify::CopyRegionFromFinished(req) => {
1276 self.handle_copy_region_from_finished(req)
1277 }
1278 }
1279 }
1280
1281 async fn set_role_state_gracefully(
1283 &mut self,
1284 region_id: RegionId,
1285 region_role_state: SettableRegionRoleState,
1286 sender: oneshot::Sender<SetRegionRoleStateResponse>,
1287 ) {
1288 if let Some(region) = self.regions.get_region(region_id) {
1289 common_runtime::spawn_global(async move {
1291 match region.set_role_state_gracefully(region_role_state).await {
1292 Ok(()) => {
1293 let last_entry_id = region.version_control.current().last_entry_id;
1294 let _ = sender.send(SetRegionRoleStateResponse::success(
1295 SetRegionRoleStateSuccess::mito(last_entry_id),
1296 ));
1297 }
1298 Err(e) => {
1299 error!(e; "Failed to set region {} role state to {:?}", region_id, region_role_state);
1300 let _ = sender.send(SetRegionRoleStateResponse::invalid_transition(
1301 BoxedError::new(e),
1302 ));
1303 }
1304 }
1305 });
1306 } else {
1307 let _ = sender.send(SetRegionRoleStateResponse::NotFound);
1308 }
1309 }
1310}
1311
1312impl<S> RegionWorkerLoop<S> {
1313 async fn clean(&self) {
1315 let regions = self.regions.list_regions();
1317 for region in regions {
1318 region.stop().await;
1319 }
1320
1321 self.regions.clear();
1322 }
1323
1324 fn notify_group(&mut self) {
1327 let _ = self.flush_sender.send(());
1329 self.flush_receiver.borrow_and_update();
1331 }
1332}
1333
1334#[derive(Default, Clone)]
1336pub(crate) struct WorkerListener {
1337 #[cfg(any(test, feature = "test"))]
1338 listener: Option<crate::engine::listener::EventListenerRef>,
1339}
1340
1341impl WorkerListener {
1342 #[cfg(any(test, feature = "test"))]
1343 pub(crate) fn new(
1344 listener: Option<crate::engine::listener::EventListenerRef>,
1345 ) -> WorkerListener {
1346 WorkerListener { listener }
1347 }
1348
1349 pub(crate) fn on_flush_success(&self, region_id: RegionId) {
1351 #[cfg(any(test, feature = "test"))]
1352 if let Some(listener) = &self.listener {
1353 listener.on_flush_success(region_id);
1354 }
1355 let _ = region_id;
1357 }
1358
1359 pub(crate) fn on_write_stall(&self) {
1361 #[cfg(any(test, feature = "test"))]
1362 if let Some(listener) = &self.listener {
1363 listener.on_write_stall();
1364 }
1365 }
1366
1367 pub(crate) async fn on_flush_begin(&self, region_id: RegionId) {
1368 #[cfg(any(test, feature = "test"))]
1369 if let Some(listener) = &self.listener {
1370 listener.on_flush_begin(region_id).await;
1371 }
1372 let _ = region_id;
1374 }
1375
1376 pub(crate) async fn on_flush_commit_begin(&self, _region_id: RegionId) {
1377 #[cfg(any(test, feature = "test"))]
1378 if let Some(listener) = &self.listener {
1379 listener.on_flush_commit_begin(_region_id).await;
1380 }
1381 }
1382
1383 pub(crate) fn on_flush_cancel_requested(&self, _region_id: RegionId) {
1384 #[cfg(any(test, feature = "test"))]
1385 if let Some(listener) = &self.listener {
1386 listener.on_flush_cancel_requested(_region_id);
1387 }
1388 }
1389
1390 pub(crate) fn on_later_drop_begin(&self, region_id: RegionId) -> Option<Duration> {
1391 #[cfg(any(test, feature = "test"))]
1392 if let Some(listener) = &self.listener {
1393 return listener.on_later_drop_begin(region_id);
1394 }
1395 let _ = region_id;
1397 None
1398 }
1399
1400 pub(crate) fn on_later_drop_end(&self, region_id: RegionId, removed: bool) {
1402 #[cfg(any(test, feature = "test"))]
1403 if let Some(listener) = &self.listener {
1404 listener.on_later_drop_end(region_id, removed);
1405 }
1406 let _ = region_id;
1408 let _ = removed;
1409 }
1410
1411 pub(crate) async fn on_merge_ssts_finished(&self, region_id: RegionId) {
1412 #[cfg(any(test, feature = "test"))]
1413 if let Some(listener) = &self.listener {
1414 listener.on_merge_ssts_finished(region_id).await;
1415 }
1416 let _ = region_id;
1418 }
1419
1420 pub(crate) fn on_recv_requests(&self, request_num: usize) {
1421 #[cfg(any(test, feature = "test"))]
1422 if let Some(listener) = &self.listener {
1423 listener.on_recv_requests(request_num);
1424 }
1425 let _ = request_num;
1427 }
1428
1429 pub(crate) fn on_file_cache_filled(&self, _file_id: FileId) {
1430 #[cfg(any(test, feature = "test"))]
1431 if let Some(listener) = &self.listener {
1432 listener.on_file_cache_filled(_file_id);
1433 }
1434 }
1435
1436 pub(crate) fn on_compaction_scheduled(&self, _region_id: RegionId) {
1437 #[cfg(any(test, feature = "test"))]
1438 if let Some(listener) = &self.listener {
1439 listener.on_compaction_scheduled(_region_id);
1440 }
1441 }
1442
1443 pub(crate) async fn on_compaction_pick_begin(&self, _region_id: RegionId) {
1444 #[cfg(any(test, feature = "test"))]
1445 if let Some(listener) = &self.listener {
1446 listener.on_compaction_pick_begin(_region_id).await;
1447 }
1448 }
1449
1450 pub(crate) async fn on_compaction_commit_begin(&self, _region_id: RegionId) {
1451 #[cfg(any(test, feature = "test"))]
1452 if let Some(listener) = &self.listener {
1453 listener.on_compaction_commit_begin(_region_id).await;
1454 }
1455 }
1456
1457 pub(crate) async fn on_compaction_result_notified(&self, _region_id: RegionId) {
1458 #[cfg(any(test, feature = "test"))]
1459 if let Some(listener) = &self.listener {
1460 listener.on_compaction_result_notified(_region_id).await;
1461 }
1462 }
1463
1464 pub(crate) fn on_compaction_cancel_requested(&self, _region_id: RegionId) {
1465 #[cfg(any(test, feature = "test"))]
1466 if let Some(listener) = &self.listener {
1467 listener.on_compaction_cancel_requested(_region_id);
1468 }
1469 }
1470
1471 pub(crate) async fn on_notify_region_change_result_begin(&self, _region_id: RegionId) {
1472 #[cfg(any(test, feature = "test"))]
1473 if let Some(listener) = &self.listener {
1474 listener
1475 .on_notify_region_change_result_begin(_region_id)
1476 .await;
1477 }
1478 }
1479
1480 pub(crate) async fn on_enter_staging_result_begin(&self, _region_id: RegionId) {
1481 #[cfg(any(test, feature = "test"))]
1482 if let Some(listener) = &self.listener {
1483 listener.on_enter_staging_result_begin(_region_id).await;
1484 }
1485 }
1486
1487 pub(crate) async fn on_index_build_finish(&self, _region_file_id: RegionFileId) {
1488 #[cfg(any(test, feature = "test"))]
1489 if let Some(listener) = &self.listener {
1490 listener.on_index_build_finish(_region_file_id).await;
1491 }
1492 }
1493
1494 pub(crate) async fn on_index_build_begin(&self, _region_file_id: RegionFileId) {
1495 #[cfg(any(test, feature = "test"))]
1496 if let Some(listener) = &self.listener {
1497 listener.on_index_build_begin(_region_file_id).await;
1498 }
1499 }
1500
1501 pub(crate) async fn on_index_build_abort(&self, _region_file_id: RegionFileId) {
1502 #[cfg(any(test, feature = "test"))]
1503 if let Some(listener) = &self.listener {
1504 listener.on_index_build_abort(_region_file_id).await;
1505 }
1506 }
1507
1508 pub(crate) async fn on_index_build_before_manifest_commit(
1509 &self,
1510 _region_file_id: RegionFileId,
1511 ) {
1512 #[cfg(any(test, feature = "test"))]
1513 if let Some(listener) = &self.listener {
1514 listener
1515 .on_index_build_before_manifest_commit(_region_file_id)
1516 .await;
1517 }
1518 }
1519
1520 pub(crate) async fn on_index_build_manifest_committed(&self, _region_file_id: RegionFileId) {
1521 #[cfg(any(test, feature = "test"))]
1522 if let Some(listener) = &self.listener {
1523 listener
1524 .on_index_build_manifest_committed(_region_file_id)
1525 .await;
1526 }
1527 }
1528}
1529
1530#[cfg(test)]
1531mod tests {
1532 use super::*;
1533 use crate::test_util::TestEnv;
1534
1535 #[test]
1536 fn test_region_id_to_index() {
1537 let num_workers = 4;
1538
1539 let region_id = RegionId::new(1, 2);
1540 let index = region_id_to_index(region_id, num_workers);
1541 assert_eq!(index, 3);
1542
1543 let region_id = RegionId::new(2, 3);
1544 let index = region_id_to_index(region_id, num_workers);
1545 assert_eq!(index, 1);
1546 }
1547
1548 #[tokio::test]
1549 async fn test_worker_group_start_stop() {
1550 let env = TestEnv::with_prefix("group-stop").await;
1551 let group = env
1552 .create_worker_group(MitoConfig {
1553 num_workers: 4,
1554 ..Default::default()
1555 })
1556 .await;
1557
1558 group.stop().await.unwrap();
1559 }
1560}