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::ObjectStore;
49use object_store::manager::ObjectStoreManagerRef;
50use prometheus::{Histogram, IntGauge};
51use rand::{Rng, rng};
52use snafu::{ResultExt, ensure};
53use store_api::logstore::LogStore;
54use store_api::region_engine::{
55 SetRegionRoleStateResponse, SetRegionRoleStateSuccess, SettableRegionRoleState,
56};
57use store_api::storage::{FileId, RegionId};
58use tokio::sync::mpsc::{Receiver, Sender};
59use tokio::sync::{Mutex, Semaphore, mpsc, oneshot, watch};
60
61use crate::access_layer::new_fs_cache_store;
62use crate::cache::write_cache::{WriteCache, WriteCacheRef};
63use crate::cache::{CacheManager, CacheManagerRef, WriteCacheUploadStoreWrapperRef};
64use crate::compaction::CompactionScheduler;
65use crate::compaction::memory_manager::{CompactionMemoryManager, new_compaction_memory_manager};
66use crate::config::MitoConfig;
67use crate::error::{self, CreateDirSnafu, JoinSnafu, Result, WorkerStoppedSnafu};
68use crate::flush::{FlushScheduler, WriteBufferManagerImpl, WriteBufferManagerRef};
69use crate::gc::{GcLimiter, GcLimiterRef};
70use crate::memtable::MemtableBuilderProvider;
71use crate::metrics::{REGION_COUNT, REQUEST_WAIT_TIME, WRITE_STALLING};
72use crate::region::opener::PartitionExprFetcherRef;
73use crate::region::{
74 CatchupRegions, CatchupRegionsRef, MitoRegionRef, OpeningRegions, OpeningRegionsRef, RegionMap,
75 RegionMapRef,
76};
77use crate::request::{
78 BackgroundNotify, BulkInsertRequest, DdlRequest, OptionOutputTx, SenderBulkRequest,
79 SenderDdlRequest, SenderWriteRequest, WorkerRequest, WorkerRequestWithTime,
80};
81use crate::schedule::scheduler::{LocalScheduler, SchedulerRef};
82use crate::series_index::{
83 IndexFilePurger, SeriesIndexTaskState, series_index_channel, spawn_series_index_tasks,
84};
85use crate::sst::file::RegionFileId;
86use crate::sst::file_ref::FileReferenceManagerRef;
87use crate::sst::index::IndexBuildScheduler;
88use crate::sst::index::intermediate::IntermediateManager;
89use crate::sst::index::puffin_manager::PuffinManagerFactory;
90use crate::time_provider::{StdTimeProvider, TimeProviderRef};
91use crate::wal::Wal;
92use crate::worker::handle_manifest::RegionEditQueues;
93
94pub(crate) type WorkerId = u32;
96
97pub(crate) const DROPPING_MARKER_FILE: &str = ".dropping";
98
99pub(crate) const CHECK_REGION_INTERVAL: Duration = Duration::from_secs(60);
101pub(crate) const MAX_INITIAL_CHECK_DELAY_SECS: u64 = 60 * 3;
103
104#[cfg_attr(doc, aquamarine::aquamarine)]
105pub(crate) struct WorkerGroup {
142 workers: Vec<RegionWorker>,
144 flush_job_pool: SchedulerRef,
146 compact_job_pool: SchedulerRef,
148 index_build_job_pool: SchedulerRef,
150 purge_scheduler: SchedulerRef,
152 cache_manager: CacheManagerRef,
154 file_ref_manager: FileReferenceManagerRef,
156 gc_limiter: GcLimiterRef,
158 object_store_manager: ObjectStoreManagerRef,
160 puffin_manager_factory: PuffinManagerFactory,
162 intermediate_manager: IntermediateManager,
164 schema_metadata_manager: SchemaMetadataManagerRef,
166}
167
168impl WorkerGroup {
169 pub(crate) async fn start<S: LogStore>(
173 config: Arc<MitoConfig>,
174 log_store: Arc<S>,
175 object_store_manager: ObjectStoreManagerRef,
176 schema_metadata_manager: SchemaMetadataManagerRef,
177 file_ref_manager: FileReferenceManagerRef,
178 partition_expr_fetcher: PartitionExprFetcherRef,
179 plugins: Plugins,
180 ) -> Result<WorkerGroup> {
181 let (flush_sender, flush_receiver) = watch::channel(());
182 let write_buffer_manager = Arc::new(
183 WriteBufferManagerImpl::new(config.global_write_buffer_size.as_bytes() as usize)
184 .with_notifier(flush_sender.clone()),
185 );
186 let puffin_manager_factory = PuffinManagerFactory::new(
187 &config.index.aux_path,
188 config.index.staging_size.as_bytes(),
189 Some(config.index.write_buffer_size.as_bytes() as _),
190 config.index.staging_ttl,
191 )
192 .await?;
193 let intermediate_manager = IntermediateManager::init_fs(&config.index.aux_path)
194 .await?
195 .with_buffer_size(Some(config.index.write_buffer_size.as_bytes() as _));
196 let index_build_job_pool =
197 Arc::new(LocalScheduler::new(config.max_background_index_builds));
198 let series_index_store = if config.experimental_series_index_root.trim().is_empty() {
199 None
200 } else {
201 Some(new_fs_cache_store(&config.experimental_series_index_root).await?)
202 };
203 let flush_job_pool = Arc::new(LocalScheduler::new(config.max_background_flushes));
204 let compact_job_pool = Arc::new(LocalScheduler::new(config.max_background_compactions));
205 let flush_semaphore = Arc::new(Semaphore::new(config.max_background_flushes));
206 let purge_scheduler = Arc::new(LocalScheduler::new(config.max_background_purges));
208 let upload_store_wrapper = plugins.get::<WriteCacheUploadStoreWrapperRef>();
209 let write_cache = write_cache_from_config(
210 &config,
211 puffin_manager_factory.clone(),
212 intermediate_manager.clone(),
213 upload_store_wrapper,
214 )
215 .await?;
216 let cache_manager = Arc::new(
217 CacheManager::builder()
218 .sst_meta_cache_size(config.sst_meta_cache_size.as_bytes())
219 .vector_cache_size(config.vector_cache_size.as_bytes())
220 .page_cache_size(config.page_cache_size.as_bytes())
221 .selector_result_cache_size(config.selector_result_cache_size.as_bytes())
222 .range_result_cache_size(config.range_result_cache_size.as_bytes())
223 .prefilter_result_cache_size(config.prefilter_result_cache_size.as_bytes())
224 .index_metadata_size(config.index.metadata_cache_size.as_bytes())
225 .index_content_size(config.index.content_cache_size.as_bytes())
226 .index_content_page_size(config.index.content_cache_page_size.as_bytes())
227 .index_result_cache_size(config.index.result_cache_size.as_bytes())
228 .puffin_metadata_size(config.index.metadata_cache_size.as_bytes())
229 .write_cache(write_cache)
230 .build(),
231 );
232 let time_provider = Arc::new(StdTimeProvider);
233 let total_memory = get_total_memory_bytes();
234 let total_memory = if total_memory > 0 {
235 total_memory as u64
236 } else {
237 0
238 };
239 let compaction_limit_bytes = config
240 .experimental_compaction_memory_limit
241 .resolve(total_memory);
242 let compaction_memory_manager =
243 Arc::new(new_compaction_memory_manager(compaction_limit_bytes));
244 let gc_limiter = Arc::new(GcLimiter::new(config.gc.max_concurrent_gc_job));
245
246 let workers = (0..config.num_workers)
247 .map(|id| {
248 WorkerStarter {
249 id: id as WorkerId,
250 config: config.clone(),
251 log_store: log_store.clone(),
252 object_store_manager: object_store_manager.clone(),
253 write_buffer_manager: write_buffer_manager.clone(),
254 index_build_job_pool: index_build_job_pool.clone(),
255 series_index_store: series_index_store.clone(),
256 flush_job_pool: flush_job_pool.clone(),
257 compact_job_pool: compact_job_pool.clone(),
258 purge_scheduler: purge_scheduler.clone(),
259 listener: WorkerListener::default(),
260 cache_manager: cache_manager.clone(),
261 compaction_memory_manager: compaction_memory_manager.clone(),
262 puffin_manager_factory: puffin_manager_factory.clone(),
263 intermediate_manager: intermediate_manager.clone(),
264 time_provider: time_provider.clone(),
265 flush_sender: flush_sender.clone(),
266 flush_receiver: flush_receiver.clone(),
267 plugins: plugins.clone(),
268 schema_metadata_manager: schema_metadata_manager.clone(),
269 file_ref_manager: file_ref_manager.clone(),
270 partition_expr_fetcher: partition_expr_fetcher.clone(),
271 flush_semaphore: flush_semaphore.clone(),
272 }
273 .start()
274 })
275 .collect::<Result<Vec<_>>>()?;
276
277 Ok(WorkerGroup {
278 workers,
279 flush_job_pool,
280 compact_job_pool,
281 index_build_job_pool,
282 purge_scheduler,
283 cache_manager,
284 file_ref_manager,
285 gc_limiter,
286 object_store_manager,
287 puffin_manager_factory,
288 intermediate_manager,
289 schema_metadata_manager,
290 })
291 }
292
293 pub(crate) async fn stop(&self) -> Result<()> {
295 info!("Stop region worker group");
296
297 self.compact_job_pool.stop(true).await?;
300 self.flush_job_pool.stop(true).await?;
302 self.purge_scheduler.stop(true).await?;
304 self.index_build_job_pool.stop(true).await?;
306
307 try_join_all(self.workers.iter().map(|worker| worker.stop())).await?;
308
309 Ok(())
310 }
311
312 pub(crate) async fn submit_to_worker(
314 &self,
315 region_id: RegionId,
316 request: WorkerRequest,
317 ) -> Result<()> {
318 self.worker(region_id).submit_request(request).await
319 }
320
321 pub(crate) fn is_region_exists(&self, region_id: RegionId) -> bool {
323 self.worker(region_id).is_region_exists(region_id)
324 }
325
326 pub(crate) fn is_region_opening(&self, region_id: RegionId) -> bool {
328 self.worker(region_id).is_region_opening(region_id)
329 }
330
331 pub(crate) fn is_region_catching_up(&self, region_id: RegionId) -> bool {
333 self.worker(region_id).is_region_catching_up(region_id)
334 }
335
336 pub(crate) fn get_region(&self, region_id: RegionId) -> Option<MitoRegionRef> {
340 self.worker(region_id).get_region(region_id)
341 }
342
343 pub(crate) fn cache_manager(&self) -> CacheManagerRef {
345 self.cache_manager.clone()
346 }
347
348 pub(crate) fn file_ref_manager(&self) -> FileReferenceManagerRef {
349 self.file_ref_manager.clone()
350 }
351
352 pub(crate) fn gc_limiter(&self) -> GcLimiterRef {
353 self.gc_limiter.clone()
354 }
355
356 pub(crate) fn worker(&self, region_id: RegionId) -> &RegionWorker {
358 let index = region_id_to_index(region_id, self.workers.len());
359
360 &self.workers[index]
361 }
362
363 pub(crate) fn all_regions(&self) -> impl Iterator<Item = MitoRegionRef> + use<'_> {
364 self.workers
365 .iter()
366 .flat_map(|worker| worker.regions.list_regions())
367 }
368
369 pub(crate) fn object_store_manager(&self) -> &ObjectStoreManagerRef {
370 &self.object_store_manager
371 }
372
373 pub(crate) fn puffin_manager_factory(&self) -> &PuffinManagerFactory {
374 &self.puffin_manager_factory
375 }
376
377 pub(crate) fn intermediate_manager(&self) -> &IntermediateManager {
378 &self.intermediate_manager
379 }
380
381 pub(crate) fn schema_metadata_manager(&self) -> &SchemaMetadataManagerRef {
382 &self.schema_metadata_manager
383 }
384}
385
386#[cfg(any(test, feature = "test"))]
388impl WorkerGroup {
389 #[allow(clippy::too_many_arguments)]
393 pub(crate) async fn start_for_test<S: LogStore>(
394 config: Arc<MitoConfig>,
395 log_store: Arc<S>,
396 object_store_manager: ObjectStoreManagerRef,
397 write_buffer_manager: Option<WriteBufferManagerRef>,
398 listener: Option<crate::engine::listener::EventListenerRef>,
399 schema_metadata_manager: SchemaMetadataManagerRef,
400 file_ref_manager: FileReferenceManagerRef,
401 time_provider: TimeProviderRef,
402 partition_expr_fetcher: PartitionExprFetcherRef,
403 ) -> Result<WorkerGroup> {
404 let (flush_sender, flush_receiver) = watch::channel(());
405 let write_buffer_manager = write_buffer_manager.unwrap_or_else(|| {
406 Arc::new(
407 WriteBufferManagerImpl::new(config.global_write_buffer_size.as_bytes() as usize)
408 .with_notifier(flush_sender.clone()),
409 )
410 });
411 let index_build_job_pool =
412 Arc::new(LocalScheduler::new(config.max_background_index_builds));
413 let series_index_store = if config.experimental_series_index_root.trim().is_empty() {
414 None
415 } else {
416 Some(new_fs_cache_store(&config.experimental_series_index_root).await?)
417 };
418 let flush_job_pool = Arc::new(LocalScheduler::new(config.max_background_flushes));
419 let compact_job_pool = Arc::new(LocalScheduler::new(config.max_background_compactions));
420 let flush_semaphore = Arc::new(Semaphore::new(config.max_background_flushes));
421 let purge_scheduler = Arc::new(LocalScheduler::new(config.max_background_flushes));
422 let puffin_manager_factory = PuffinManagerFactory::new(
423 &config.index.aux_path,
424 config.index.staging_size.as_bytes(),
425 Some(config.index.write_buffer_size.as_bytes() as _),
426 config.index.staging_ttl,
427 )
428 .await?;
429 let intermediate_manager = IntermediateManager::init_fs(&config.index.aux_path)
430 .await?
431 .with_buffer_size(Some(config.index.write_buffer_size.as_bytes() as _));
432 let write_cache = write_cache_from_config(
433 &config,
434 puffin_manager_factory.clone(),
435 intermediate_manager.clone(),
436 None,
437 )
438 .await?;
439 let cache_manager = Arc::new(
440 CacheManager::builder()
441 .sst_meta_cache_size(config.sst_meta_cache_size.as_bytes())
442 .vector_cache_size(config.vector_cache_size.as_bytes())
443 .page_cache_size(config.page_cache_size.as_bytes())
444 .selector_result_cache_size(config.selector_result_cache_size.as_bytes())
445 .range_result_cache_size(config.range_result_cache_size.as_bytes())
446 .prefilter_result_cache_size(config.prefilter_result_cache_size.as_bytes())
447 .write_cache(write_cache)
448 .build(),
449 );
450 let total_memory = get_total_memory_bytes();
451 let total_memory = if total_memory > 0 {
452 total_memory as u64
453 } else {
454 0
455 };
456 let compaction_limit_bytes = config
457 .experimental_compaction_memory_limit
458 .resolve(total_memory);
459 let compaction_memory_manager =
460 Arc::new(new_compaction_memory_manager(compaction_limit_bytes));
461 let gc_limiter = Arc::new(GcLimiter::new(config.gc.max_concurrent_gc_job));
462 let workers = (0..config.num_workers)
463 .map(|id| {
464 WorkerStarter {
465 id: id as WorkerId,
466 config: config.clone(),
467 log_store: log_store.clone(),
468 object_store_manager: object_store_manager.clone(),
469 write_buffer_manager: write_buffer_manager.clone(),
470 index_build_job_pool: index_build_job_pool.clone(),
471 series_index_store: series_index_store.clone(),
472 flush_job_pool: flush_job_pool.clone(),
473 compact_job_pool: compact_job_pool.clone(),
474 purge_scheduler: purge_scheduler.clone(),
475 listener: WorkerListener::new(listener.clone()),
476 cache_manager: cache_manager.clone(),
477 compaction_memory_manager: compaction_memory_manager.clone(),
478 puffin_manager_factory: puffin_manager_factory.clone(),
479 intermediate_manager: intermediate_manager.clone(),
480 time_provider: time_provider.clone(),
481 flush_sender: flush_sender.clone(),
482 flush_receiver: flush_receiver.clone(),
483 plugins: Plugins::new(),
484 schema_metadata_manager: schema_metadata_manager.clone(),
485 file_ref_manager: file_ref_manager.clone(),
486 partition_expr_fetcher: partition_expr_fetcher.clone(),
487 flush_semaphore: flush_semaphore.clone(),
488 }
489 .start()
490 })
491 .collect::<Result<Vec<_>>>()?;
492
493 Ok(WorkerGroup {
494 workers,
495 flush_job_pool,
496 compact_job_pool,
497 index_build_job_pool,
498 purge_scheduler,
499 cache_manager,
500 file_ref_manager,
501 gc_limiter,
502 object_store_manager,
503 puffin_manager_factory,
504 intermediate_manager,
505 schema_metadata_manager,
506 })
507 }
508
509 pub(crate) fn purge_scheduler(&self) -> &SchedulerRef {
511 &self.purge_scheduler
512 }
513}
514
515fn region_id_to_index(id: RegionId, num_workers: usize) -> usize {
516 ((id.table_id() as usize % num_workers) + (id.region_number() as usize % num_workers))
517 % num_workers
518}
519
520pub async fn write_cache_from_config(
521 config: &MitoConfig,
522 puffin_manager_factory: PuffinManagerFactory,
523 intermediate_manager: IntermediateManager,
524 upload_store_wrapper: Option<WriteCacheUploadStoreWrapperRef>,
525) -> Result<Option<WriteCacheRef>> {
526 if !config.enable_write_cache {
527 return Ok(None);
528 }
529
530 tokio::fs::create_dir_all(Path::new(&config.write_cache_path))
531 .await
532 .context(CreateDirSnafu {
533 dir: &config.write_cache_path,
534 })?;
535
536 let cache = WriteCache::new_fs(
537 &config.write_cache_path,
538 config.write_cache_size,
539 config.write_cache_ttl,
540 Some(config.index_cache_percent),
541 config.enable_refill_cache_on_read,
542 puffin_manager_factory,
543 intermediate_manager,
544 config.manifest_cache_size,
545 )
546 .await?
547 .with_upload_store_wrapper(upload_store_wrapper);
548 Ok(Some(Arc::new(cache)))
549}
550
551pub(crate) fn worker_init_check_delay() -> Duration {
553 let init_check_delay = rng().random_range(0..MAX_INITIAL_CHECK_DELAY_SECS);
554 Duration::from_secs(init_check_delay)
555}
556
557struct WorkerStarter<S> {
559 id: WorkerId,
560 config: Arc<MitoConfig>,
561 log_store: Arc<S>,
562 object_store_manager: ObjectStoreManagerRef,
563 write_buffer_manager: WriteBufferManagerRef,
564 compact_job_pool: SchedulerRef,
565 index_build_job_pool: SchedulerRef,
566 series_index_store: Option<ObjectStore>,
567 flush_job_pool: SchedulerRef,
568 purge_scheduler: SchedulerRef,
569 listener: WorkerListener,
570 cache_manager: CacheManagerRef,
571 compaction_memory_manager: Arc<CompactionMemoryManager>,
572 puffin_manager_factory: PuffinManagerFactory,
573 intermediate_manager: IntermediateManager,
574 time_provider: TimeProviderRef,
575 flush_sender: watch::Sender<()>,
577 flush_receiver: watch::Receiver<()>,
579 plugins: Plugins,
580 schema_metadata_manager: SchemaMetadataManagerRef,
581 file_ref_manager: FileReferenceManagerRef,
582 partition_expr_fetcher: PartitionExprFetcherRef,
583 flush_semaphore: Arc<Semaphore>,
584}
585
586impl<S: LogStore> WorkerStarter<S> {
587 fn start(self) -> Result<RegionWorker> {
589 let regions = Arc::new(RegionMap::default());
590 let opening_regions = Arc::new(OpeningRegions::default());
591 let catchup_regions = Arc::new(CatchupRegions::default());
592 let (sender, receiver) = mpsc::channel(self.config.worker_channel_size);
593
594 let running = Arc::new(AtomicBool::new(true));
595 let series_index_task_state = self
596 .series_index_store
597 .as_ref()
598 .map(|_| Arc::new(SeriesIndexTaskState::new()));
599 let mut series_index_purger = None;
600 let series_index_handle = self
601 .series_index_store
602 .clone()
603 .zip(series_index_task_state.clone())
604 .map(|(store, state)| {
605 let (purger, purge_receiver) = series_index_channel(store.clone());
606 series_index_purger = Some(purger);
607 spawn_series_index_tasks(
608 self.id,
609 store,
610 state,
611 purge_receiver,
612 self.config.experimental_series_index_maintenance_interval,
613 )
614 });
615 let now = self.time_provider.current_time_millis();
616 let id_string = self.id.to_string();
617 let mut worker_thread = RegionWorkerLoop {
618 id: self.id,
619 config: self.config.clone(),
620 regions: regions.clone(),
621 catchup_regions: catchup_regions.clone(),
622 dropping_regions: Arc::new(RegionMap::default()),
623 opening_regions: opening_regions.clone(),
624 sender: sender.clone(),
625 receiver,
626 wal: Wal::new(self.log_store),
627 object_store_manager: self.object_store_manager.clone(),
628 running: running.clone(),
629 memtable_builder_provider: MemtableBuilderProvider::new(
630 Some(self.write_buffer_manager.clone()),
631 self.config.clone(),
632 ),
633 purge_scheduler: self.purge_scheduler.clone(),
634 write_buffer_manager: self.write_buffer_manager,
635 index_build_scheduler: IndexBuildScheduler::new(
636 self.index_build_job_pool,
637 self.config.max_background_index_builds,
638 ),
639 series_index_store: self.series_index_store,
640 series_index_purger,
641 flush_scheduler: FlushScheduler::new(self.flush_job_pool),
642 compaction_scheduler: CompactionScheduler::new(
643 self.compact_job_pool,
644 sender.clone(),
645 self.cache_manager.clone(),
646 self.config.clone(),
647 self.listener.clone(),
648 self.plugins.clone(),
649 self.compaction_memory_manager.clone(),
650 self.config.experimental_compaction_on_exhausted,
651 ),
652 stalled_requests: StalledRequests::default(),
653 listener: self.listener,
654 cache_manager: self.cache_manager,
655 puffin_manager_factory: self.puffin_manager_factory,
656 intermediate_manager: self.intermediate_manager,
657 time_provider: self.time_provider,
658 last_periodical_check_millis: now,
659 flush_sender: self.flush_sender,
660 flush_receiver: self.flush_receiver,
661 stalling_count: WRITE_STALLING.with_label_values(&[&id_string]),
662 region_count: REGION_COUNT.with_label_values(&[&id_string]),
663 request_wait_time: REQUEST_WAIT_TIME.with_label_values(&[&id_string]),
664 region_edit_queues: RegionEditQueues::default(),
665 schema_metadata_manager: self.schema_metadata_manager,
666 file_ref_manager: self.file_ref_manager.clone(),
667 partition_expr_fetcher: self.partition_expr_fetcher,
668 flush_semaphore: self.flush_semaphore,
669 plugins: self.plugins,
670 };
671 let handle = common_runtime::spawn_global(async move {
672 worker_thread.run().await;
673 });
674
675 Ok(RegionWorker {
676 id: self.id,
677 regions,
678 opening_regions,
679 catchup_regions,
680 sender,
681 handle: Mutex::new(Some(handle)),
682 series_index_handle: Mutex::new(series_index_handle),
683 series_index_task_state,
684 running,
685 })
686 }
687}
688
689pub(crate) struct RegionWorker {
691 id: WorkerId,
693 regions: RegionMapRef,
695 opening_regions: OpeningRegionsRef,
697 catchup_regions: CatchupRegionsRef,
699 sender: Sender<WorkerRequestWithTime>,
701 handle: Mutex<Option<JoinHandle<()>>>,
703 series_index_handle: Mutex<Option<JoinHandle<()>>>,
705 series_index_task_state: Option<Arc<SeriesIndexTaskState>>,
707 running: Arc<AtomicBool>,
709}
710
711impl RegionWorker {
712 async fn submit_request(&self, request: WorkerRequest) -> Result<()> {
714 ensure!(self.is_running(), WorkerStoppedSnafu { id: self.id });
715 let request_with_time = WorkerRequestWithTime::new(request);
716 if self.sender.send(request_with_time).await.is_err() {
717 warn!(
718 "Worker {} is already exited but the running flag is still true",
719 self.id
720 );
721 self.set_running(false);
723 if let Some(state) = &self.series_index_task_state {
724 state.stop();
725 }
726 return WorkerStoppedSnafu { id: self.id }.fail();
727 }
728
729 Ok(())
730 }
731
732 async fn stop(&self) -> Result<()> {
736 let handle = self.handle.lock().await.take();
737 self.set_running(false);
738 if let Some(state) = &self.series_index_task_state {
739 state.stop();
740 }
741
742 let mut worker_result = Ok(());
743 if let Some(handle) = handle {
744 info!("Stop region worker {}", self.id);
745
746 if self
747 .sender
748 .send(WorkerRequestWithTime::new(WorkerRequest::Stop))
749 .await
750 .is_err()
751 {
752 warn!("Worker {} is already exited before stop", self.id);
753 }
754
755 worker_result = handle.await.context(JoinSnafu);
756 }
757 let series_index_result = if let Some(handle) = self.series_index_handle.lock().await.take()
758 {
759 handle.await.context(JoinSnafu)
760 } else {
761 Ok(())
762 };
763 worker_result?;
764 series_index_result?;
765
766 Ok(())
767 }
768
769 fn is_running(&self) -> bool {
771 self.running.load(Ordering::Relaxed)
772 }
773
774 fn set_running(&self, value: bool) {
776 self.running.store(value, Ordering::Relaxed)
777 }
778
779 fn is_region_exists(&self, region_id: RegionId) -> bool {
781 self.regions.is_region_exists(region_id)
782 }
783
784 fn is_region_opening(&self, region_id: RegionId) -> bool {
786 self.opening_regions.is_region_exists(region_id)
787 }
788
789 fn is_region_catching_up(&self, region_id: RegionId) -> bool {
791 self.catchup_regions.is_region_exists(region_id)
792 }
793
794 fn get_region(&self, region_id: RegionId) -> Option<MitoRegionRef> {
796 self.regions.get_region(region_id)
797 }
798
799 #[cfg(test)]
800 pub(crate) fn opening_regions(&self) -> &OpeningRegionsRef {
802 &self.opening_regions
803 }
804
805 #[cfg(test)]
806 pub(crate) fn catchup_regions(&self) -> &CatchupRegionsRef {
808 &self.catchup_regions
809 }
810}
811
812impl Drop for RegionWorker {
813 fn drop(&mut self) {
814 if let Some(state) = &self.series_index_task_state {
815 state.stop();
816 }
817 if self.is_running() {
818 self.set_running(false);
819 }
821 }
822}
823
824type RequestBuffer = Vec<WorkerRequest>;
825
826#[derive(Default)]
830pub(crate) struct StalledRequests {
831 pub(crate) requests:
838 HashMap<RegionId, (usize, Vec<SenderWriteRequest>, Vec<SenderBulkRequest>)>,
839 pub(crate) estimated_size: usize,
841}
842
843impl StalledRequests {
844 pub(crate) fn estimated_size(&self, region_id: &RegionId) -> usize {
846 self.requests
847 .get(region_id)
848 .map(|(size, _, _)| *size)
849 .unwrap_or_default()
850 }
851
852 pub(crate) fn append(
854 &mut self,
855 requests: &mut Vec<SenderWriteRequest>,
856 bulk_requests: &mut Vec<SenderBulkRequest>,
857 ) {
858 for req in requests.drain(..) {
859 self.push(req);
860 }
861 for req in bulk_requests.drain(..) {
862 self.push_bulk(req);
863 }
864 }
865
866 pub(crate) fn push(&mut self, req: SenderWriteRequest) {
868 let (size, requests, _) = self.requests.entry(req.request.region_id).or_default();
869 let req_size = req.request.estimated_size();
870 *size += req_size;
871 self.estimated_size += req_size;
872 requests.push(req);
873 }
874
875 pub(crate) fn push_bulk(&mut self, req: SenderBulkRequest) {
876 let region_id = req.region_id;
877 let (size, _, requests) = self.requests.entry(region_id).or_default();
878 let req_size = req.request.estimated_size();
879 *size += req_size;
880 self.estimated_size += req_size;
881 requests.push(req);
882 }
883
884 pub(crate) fn remove(
886 &mut self,
887 region_id: &RegionId,
888 ) -> (Vec<SenderWriteRequest>, Vec<SenderBulkRequest>) {
889 if let Some((size, write_reqs, bulk_reqs)) = self.requests.remove(region_id) {
890 self.estimated_size -= size;
891 (write_reqs, bulk_reqs)
892 } else {
893 (vec![], vec![])
894 }
895 }
896
897 pub(crate) fn stalled_count(&self) -> usize {
899 self.requests
900 .values()
901 .map(|(_, reqs, bulk_reqs)| reqs.len() + bulk_reqs.len())
902 .sum()
903 }
904}
905
906struct RegionWorkerLoop<S> {
908 id: WorkerId,
910 config: Arc<MitoConfig>,
912 regions: RegionMapRef,
914 dropping_regions: RegionMapRef,
916 opening_regions: OpeningRegionsRef,
918 catchup_regions: CatchupRegionsRef,
920 sender: Sender<WorkerRequestWithTime>,
922 receiver: Receiver<WorkerRequestWithTime>,
924 wal: Wal<S>,
926 object_store_manager: ObjectStoreManagerRef,
928 running: Arc<AtomicBool>,
930 memtable_builder_provider: MemtableBuilderProvider,
932 purge_scheduler: SchedulerRef,
934 write_buffer_manager: WriteBufferManagerRef,
936 index_build_scheduler: IndexBuildScheduler,
938 series_index_store: Option<ObjectStore>,
940 series_index_purger: Option<IndexFilePurger>,
941 flush_scheduler: FlushScheduler,
943 compaction_scheduler: CompactionScheduler,
945 stalled_requests: StalledRequests,
947 listener: WorkerListener,
949 cache_manager: CacheManagerRef,
951 puffin_manager_factory: PuffinManagerFactory,
953 intermediate_manager: IntermediateManager,
955 time_provider: TimeProviderRef,
957 last_periodical_check_millis: i64,
959 flush_sender: watch::Sender<()>,
961 flush_receiver: watch::Receiver<()>,
963 stalling_count: IntGauge,
965 region_count: IntGauge,
967 request_wait_time: Histogram,
969 region_edit_queues: RegionEditQueues,
971 schema_metadata_manager: SchemaMetadataManagerRef,
973 file_ref_manager: FileReferenceManagerRef,
975 partition_expr_fetcher: PartitionExprFetcherRef,
977 flush_semaphore: Arc<Semaphore>,
979 plugins: Plugins,
981}
982
983impl<S: LogStore> RegionWorkerLoop<S> {
984 async fn run(&mut self) {
986 let init_check_delay = worker_init_check_delay();
987 info!(
988 "Start region worker thread {}, init_check_delay: {:?}",
989 self.id, init_check_delay
990 );
991 self.last_periodical_check_millis += init_check_delay.as_millis() as i64;
992
993 let mut write_req_buffer: Vec<SenderWriteRequest> =
995 Vec::with_capacity(self.config.worker_request_batch_size);
996 let mut bulk_req_buffer: Vec<SenderBulkRequest> =
997 Vec::with_capacity(self.config.worker_request_batch_size);
998 let mut ddl_req_buffer: Vec<SenderDdlRequest> =
999 Vec::with_capacity(self.config.worker_request_batch_size);
1000 let mut general_req_buffer: Vec<WorkerRequest> =
1001 RequestBuffer::with_capacity(self.config.worker_request_batch_size);
1002
1003 while self.running.load(Ordering::Relaxed) {
1004 write_req_buffer.clear();
1006 ddl_req_buffer.clear();
1007 general_req_buffer.clear();
1008 let mut bulk_insert_req_num = 0;
1009
1010 let max_wait_time = self.time_provider.wait_duration(CHECK_REGION_INTERVAL);
1011 let sleep = tokio::time::sleep(max_wait_time);
1012 tokio::pin!(sleep);
1013
1014 tokio::select! {
1015 request_opt = self.receiver.recv() => {
1016 match request_opt {
1017 Some(request_with_time) => {
1018 let wait_time = request_with_time.created_at.elapsed();
1020 self.request_wait_time.observe(wait_time.as_secs_f64());
1021
1022 match request_with_time.request {
1023 WorkerRequest::Write(sender_req) => write_req_buffer.push(sender_req),
1024 WorkerRequest::Ddl(sender_req) => ddl_req_buffer.push(sender_req),
1025 WorkerRequest::BulkInserts(bulk_insert) => {
1026 bulk_insert_req_num += 1;
1027 self.buffer_bulk_insert_request(bulk_insert, &mut bulk_req_buffer)
1028 .await;
1029 }
1030 req => general_req_buffer.push(req),
1031 }
1032 },
1033 None => break,
1035 }
1036 }
1037 recv_res = self.flush_receiver.changed() => {
1038 if recv_res.is_err() {
1039 break;
1041 } else {
1042 self.maybe_flush_worker();
1047 self.handle_stalled_requests().await;
1049 continue;
1050 }
1051 }
1052 _ = &mut sleep => {
1053 self.handle_periodical_tasks();
1055 continue;
1056 }
1057 }
1058
1059 if self.flush_receiver.has_changed().unwrap_or(false) {
1060 self.handle_stalled_requests().await;
1064 }
1065
1066 for _ in 1..self.config.worker_request_batch_size {
1068 match self.receiver.try_recv() {
1070 Ok(request_with_time) => {
1071 let wait_time = request_with_time.created_at.elapsed();
1073 self.request_wait_time.observe(wait_time.as_secs_f64());
1074
1075 match request_with_time.request {
1076 WorkerRequest::Write(sender_req) => write_req_buffer.push(sender_req),
1077 WorkerRequest::Ddl(sender_req) => ddl_req_buffer.push(sender_req),
1078 WorkerRequest::BulkInserts(bulk_insert) => {
1079 bulk_insert_req_num += 1;
1080 self.buffer_bulk_insert_request(bulk_insert, &mut bulk_req_buffer)
1081 .await
1082 }
1083 req => general_req_buffer.push(req),
1084 }
1085 }
1086 Err(_) => break,
1088 }
1089 }
1090
1091 self.listener.on_recv_requests(
1092 write_req_buffer.len()
1093 + ddl_req_buffer.len()
1094 + general_req_buffer.len()
1095 + bulk_insert_req_num,
1096 );
1097
1098 self.handle_requests(
1099 &mut write_req_buffer,
1100 &mut ddl_req_buffer,
1101 &mut general_req_buffer,
1102 &mut bulk_req_buffer,
1103 )
1104 .await;
1105
1106 self.handle_periodical_tasks();
1107 }
1108
1109 self.clean().await;
1110
1111 info!("Exit region worker thread {}", self.id);
1112 }
1113
1114 async fn buffer_bulk_insert_request(
1115 &mut self,
1116 bulk_insert: BulkInsertRequest,
1117 bulk_requests: &mut Vec<SenderBulkRequest>,
1118 ) {
1119 let BulkInsertRequest {
1120 metadata,
1121 request,
1122 sender,
1123 } = bulk_insert;
1124
1125 if let Some(region_metadata) = metadata {
1126 self.handle_bulk_insert_batch(region_metadata, request, bulk_requests, sender)
1127 .await;
1128 } else {
1129 error!("Cannot find region metadata for {}", request.region_id);
1130 sender.send(
1131 error::RegionNotFoundSnafu {
1132 region_id: request.region_id,
1133 }
1134 .fail(),
1135 );
1136 }
1137 }
1138
1139 async fn handle_requests(
1143 &mut self,
1144 write_requests: &mut Vec<SenderWriteRequest>,
1145 ddl_requests: &mut Vec<SenderDdlRequest>,
1146 general_requests: &mut Vec<WorkerRequest>,
1147 bulk_requests: &mut Vec<SenderBulkRequest>,
1148 ) {
1149 for worker_req in general_requests.drain(..) {
1150 match worker_req {
1151 WorkerRequest::Write(_) | WorkerRequest::Ddl(_) => {
1152 continue;
1154 }
1155 WorkerRequest::BulkInserts(_) => unreachable!("bulk inserts are buffered"),
1156 WorkerRequest::Background { region_id, notify } => {
1157 if matches!(
1158 ¬ify,
1159 BackgroundNotify::RegionEdit(edit_result)
1160 if edit_result.update_region_state
1161 ) {
1162 self.handle_buffered_region_write_requests(
1171 ®ion_id,
1172 write_requests,
1173 bulk_requests,
1174 )
1175 .await;
1176 }
1177 self.handle_background_notify(region_id, notify).await;
1179 }
1180 WorkerRequest::SetRegionRoleStateGracefully {
1181 region_id,
1182 region_role_state,
1183 sender,
1184 } => {
1185 self.set_role_state_gracefully(region_id, region_role_state, sender)
1186 .await;
1187 }
1188 WorkerRequest::EditRegion(request) => {
1189 self.handle_region_edit(request);
1190 }
1191 WorkerRequest::Stop => {
1192 debug_assert!(!self.running.load(Ordering::Relaxed));
1193 }
1194 WorkerRequest::SyncRegion(req) => {
1195 self.handle_region_sync(req).await;
1196 }
1197 WorkerRequest::RemapManifests(req) => {
1198 self.handle_remap_manifests_request(req);
1199 }
1200 WorkerRequest::CopyRegionFrom(req) => {
1201 self.handle_copy_region_from_request(req);
1202 }
1203 }
1204 }
1205
1206 self.handle_write_requests(write_requests, bulk_requests, true)
1209 .await;
1210
1211 self.handle_ddl_requests(ddl_requests).await;
1212 }
1213
1214 async fn handle_ddl_requests(&mut self, ddl_requests: &mut Vec<SenderDdlRequest>) {
1216 if ddl_requests.is_empty() {
1217 return;
1218 }
1219
1220 for ddl in ddl_requests.drain(..) {
1221 let res = match ddl.request {
1222 DdlRequest::Create(req) => self.handle_create_request(ddl.region_id, req).await,
1223 DdlRequest::Drop(req) => {
1224 self.handle_drop_request(ddl.region_id, req, ddl.sender)
1225 .await;
1226 continue;
1227 }
1228 DdlRequest::Open((req, wal_entry_receiver)) => {
1229 self.handle_open_request(ddl.region_id, req, wal_entry_receiver, ddl.sender)
1230 .await;
1231 continue;
1232 }
1233 DdlRequest::OfflineCleanup(req) => {
1234 self.handle_offline_cleanup_request(ddl.region_id, req)
1235 .await
1236 }
1237 DdlRequest::Close(req) => {
1238 self.handle_close_request(ddl.region_id, req, ddl.sender)
1239 .await;
1240 continue;
1241 }
1242 DdlRequest::Alter(req) => {
1243 self.handle_alter_request(ddl.region_id, req, ddl.sender)
1244 .await;
1245 continue;
1246 }
1247 DdlRequest::Flush(req) => {
1248 self.handle_flush_request(ddl.region_id, req, ddl.sender);
1249 continue;
1250 }
1251 DdlRequest::Compact(req) => {
1252 self.handle_compaction_request(ddl.region_id, req, ddl.sender)
1253 .await;
1254 continue;
1255 }
1256 DdlRequest::BuildIndex(req) => {
1257 self.handle_build_index_request(ddl.region_id, req, ddl.sender)
1258 .await;
1259 continue;
1260 }
1261 DdlRequest::Truncate(req) => {
1262 self.handle_truncate_request(ddl.region_id, req, ddl.sender)
1263 .await;
1264 continue;
1265 }
1266 DdlRequest::Catchup((req, wal_entry_receiver)) => {
1267 self.handle_catchup_request(ddl.region_id, req, wal_entry_receiver, ddl.sender)
1268 .await;
1269 continue;
1270 }
1271 DdlRequest::EnterStaging(req) => {
1272 self.handle_enter_staging_request(
1273 ddl.region_id,
1274 req.partition_directive,
1275 ddl.sender,
1276 )
1277 .await;
1278 continue;
1279 }
1280 DdlRequest::ApplyStagingManifest(req) => {
1281 self.handle_apply_staging_manifest_request(ddl.region_id, req, ddl.sender)
1282 .await;
1283 continue;
1284 }
1285 };
1286
1287 ddl.sender.send(res);
1288 }
1289 }
1290
1291 fn handle_periodical_tasks(&mut self) {
1293 let interval = CHECK_REGION_INTERVAL.as_millis() as i64;
1294 if self
1295 .time_provider
1296 .elapsed_since(self.last_periodical_check_millis)
1297 < interval
1298 {
1299 return;
1300 }
1301
1302 self.last_periodical_check_millis = self.time_provider.current_time_millis();
1303
1304 if let Err(e) = self.flush_periodically() {
1305 error!(e; "Failed to flush regions periodically");
1306 }
1307 }
1308
1309 async fn handle_background_notify(&mut self, region_id: RegionId, notify: BackgroundNotify) {
1311 match notify {
1312 BackgroundNotify::CompactionPickFinished(req) => {
1313 self.handle_compaction_pick_finished(region_id, req).await
1314 }
1315 BackgroundNotify::FlushFinished(req) => {
1316 self.handle_flush_finished(region_id, req).await
1317 }
1318 BackgroundNotify::FlushFailed(req) => self.handle_flush_failed(region_id, req).await,
1319 BackgroundNotify::IndexBuildFinished(req) => {
1320 self.handle_index_build_finished(region_id, req).await
1321 }
1322 BackgroundNotify::IndexBuildStopped(req) => {
1323 self.handle_index_build_stopped(region_id, req).await
1324 }
1325 BackgroundNotify::IndexBuildFailed(req) => {
1326 self.handle_index_build_failed(region_id, req).await
1327 }
1328 BackgroundNotify::IndexBuildRetry(req) => {
1329 self.handle_rebuild_index(req, OptionOutputTx::new(None))
1330 .await
1331 }
1332 BackgroundNotify::CompactionFinished(req) => {
1333 self.handle_compaction_finished(region_id, req).await
1334 }
1335 BackgroundNotify::CompactionCancelled(req) => {
1336 self.handle_compaction_cancelled(region_id, req).await
1337 }
1338 BackgroundNotify::CompactionFailed(req) => self.handle_compaction_failure(req).await,
1339 BackgroundNotify::Truncate(req) => self.handle_truncate_result(req).await,
1340 BackgroundNotify::DiscardUnflushed(req) => {
1341 self.handle_discard_unflushed_result(req).await
1342 }
1343 BackgroundNotify::RegionChange(req) => {
1344 self.handle_manifest_region_change_result(req).await
1345 }
1346 BackgroundNotify::EnterStaging(req) => self.handle_enter_staging_result(req).await,
1347 BackgroundNotify::RegionEdit(req) => self.handle_region_edit_result(req).await,
1348 BackgroundNotify::CopyRegionFromFinished(req) => {
1349 self.handle_copy_region_from_finished(req)
1350 }
1351 }
1352 }
1353
1354 async fn set_role_state_gracefully(
1356 &mut self,
1357 region_id: RegionId,
1358 region_role_state: SettableRegionRoleState,
1359 sender: oneshot::Sender<SetRegionRoleStateResponse>,
1360 ) {
1361 if let Some(region) = self.regions.get_region(region_id) {
1362 common_runtime::spawn_global(async move {
1364 match region.set_role_state_gracefully(region_role_state).await {
1365 Ok(()) => {
1366 let last_entry_id = region.version_control.current().last_entry_id;
1367 let _ = sender.send(SetRegionRoleStateResponse::success(
1368 SetRegionRoleStateSuccess::mito(last_entry_id),
1369 ));
1370 }
1371 Err(e) => {
1372 error!(e; "Failed to set region {} role state to {:?}", region_id, region_role_state);
1373 let _ = sender.send(SetRegionRoleStateResponse::invalid_transition(
1374 BoxedError::new(e),
1375 ));
1376 }
1377 }
1378 });
1379 } else {
1380 let _ = sender.send(SetRegionRoleStateResponse::NotFound);
1381 }
1382 }
1383}
1384
1385impl<S> RegionWorkerLoop<S> {
1386 async fn clean(&self) {
1388 let regions = self.regions.list_regions();
1390 for region in regions {
1391 region.stop().await;
1392 }
1393
1394 self.regions.clear();
1395 }
1396
1397 fn notify_group(&mut self) {
1400 let _ = self.flush_sender.send(());
1402 self.flush_receiver.borrow_and_update();
1404 }
1405}
1406
1407#[derive(Default, Clone)]
1409pub(crate) struct WorkerListener {
1410 #[cfg(any(test, feature = "test"))]
1411 listener: Option<crate::engine::listener::EventListenerRef>,
1412}
1413
1414impl WorkerListener {
1415 #[cfg(any(test, feature = "test"))]
1416 pub(crate) fn new(
1417 listener: Option<crate::engine::listener::EventListenerRef>,
1418 ) -> WorkerListener {
1419 WorkerListener { listener }
1420 }
1421
1422 pub(crate) fn on_flush_success(&self, region_id: RegionId) {
1424 #[cfg(any(test, feature = "test"))]
1425 if let Some(listener) = &self.listener {
1426 listener.on_flush_success(region_id);
1427 }
1428 let _ = region_id;
1430 }
1431
1432 pub(crate) fn on_write_stall(&self) {
1434 #[cfg(any(test, feature = "test"))]
1435 if let Some(listener) = &self.listener {
1436 listener.on_write_stall();
1437 }
1438 }
1439
1440 pub(crate) async fn on_flush_begin(&self, region_id: RegionId) {
1441 #[cfg(any(test, feature = "test"))]
1442 if let Some(listener) = &self.listener {
1443 listener.on_flush_begin(region_id).await;
1444 }
1445 let _ = region_id;
1447 }
1448
1449 pub(crate) async fn on_flush_commit_begin(&self, _region_id: RegionId) {
1450 #[cfg(any(test, feature = "test"))]
1451 if let Some(listener) = &self.listener {
1452 listener.on_flush_commit_begin(_region_id).await;
1453 }
1454 }
1455
1456 pub(crate) fn on_flush_cancel_requested(&self, _region_id: RegionId) {
1457 #[cfg(any(test, feature = "test"))]
1458 if let Some(listener) = &self.listener {
1459 listener.on_flush_cancel_requested(_region_id);
1460 }
1461 }
1462
1463 pub(crate) fn on_later_drop_begin(&self, region_id: RegionId) -> Option<Duration> {
1464 #[cfg(any(test, feature = "test"))]
1465 if let Some(listener) = &self.listener {
1466 return listener.on_later_drop_begin(region_id);
1467 }
1468 let _ = region_id;
1470 None
1471 }
1472
1473 pub(crate) fn on_later_drop_end(&self, region_id: RegionId, removed: bool) {
1475 #[cfg(any(test, feature = "test"))]
1476 if let Some(listener) = &self.listener {
1477 listener.on_later_drop_end(region_id, removed);
1478 }
1479 let _ = region_id;
1481 let _ = removed;
1482 }
1483
1484 pub(crate) async fn on_merge_ssts_finished(&self, region_id: RegionId) {
1485 #[cfg(any(test, feature = "test"))]
1486 if let Some(listener) = &self.listener {
1487 listener.on_merge_ssts_finished(region_id).await;
1488 }
1489 let _ = region_id;
1491 }
1492
1493 pub(crate) fn on_recv_requests(&self, request_num: usize) {
1494 #[cfg(any(test, feature = "test"))]
1495 if let Some(listener) = &self.listener {
1496 listener.on_recv_requests(request_num);
1497 }
1498 let _ = request_num;
1500 }
1501
1502 pub(crate) fn on_file_cache_filled(&self, _file_id: FileId) {
1503 #[cfg(any(test, feature = "test"))]
1504 if let Some(listener) = &self.listener {
1505 listener.on_file_cache_filled(_file_id);
1506 }
1507 }
1508
1509 pub(crate) fn on_compaction_scheduled(&self, _region_id: RegionId) {
1510 #[cfg(any(test, feature = "test"))]
1511 if let Some(listener) = &self.listener {
1512 listener.on_compaction_scheduled(_region_id);
1513 }
1514 }
1515
1516 pub(crate) async fn on_compaction_pick_begin(&self, _region_id: RegionId) {
1517 #[cfg(any(test, feature = "test"))]
1518 if let Some(listener) = &self.listener {
1519 listener.on_compaction_pick_begin(_region_id).await;
1520 }
1521 }
1522
1523 pub(crate) async fn on_compaction_commit_begin(&self, _region_id: RegionId) {
1524 #[cfg(any(test, feature = "test"))]
1525 if let Some(listener) = &self.listener {
1526 listener.on_compaction_commit_begin(_region_id).await;
1527 }
1528 }
1529
1530 pub(crate) async fn on_compaction_result_notified(&self, _region_id: RegionId) {
1531 #[cfg(any(test, feature = "test"))]
1532 if let Some(listener) = &self.listener {
1533 listener.on_compaction_result_notified(_region_id).await;
1534 }
1535 }
1536
1537 pub(crate) fn on_compaction_cancel_requested(&self, _region_id: RegionId) {
1538 #[cfg(any(test, feature = "test"))]
1539 if let Some(listener) = &self.listener {
1540 listener.on_compaction_cancel_requested(_region_id);
1541 }
1542 }
1543
1544 pub(crate) async fn on_notify_region_change_result_begin(&self, _region_id: RegionId) {
1545 #[cfg(any(test, feature = "test"))]
1546 if let Some(listener) = &self.listener {
1547 listener
1548 .on_notify_region_change_result_begin(_region_id)
1549 .await;
1550 }
1551 }
1552
1553 pub(crate) async fn on_enter_staging_result_begin(&self, _region_id: RegionId) {
1554 #[cfg(any(test, feature = "test"))]
1555 if let Some(listener) = &self.listener {
1556 listener.on_enter_staging_result_begin(_region_id).await;
1557 }
1558 }
1559
1560 pub(crate) async fn on_index_build_finish(&self, _region_file_id: RegionFileId) {
1561 #[cfg(any(test, feature = "test"))]
1562 if let Some(listener) = &self.listener {
1563 listener.on_index_build_finish(_region_file_id).await;
1564 }
1565 }
1566
1567 pub(crate) async fn on_index_build_begin(&self, _region_file_id: RegionFileId) {
1568 #[cfg(any(test, feature = "test"))]
1569 if let Some(listener) = &self.listener {
1570 listener.on_index_build_begin(_region_file_id).await;
1571 }
1572 }
1573
1574 pub(crate) async fn on_index_build_abort(&self, _region_file_id: RegionFileId) {
1575 #[cfg(any(test, feature = "test"))]
1576 if let Some(listener) = &self.listener {
1577 listener.on_index_build_abort(_region_file_id).await;
1578 }
1579 }
1580
1581 pub(crate) async fn on_index_build_before_manifest_commit(
1582 &self,
1583 _region_file_id: RegionFileId,
1584 ) {
1585 #[cfg(any(test, feature = "test"))]
1586 if let Some(listener) = &self.listener {
1587 listener
1588 .on_index_build_before_manifest_commit(_region_file_id)
1589 .await;
1590 }
1591 }
1592
1593 pub(crate) async fn on_index_build_manifest_committed(&self, _region_file_id: RegionFileId) {
1594 #[cfg(any(test, feature = "test"))]
1595 if let Some(listener) = &self.listener {
1596 listener
1597 .on_index_build_manifest_committed(_region_file_id)
1598 .await;
1599 }
1600 }
1601}
1602
1603#[cfg(test)]
1604mod tests {
1605 use super::*;
1606 use crate::test_util::TestEnv;
1607
1608 #[test]
1609 fn test_region_id_to_index() {
1610 let num_workers = 4;
1611
1612 let region_id = RegionId::new(1, 2);
1613 let index = region_id_to_index(region_id, num_workers);
1614 assert_eq!(index, 3);
1615
1616 let region_id = RegionId::new(2, 3);
1617 let index = region_id_to_index(region_id, num_workers);
1618 assert_eq!(index, 1);
1619 }
1620
1621 #[tokio::test]
1622 async fn test_worker_group_start_stop() {
1623 let env = TestEnv::with_prefix("group-stop").await;
1624 let group = env
1625 .create_worker_group(MitoConfig {
1626 num_workers: 4,
1627 experimental_series_index_root: "series-index".to_string(),
1628 ..Default::default()
1629 })
1630 .await;
1631
1632 tokio::time::timeout(Duration::from_secs(5), group.stop())
1633 .await
1634 .expect("series-index tasks should stop without waiting for their interval")
1635 .unwrap();
1636 }
1637}