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, AtomicU64, 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 #[allow(clippy::too_many_arguments)]
173 pub(crate) async fn start<S: LogStore>(
174 data_home: &str,
175 config: Arc<MitoConfig>,
176 log_store: Arc<S>,
177 object_store_manager: ObjectStoreManagerRef,
178 schema_metadata_manager: SchemaMetadataManagerRef,
179 file_ref_manager: FileReferenceManagerRef,
180 partition_expr_fetcher: PartitionExprFetcherRef,
181 plugins: Plugins,
182 ) -> Result<WorkerGroup> {
183 let (flush_sender, flush_receiver) = watch::channel(());
184 let write_buffer_manager = Arc::new(
185 WriteBufferManagerImpl::new(config.global_write_buffer_size.as_bytes() as usize)
186 .with_notifier(flush_sender.clone()),
187 );
188 let puffin_manager_factory = PuffinManagerFactory::new(
189 &config.index.aux_path,
190 config.index.staging_size.as_bytes(),
191 Some(config.index.write_buffer_size.as_bytes() as _),
192 config.index.staging_ttl,
193 )
194 .await?;
195 let intermediate_manager = IntermediateManager::init_fs(&config.index.aux_path)
196 .await?
197 .with_buffer_size(Some(config.index.write_buffer_size.as_bytes() as _));
198 let index_build_job_pool =
199 Arc::new(LocalScheduler::new(config.max_background_index_builds));
200 let series_index_store = series_index_store_from_config(&config, data_home).await?;
201 let series_index_disk_usage = Arc::new(AtomicU64::new(0));
202 let flush_job_pool = Arc::new(LocalScheduler::new(config.max_background_flushes));
203 let compact_job_pool = Arc::new(LocalScheduler::new(config.max_background_compactions));
204 let flush_semaphore = Arc::new(Semaphore::new(config.max_background_flushes));
205 let purge_scheduler = Arc::new(LocalScheduler::new(config.max_background_purges));
207 let upload_store_wrapper = plugins.get::<WriteCacheUploadStoreWrapperRef>();
208 let write_cache = write_cache_from_config(
209 &config,
210 puffin_manager_factory.clone(),
211 intermediate_manager.clone(),
212 upload_store_wrapper,
213 )
214 .await?;
215 let cache_manager = Arc::new(
216 CacheManager::builder()
217 .sst_meta_cache_size(config.sst_meta_cache_size.as_bytes())
218 .vector_cache_size(config.vector_cache_size.as_bytes())
219 .page_cache_size(config.page_cache_size.as_bytes())
220 .selector_result_cache_size(config.selector_result_cache_size.as_bytes())
221 .range_result_cache_size(config.range_result_cache_size.as_bytes())
222 .prefilter_result_cache_size(config.prefilter_result_cache_size.as_bytes())
223 .index_metadata_size(config.index.metadata_cache_size.as_bytes())
224 .index_content_size(config.index.content_cache_size.as_bytes())
225 .index_content_page_size(config.index.content_cache_page_size.as_bytes())
226 .index_result_cache_size(config.index.result_cache_size.as_bytes())
227 .puffin_metadata_size(config.index.metadata_cache_size.as_bytes())
228 .write_cache(write_cache)
229 .build(),
230 );
231 let time_provider = Arc::new(StdTimeProvider);
232 let total_memory = get_total_memory_bytes();
233 let total_memory = if total_memory > 0 {
234 total_memory as u64
235 } else {
236 0
237 };
238 let compaction_limit_bytes = config
239 .experimental_compaction_memory_limit
240 .resolve(total_memory);
241 let compaction_memory_manager =
242 Arc::new(new_compaction_memory_manager(compaction_limit_bytes));
243 let gc_limiter = Arc::new(GcLimiter::new(config.gc.max_concurrent_gc_job));
244
245 let workers = (0..config.num_workers)
246 .map(|id| {
247 WorkerStarter {
248 id: id as WorkerId,
249 config: config.clone(),
250 log_store: log_store.clone(),
251 object_store_manager: object_store_manager.clone(),
252 write_buffer_manager: write_buffer_manager.clone(),
253 index_build_job_pool: index_build_job_pool.clone(),
254 series_index_store: series_index_store.clone(),
255 series_index_disk_usage: series_index_disk_usage.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 data_home: &str,
395 config: Arc<MitoConfig>,
396 log_store: Arc<S>,
397 object_store_manager: ObjectStoreManagerRef,
398 write_buffer_manager: Option<WriteBufferManagerRef>,
399 listener: Option<crate::engine::listener::EventListenerRef>,
400 schema_metadata_manager: SchemaMetadataManagerRef,
401 file_ref_manager: FileReferenceManagerRef,
402 time_provider: TimeProviderRef,
403 partition_expr_fetcher: PartitionExprFetcherRef,
404 ) -> Result<WorkerGroup> {
405 let (flush_sender, flush_receiver) = watch::channel(());
406 let write_buffer_manager = write_buffer_manager.unwrap_or_else(|| {
407 Arc::new(
408 WriteBufferManagerImpl::new(config.global_write_buffer_size.as_bytes() as usize)
409 .with_notifier(flush_sender.clone()),
410 )
411 });
412 let index_build_job_pool =
413 Arc::new(LocalScheduler::new(config.max_background_index_builds));
414 let series_index_store = series_index_store_from_config(&config, data_home).await?;
415 let series_index_disk_usage = Arc::new(AtomicU64::new(0));
416 let flush_job_pool = Arc::new(LocalScheduler::new(config.max_background_flushes));
417 let compact_job_pool = Arc::new(LocalScheduler::new(config.max_background_compactions));
418 let flush_semaphore = Arc::new(Semaphore::new(config.max_background_flushes));
419 let purge_scheduler = Arc::new(LocalScheduler::new(config.max_background_flushes));
420 let puffin_manager_factory = PuffinManagerFactory::new(
421 &config.index.aux_path,
422 config.index.staging_size.as_bytes(),
423 Some(config.index.write_buffer_size.as_bytes() as _),
424 config.index.staging_ttl,
425 )
426 .await?;
427 let intermediate_manager = IntermediateManager::init_fs(&config.index.aux_path)
428 .await?
429 .with_buffer_size(Some(config.index.write_buffer_size.as_bytes() as _));
430 let write_cache = write_cache_from_config(
431 &config,
432 puffin_manager_factory.clone(),
433 intermediate_manager.clone(),
434 None,
435 )
436 .await?;
437 let cache_manager = Arc::new(
438 CacheManager::builder()
439 .sst_meta_cache_size(config.sst_meta_cache_size.as_bytes())
440 .vector_cache_size(config.vector_cache_size.as_bytes())
441 .page_cache_size(config.page_cache_size.as_bytes())
442 .selector_result_cache_size(config.selector_result_cache_size.as_bytes())
443 .range_result_cache_size(config.range_result_cache_size.as_bytes())
444 .prefilter_result_cache_size(config.prefilter_result_cache_size.as_bytes())
445 .write_cache(write_cache)
446 .build(),
447 );
448 let total_memory = get_total_memory_bytes();
449 let total_memory = if total_memory > 0 {
450 total_memory as u64
451 } else {
452 0
453 };
454 let compaction_limit_bytes = config
455 .experimental_compaction_memory_limit
456 .resolve(total_memory);
457 let compaction_memory_manager =
458 Arc::new(new_compaction_memory_manager(compaction_limit_bytes));
459 let gc_limiter = Arc::new(GcLimiter::new(config.gc.max_concurrent_gc_job));
460 let workers = (0..config.num_workers)
461 .map(|id| {
462 WorkerStarter {
463 id: id as WorkerId,
464 config: config.clone(),
465 log_store: log_store.clone(),
466 object_store_manager: object_store_manager.clone(),
467 write_buffer_manager: write_buffer_manager.clone(),
468 index_build_job_pool: index_build_job_pool.clone(),
469 series_index_store: series_index_store.clone(),
470 series_index_disk_usage: series_index_disk_usage.clone(),
471 flush_job_pool: flush_job_pool.clone(),
472 compact_job_pool: compact_job_pool.clone(),
473 purge_scheduler: purge_scheduler.clone(),
474 listener: WorkerListener::new(listener.clone()),
475 cache_manager: cache_manager.clone(),
476 compaction_memory_manager: compaction_memory_manager.clone(),
477 puffin_manager_factory: puffin_manager_factory.clone(),
478 intermediate_manager: intermediate_manager.clone(),
479 time_provider: time_provider.clone(),
480 flush_sender: flush_sender.clone(),
481 flush_receiver: flush_receiver.clone(),
482 plugins: Plugins::new(),
483 schema_metadata_manager: schema_metadata_manager.clone(),
484 file_ref_manager: file_ref_manager.clone(),
485 partition_expr_fetcher: partition_expr_fetcher.clone(),
486 flush_semaphore: flush_semaphore.clone(),
487 }
488 .start()
489 })
490 .collect::<Result<Vec<_>>>()?;
491
492 Ok(WorkerGroup {
493 workers,
494 flush_job_pool,
495 compact_job_pool,
496 index_build_job_pool,
497 purge_scheduler,
498 cache_manager,
499 file_ref_manager,
500 gc_limiter,
501 object_store_manager,
502 puffin_manager_factory,
503 intermediate_manager,
504 schema_metadata_manager,
505 })
506 }
507
508 pub(crate) fn purge_scheduler(&self) -> &SchedulerRef {
510 &self.purge_scheduler
511 }
512}
513
514fn region_id_to_index(id: RegionId, num_workers: usize) -> usize {
515 ((id.table_id() as usize % num_workers) + (id.region_number() as usize % num_workers))
516 % num_workers
517}
518
519async fn series_index_store_from_config(
521 config: &MitoConfig,
522 data_home: &str,
523) -> Result<Option<ObjectStore>> {
524 if !config.experimental_enable_series_index {
525 return Ok(None);
526 }
527
528 let root = Path::new(data_home).join("series_index");
529 new_fs_cache_store(&root.to_string_lossy()).await.map(Some)
530}
531
532pub async fn write_cache_from_config(
533 config: &MitoConfig,
534 puffin_manager_factory: PuffinManagerFactory,
535 intermediate_manager: IntermediateManager,
536 upload_store_wrapper: Option<WriteCacheUploadStoreWrapperRef>,
537) -> Result<Option<WriteCacheRef>> {
538 if !config.enable_write_cache {
539 return Ok(None);
540 }
541
542 tokio::fs::create_dir_all(Path::new(&config.write_cache_path))
543 .await
544 .context(CreateDirSnafu {
545 dir: &config.write_cache_path,
546 })?;
547
548 let cache = WriteCache::new_fs(
549 &config.write_cache_path,
550 config.write_cache_size,
551 config.write_cache_ttl,
552 Some(config.index_cache_percent),
553 config.enable_refill_cache_on_read,
554 puffin_manager_factory,
555 intermediate_manager,
556 config.manifest_cache_size,
557 )
558 .await?
559 .with_upload_store_wrapper(upload_store_wrapper);
560 Ok(Some(Arc::new(cache)))
561}
562
563pub(crate) fn worker_init_check_delay() -> Duration {
565 let init_check_delay = rng().random_range(0..MAX_INITIAL_CHECK_DELAY_SECS);
566 Duration::from_secs(init_check_delay)
567}
568
569struct WorkerStarter<S> {
571 id: WorkerId,
572 config: Arc<MitoConfig>,
573 log_store: Arc<S>,
574 object_store_manager: ObjectStoreManagerRef,
575 write_buffer_manager: WriteBufferManagerRef,
576 compact_job_pool: SchedulerRef,
577 index_build_job_pool: SchedulerRef,
578 series_index_store: Option<ObjectStore>,
579 series_index_disk_usage: Arc<AtomicU64>,
580 flush_job_pool: SchedulerRef,
581 purge_scheduler: SchedulerRef,
582 listener: WorkerListener,
583 cache_manager: CacheManagerRef,
584 compaction_memory_manager: Arc<CompactionMemoryManager>,
585 puffin_manager_factory: PuffinManagerFactory,
586 intermediate_manager: IntermediateManager,
587 time_provider: TimeProviderRef,
588 flush_sender: watch::Sender<()>,
590 flush_receiver: watch::Receiver<()>,
592 plugins: Plugins,
593 schema_metadata_manager: SchemaMetadataManagerRef,
594 file_ref_manager: FileReferenceManagerRef,
595 partition_expr_fetcher: PartitionExprFetcherRef,
596 flush_semaphore: Arc<Semaphore>,
597}
598
599impl<S: LogStore> WorkerStarter<S> {
600 fn start(self) -> Result<RegionWorker> {
602 let regions = Arc::new(RegionMap::default());
603 let opening_regions = Arc::new(OpeningRegions::default());
604 let catchup_regions = Arc::new(CatchupRegions::default());
605 let (sender, receiver) = mpsc::channel(self.config.worker_channel_size);
606
607 let running = Arc::new(AtomicBool::new(true));
608 let (series_index_task_state, series_index_receiver) = if self.series_index_store.is_some()
609 {
610 let (state, receiver) =
611 SeriesIndexTaskState::new(self.id, self.config.worker_channel_size);
612 (Some(Arc::new(state)), Some(receiver))
613 } else {
614 (None, None)
615 };
616 let mut series_index_purger = None;
617 let series_index_handle = self
618 .series_index_store
619 .clone()
620 .zip(series_index_task_state.clone())
621 .zip(series_index_receiver)
622 .map(|((store, state), receiver)| {
623 let (purger, purge_receiver) = series_index_channel(store.clone());
624 series_index_purger = Some(purger.clone());
625 spawn_series_index_tasks(
626 self.id,
627 store,
628 regions.clone(),
629 state,
630 receiver,
631 self.config.experimental_series_index_bucket_width,
632 purger,
633 purge_receiver,
634 self.series_index_disk_usage.clone(),
635 self.config.experimental_series_index_max_size.as_bytes(),
636 self.config.experimental_series_index_maintenance_interval,
637 self.time_provider.clone(),
638 self.config.experimental_enable_range_index,
639 )
640 });
641 let now = self.time_provider.current_time_millis();
642 let id_string = self.id.to_string();
643 let mut worker_thread = RegionWorkerLoop {
644 id: self.id,
645 config: self.config.clone(),
646 regions: regions.clone(),
647 catchup_regions: catchup_regions.clone(),
648 dropping_regions: Arc::new(RegionMap::default()),
649 opening_regions: opening_regions.clone(),
650 sender: sender.clone(),
651 receiver,
652 wal: Wal::new(self.log_store),
653 object_store_manager: self.object_store_manager.clone(),
654 running: running.clone(),
655 memtable_builder_provider: MemtableBuilderProvider::new(
656 Some(self.write_buffer_manager.clone()),
657 self.config.clone(),
658 ),
659 purge_scheduler: self.purge_scheduler.clone(),
660 write_buffer_manager: self.write_buffer_manager,
661 index_build_scheduler: IndexBuildScheduler::new(
662 self.index_build_job_pool,
663 self.config.max_background_index_builds,
664 ),
665 series_index_task_state: series_index_task_state.clone(),
666 series_index_store: self.series_index_store,
667 series_index_purger,
668 flush_scheduler: FlushScheduler::new(self.flush_job_pool),
669 compaction_scheduler: CompactionScheduler::new(
670 self.compact_job_pool,
671 sender.clone(),
672 self.cache_manager.clone(),
673 self.config.clone(),
674 self.listener.clone(),
675 self.plugins.clone(),
676 self.compaction_memory_manager.clone(),
677 self.config.experimental_compaction_on_exhausted,
678 ),
679 stalled_requests: StalledRequests::default(),
680 listener: self.listener,
681 cache_manager: self.cache_manager,
682 puffin_manager_factory: self.puffin_manager_factory,
683 intermediate_manager: self.intermediate_manager,
684 time_provider: self.time_provider,
685 last_periodical_check_millis: now,
686 flush_sender: self.flush_sender,
687 flush_receiver: self.flush_receiver,
688 stalling_count: WRITE_STALLING.with_label_values(&[&id_string]),
689 region_count: REGION_COUNT.with_label_values(&[&id_string]),
690 request_wait_time: REQUEST_WAIT_TIME.with_label_values(&[&id_string]),
691 region_edit_queues: RegionEditQueues::default(),
692 schema_metadata_manager: self.schema_metadata_manager,
693 file_ref_manager: self.file_ref_manager.clone(),
694 partition_expr_fetcher: self.partition_expr_fetcher,
695 flush_semaphore: self.flush_semaphore,
696 plugins: self.plugins,
697 };
698 let handle = common_runtime::spawn_global(async move {
699 worker_thread.run().await;
700 });
701
702 Ok(RegionWorker {
703 id: self.id,
704 regions,
705 opening_regions,
706 catchup_regions,
707 sender,
708 handle: Mutex::new(Some(handle)),
709 series_index_handle: Mutex::new(series_index_handle),
710 series_index_task_state,
711 running,
712 })
713 }
714}
715
716pub(crate) struct RegionWorker {
718 id: WorkerId,
720 regions: RegionMapRef,
722 opening_regions: OpeningRegionsRef,
724 catchup_regions: CatchupRegionsRef,
726 sender: Sender<WorkerRequestWithTime>,
728 handle: Mutex<Option<JoinHandle<()>>>,
730 series_index_handle: Mutex<Option<JoinHandle<()>>>,
732 series_index_task_state: Option<Arc<SeriesIndexTaskState>>,
734 running: Arc<AtomicBool>,
736}
737
738impl RegionWorker {
739 async fn submit_request(&self, request: WorkerRequest) -> Result<()> {
741 ensure!(self.is_running(), WorkerStoppedSnafu { id: self.id });
742 let request_with_time = WorkerRequestWithTime::new(request);
743 if self.sender.send(request_with_time).await.is_err() {
744 warn!(
745 "Worker {} is already exited but the running flag is still true",
746 self.id
747 );
748 self.set_running(false);
750 if let Some(state) = &self.series_index_task_state {
751 state.stop();
752 }
753 return WorkerStoppedSnafu { id: self.id }.fail();
754 }
755
756 Ok(())
757 }
758
759 async fn stop(&self) -> Result<()> {
763 let handle = self.handle.lock().await.take();
764 self.set_running(false);
765 if let Some(state) = &self.series_index_task_state {
766 state.stop();
767 }
768
769 let mut worker_result = Ok(());
770 if let Some(handle) = handle {
771 info!("Stop region worker {}", self.id);
772
773 if self
774 .sender
775 .send(WorkerRequestWithTime::new(WorkerRequest::Stop))
776 .await
777 .is_err()
778 {
779 warn!("Worker {} is already exited before stop", self.id);
780 }
781
782 worker_result = handle.await.context(JoinSnafu);
783 }
784 let series_index_result = if let Some(handle) = self.series_index_handle.lock().await.take()
785 {
786 handle.await.context(JoinSnafu)
787 } else {
788 Ok(())
789 };
790 worker_result?;
791 series_index_result?;
792
793 Ok(())
794 }
795
796 fn is_running(&self) -> bool {
798 self.running.load(Ordering::Relaxed)
799 }
800
801 fn set_running(&self, value: bool) {
803 self.running.store(value, Ordering::Relaxed)
804 }
805
806 fn is_region_exists(&self, region_id: RegionId) -> bool {
808 self.regions.is_region_exists(region_id)
809 }
810
811 fn is_region_opening(&self, region_id: RegionId) -> bool {
813 self.opening_regions.is_region_exists(region_id)
814 }
815
816 fn is_region_catching_up(&self, region_id: RegionId) -> bool {
818 self.catchup_regions.is_region_exists(region_id)
819 }
820
821 fn get_region(&self, region_id: RegionId) -> Option<MitoRegionRef> {
823 self.regions.get_region(region_id)
824 }
825
826 #[cfg(test)]
827 pub(crate) fn opening_regions(&self) -> &OpeningRegionsRef {
829 &self.opening_regions
830 }
831
832 #[cfg(test)]
833 pub(crate) fn catchup_regions(&self) -> &CatchupRegionsRef {
835 &self.catchup_regions
836 }
837}
838
839impl Drop for RegionWorker {
840 fn drop(&mut self) {
841 if let Some(state) = &self.series_index_task_state {
842 state.stop();
843 }
844 if self.is_running() {
845 self.set_running(false);
846 }
848 }
849}
850
851type RequestBuffer = Vec<WorkerRequest>;
852
853#[derive(Default)]
857pub(crate) struct StalledRequests {
858 pub(crate) requests:
865 HashMap<RegionId, (usize, Vec<SenderWriteRequest>, Vec<SenderBulkRequest>)>,
866 pub(crate) estimated_size: usize,
868}
869
870impl StalledRequests {
871 pub(crate) fn estimated_size(&self, region_id: &RegionId) -> usize {
873 self.requests
874 .get(region_id)
875 .map(|(size, _, _)| *size)
876 .unwrap_or_default()
877 }
878
879 pub(crate) fn append(
881 &mut self,
882 requests: &mut Vec<SenderWriteRequest>,
883 bulk_requests: &mut Vec<SenderBulkRequest>,
884 ) {
885 for req in requests.drain(..) {
886 self.push(req);
887 }
888 for req in bulk_requests.drain(..) {
889 self.push_bulk(req);
890 }
891 }
892
893 pub(crate) fn push(&mut self, req: SenderWriteRequest) {
895 let (size, requests, _) = self.requests.entry(req.request.region_id).or_default();
896 let req_size = req.request.estimated_size();
897 *size += req_size;
898 self.estimated_size += req_size;
899 requests.push(req);
900 }
901
902 pub(crate) fn push_bulk(&mut self, req: SenderBulkRequest) {
903 let region_id = req.region_id;
904 let (size, _, requests) = self.requests.entry(region_id).or_default();
905 let req_size = req.request.estimated_size();
906 *size += req_size;
907 self.estimated_size += req_size;
908 requests.push(req);
909 }
910
911 pub(crate) fn remove(
913 &mut self,
914 region_id: &RegionId,
915 ) -> (Vec<SenderWriteRequest>, Vec<SenderBulkRequest>) {
916 if let Some((size, write_reqs, bulk_reqs)) = self.requests.remove(region_id) {
917 self.estimated_size -= size;
918 (write_reqs, bulk_reqs)
919 } else {
920 (vec![], vec![])
921 }
922 }
923
924 pub(crate) fn stalled_count(&self) -> usize {
926 self.requests
927 .values()
928 .map(|(_, reqs, bulk_reqs)| reqs.len() + bulk_reqs.len())
929 .sum()
930 }
931}
932
933struct RegionWorkerLoop<S> {
935 id: WorkerId,
937 config: Arc<MitoConfig>,
939 regions: RegionMapRef,
941 dropping_regions: RegionMapRef,
943 opening_regions: OpeningRegionsRef,
945 catchup_regions: CatchupRegionsRef,
947 sender: Sender<WorkerRequestWithTime>,
949 receiver: Receiver<WorkerRequestWithTime>,
951 wal: Wal<S>,
953 object_store_manager: ObjectStoreManagerRef,
955 running: Arc<AtomicBool>,
957 memtable_builder_provider: MemtableBuilderProvider,
959 purge_scheduler: SchedulerRef,
961 write_buffer_manager: WriteBufferManagerRef,
963 index_build_scheduler: IndexBuildScheduler,
965 series_index_task_state: Option<Arc<SeriesIndexTaskState>>,
967 series_index_store: Option<ObjectStore>,
969 series_index_purger: Option<IndexFilePurger>,
970 flush_scheduler: FlushScheduler,
972 compaction_scheduler: CompactionScheduler,
974 stalled_requests: StalledRequests,
976 listener: WorkerListener,
978 cache_manager: CacheManagerRef,
980 puffin_manager_factory: PuffinManagerFactory,
982 intermediate_manager: IntermediateManager,
984 time_provider: TimeProviderRef,
986 last_periodical_check_millis: i64,
988 flush_sender: watch::Sender<()>,
990 flush_receiver: watch::Receiver<()>,
992 stalling_count: IntGauge,
994 region_count: IntGauge,
996 request_wait_time: Histogram,
998 region_edit_queues: RegionEditQueues,
1000 schema_metadata_manager: SchemaMetadataManagerRef,
1002 file_ref_manager: FileReferenceManagerRef,
1004 partition_expr_fetcher: PartitionExprFetcherRef,
1006 flush_semaphore: Arc<Semaphore>,
1008 plugins: Plugins,
1010}
1011
1012impl<S: LogStore> RegionWorkerLoop<S> {
1013 async fn run(&mut self) {
1015 let init_check_delay = worker_init_check_delay();
1016 info!(
1017 "Start region worker thread {}, init_check_delay: {:?}",
1018 self.id, init_check_delay
1019 );
1020 self.last_periodical_check_millis += init_check_delay.as_millis() as i64;
1021
1022 let mut write_req_buffer: Vec<SenderWriteRequest> =
1024 Vec::with_capacity(self.config.worker_request_batch_size);
1025 let mut bulk_req_buffer: Vec<SenderBulkRequest> =
1026 Vec::with_capacity(self.config.worker_request_batch_size);
1027 let mut ddl_req_buffer: Vec<SenderDdlRequest> =
1028 Vec::with_capacity(self.config.worker_request_batch_size);
1029 let mut general_req_buffer: Vec<WorkerRequest> =
1030 RequestBuffer::with_capacity(self.config.worker_request_batch_size);
1031
1032 while self.running.load(Ordering::Relaxed) {
1033 write_req_buffer.clear();
1035 ddl_req_buffer.clear();
1036 general_req_buffer.clear();
1037 let mut bulk_insert_req_num = 0;
1038
1039 let max_wait_time = self.time_provider.wait_duration(CHECK_REGION_INTERVAL);
1040 let sleep = tokio::time::sleep(max_wait_time);
1041 tokio::pin!(sleep);
1042
1043 tokio::select! {
1044 request_opt = self.receiver.recv() => {
1045 match request_opt {
1046 Some(request_with_time) => {
1047 let wait_time = request_with_time.created_at.elapsed();
1049 self.request_wait_time.observe(wait_time.as_secs_f64());
1050
1051 match request_with_time.request {
1052 WorkerRequest::Write(sender_req) => write_req_buffer.push(sender_req),
1053 WorkerRequest::Ddl(sender_req) => ddl_req_buffer.push(sender_req),
1054 WorkerRequest::BulkInserts(bulk_insert) => {
1055 bulk_insert_req_num += 1;
1056 self.buffer_bulk_insert_request(bulk_insert, &mut bulk_req_buffer)
1057 .await;
1058 }
1059 req => general_req_buffer.push(req),
1060 }
1061 },
1062 None => break,
1064 }
1065 }
1066 recv_res = self.flush_receiver.changed() => {
1067 if recv_res.is_err() {
1068 break;
1070 } else {
1071 self.maybe_flush_worker();
1076 self.handle_stalled_requests().await;
1078 continue;
1079 }
1080 }
1081 _ = &mut sleep => {
1082 self.handle_periodical_tasks();
1084 continue;
1085 }
1086 }
1087
1088 if self.flush_receiver.has_changed().unwrap_or(false) {
1089 self.handle_stalled_requests().await;
1093 }
1094
1095 for _ in 1..self.config.worker_request_batch_size {
1097 match self.receiver.try_recv() {
1099 Ok(request_with_time) => {
1100 let wait_time = request_with_time.created_at.elapsed();
1102 self.request_wait_time.observe(wait_time.as_secs_f64());
1103
1104 match request_with_time.request {
1105 WorkerRequest::Write(sender_req) => write_req_buffer.push(sender_req),
1106 WorkerRequest::Ddl(sender_req) => ddl_req_buffer.push(sender_req),
1107 WorkerRequest::BulkInserts(bulk_insert) => {
1108 bulk_insert_req_num += 1;
1109 self.buffer_bulk_insert_request(bulk_insert, &mut bulk_req_buffer)
1110 .await
1111 }
1112 req => general_req_buffer.push(req),
1113 }
1114 }
1115 Err(_) => break,
1117 }
1118 }
1119
1120 self.listener.on_recv_requests(
1121 write_req_buffer.len()
1122 + ddl_req_buffer.len()
1123 + general_req_buffer.len()
1124 + bulk_insert_req_num,
1125 );
1126
1127 self.handle_requests(
1128 &mut write_req_buffer,
1129 &mut ddl_req_buffer,
1130 &mut general_req_buffer,
1131 &mut bulk_req_buffer,
1132 )
1133 .await;
1134
1135 self.handle_periodical_tasks();
1136 }
1137
1138 self.clean().await;
1139
1140 info!("Exit region worker thread {}", self.id);
1141 }
1142
1143 async fn buffer_bulk_insert_request(
1144 &mut self,
1145 bulk_insert: BulkInsertRequest,
1146 bulk_requests: &mut Vec<SenderBulkRequest>,
1147 ) {
1148 let BulkInsertRequest {
1149 metadata,
1150 request,
1151 sender,
1152 } = bulk_insert;
1153
1154 if let Some(region_metadata) = metadata {
1155 self.handle_bulk_insert_batch(region_metadata, request, bulk_requests, sender)
1156 .await;
1157 } else {
1158 error!("Cannot find region metadata for {}", request.region_id);
1159 sender.send(
1160 error::RegionNotFoundSnafu {
1161 region_id: request.region_id,
1162 }
1163 .fail(),
1164 );
1165 }
1166 }
1167
1168 async fn handle_requests(
1172 &mut self,
1173 write_requests: &mut Vec<SenderWriteRequest>,
1174 ddl_requests: &mut Vec<SenderDdlRequest>,
1175 general_requests: &mut Vec<WorkerRequest>,
1176 bulk_requests: &mut Vec<SenderBulkRequest>,
1177 ) {
1178 for worker_req in general_requests.drain(..) {
1179 match worker_req {
1180 WorkerRequest::Write(_) | WorkerRequest::Ddl(_) => {
1181 continue;
1183 }
1184 WorkerRequest::BulkInserts(_) => unreachable!("bulk inserts are buffered"),
1185 WorkerRequest::Background { region_id, notify } => {
1186 if matches!(
1187 ¬ify,
1188 BackgroundNotify::RegionEdit(edit_result)
1189 if edit_result.update_region_state
1190 ) {
1191 self.handle_buffered_region_write_requests(
1200 ®ion_id,
1201 write_requests,
1202 bulk_requests,
1203 )
1204 .await;
1205 }
1206 self.handle_background_notify(region_id, notify).await;
1208 }
1209 WorkerRequest::SetRegionRoleStateGracefully {
1210 region_id,
1211 region_role_state,
1212 sender,
1213 } => {
1214 self.set_role_state_gracefully(region_id, region_role_state, sender)
1215 .await;
1216 }
1217 WorkerRequest::EditRegion(request) => {
1218 self.handle_region_edit(request);
1219 }
1220 WorkerRequest::Stop => {
1221 debug_assert!(!self.running.load(Ordering::Relaxed));
1222 }
1223 WorkerRequest::SyncRegion(req) => {
1224 self.handle_region_sync(req).await;
1225 }
1226 WorkerRequest::RemapManifests(req) => {
1227 self.handle_remap_manifests_request(req);
1228 }
1229 WorkerRequest::CopyRegionFrom(req) => {
1230 self.handle_copy_region_from_request(req);
1231 }
1232 }
1233 }
1234
1235 self.handle_write_requests(write_requests, bulk_requests, true)
1238 .await;
1239
1240 self.handle_ddl_requests(ddl_requests).await;
1241 }
1242
1243 async fn handle_ddl_requests(&mut self, ddl_requests: &mut Vec<SenderDdlRequest>) {
1245 if ddl_requests.is_empty() {
1246 return;
1247 }
1248
1249 for ddl in ddl_requests.drain(..) {
1250 let res = match ddl.request {
1251 DdlRequest::Create(req) => self.handle_create_request(ddl.region_id, req).await,
1252 DdlRequest::Drop(req) => {
1253 self.handle_drop_request(ddl.region_id, req, ddl.sender)
1254 .await;
1255 continue;
1256 }
1257 DdlRequest::Open((req, wal_entry_receiver)) => {
1258 self.handle_open_request(ddl.region_id, req, wal_entry_receiver, ddl.sender)
1259 .await;
1260 continue;
1261 }
1262 DdlRequest::OfflineCleanup(req) => {
1263 self.handle_offline_cleanup_request(ddl.region_id, req)
1264 .await
1265 }
1266 DdlRequest::Close(req) => {
1267 self.handle_close_request(ddl.region_id, req, ddl.sender)
1268 .await;
1269 continue;
1270 }
1271 DdlRequest::Alter(req) => {
1272 self.handle_alter_request(ddl.region_id, req, ddl.sender)
1273 .await;
1274 continue;
1275 }
1276 DdlRequest::Flush(req) => {
1277 self.handle_flush_request(ddl.region_id, req, ddl.sender);
1278 continue;
1279 }
1280 DdlRequest::Compact(req) => {
1281 self.handle_compaction_request(ddl.region_id, req, ddl.sender)
1282 .await;
1283 continue;
1284 }
1285 DdlRequest::BuildIndex(req) => {
1286 self.handle_build_index_request(ddl.region_id, req, ddl.sender)
1287 .await;
1288 continue;
1289 }
1290 DdlRequest::Truncate(req) => {
1291 self.handle_truncate_request(ddl.region_id, req, ddl.sender)
1292 .await;
1293 continue;
1294 }
1295 DdlRequest::Catchup((req, wal_entry_receiver)) => {
1296 self.handle_catchup_request(ddl.region_id, req, wal_entry_receiver, ddl.sender)
1297 .await;
1298 continue;
1299 }
1300 DdlRequest::EnterStaging(req) => {
1301 self.handle_enter_staging_request(
1302 ddl.region_id,
1303 req.partition_directive,
1304 ddl.sender,
1305 )
1306 .await;
1307 continue;
1308 }
1309 DdlRequest::ApplyStagingManifest(req) => {
1310 self.handle_apply_staging_manifest_request(ddl.region_id, req, ddl.sender)
1311 .await;
1312 continue;
1313 }
1314 };
1315
1316 ddl.sender.send(res);
1317 }
1318 }
1319
1320 fn handle_periodical_tasks(&mut self) {
1322 let interval = CHECK_REGION_INTERVAL.as_millis() as i64;
1323 if self
1324 .time_provider
1325 .elapsed_since(self.last_periodical_check_millis)
1326 < interval
1327 {
1328 return;
1329 }
1330
1331 self.last_periodical_check_millis = self.time_provider.current_time_millis();
1332
1333 if let Err(e) = self.flush_periodically() {
1334 error!(e; "Failed to flush regions periodically");
1335 }
1336 }
1337
1338 async fn handle_background_notify(&mut self, region_id: RegionId, notify: BackgroundNotify) {
1340 match notify {
1341 BackgroundNotify::CompactionPickFinished(req) => {
1342 self.handle_compaction_pick_finished(region_id, req).await
1343 }
1344 BackgroundNotify::FlushFinished(req) => {
1345 self.handle_flush_finished(region_id, req).await
1346 }
1347 BackgroundNotify::FlushFailed(req) => self.handle_flush_failed(region_id, req).await,
1348 BackgroundNotify::IndexBuildFinished(req) => {
1349 self.handle_index_build_finished(region_id, req).await
1350 }
1351 BackgroundNotify::IndexBuildStopped(req) => {
1352 self.handle_index_build_stopped(region_id, req).await
1353 }
1354 BackgroundNotify::IndexBuildFailed(req) => {
1355 self.handle_index_build_failed(region_id, req).await
1356 }
1357 BackgroundNotify::IndexBuildRetry(req) => {
1358 self.handle_rebuild_index(req, OptionOutputTx::new(None))
1359 .await
1360 }
1361 BackgroundNotify::CompactionFinished(req) => {
1362 self.handle_compaction_finished(region_id, req).await
1363 }
1364 BackgroundNotify::CompactionCancelled(req) => {
1365 self.handle_compaction_cancelled(region_id, req).await
1366 }
1367 BackgroundNotify::CompactionFailed(req) => self.handle_compaction_failure(req).await,
1368 BackgroundNotify::Truncate(req) => self.handle_truncate_result(req).await,
1369 BackgroundNotify::DiscardUnflushed(req) => {
1370 self.handle_discard_unflushed_result(req).await
1371 }
1372 BackgroundNotify::RegionChange(req) => {
1373 self.handle_manifest_region_change_result(req).await
1374 }
1375 BackgroundNotify::EnterStaging(req) => self.handle_enter_staging_result(req).await,
1376 BackgroundNotify::RegionEdit(req) => self.handle_region_edit_result(req).await,
1377 BackgroundNotify::CopyRegionFromFinished(req) => {
1378 self.handle_copy_region_from_finished(req)
1379 }
1380 }
1381 }
1382
1383 async fn set_role_state_gracefully(
1385 &mut self,
1386 region_id: RegionId,
1387 region_role_state: SettableRegionRoleState,
1388 sender: oneshot::Sender<SetRegionRoleStateResponse>,
1389 ) {
1390 if let Some(region) = self.regions.get_region(region_id) {
1391 common_runtime::spawn_global(async move {
1393 match region.set_role_state_gracefully(region_role_state).await {
1394 Ok(()) => {
1395 let last_entry_id = region.version_control.current().last_entry_id;
1396 let _ = sender.send(SetRegionRoleStateResponse::success(
1397 SetRegionRoleStateSuccess::mito(last_entry_id),
1398 ));
1399 }
1400 Err(e) => {
1401 error!(e; "Failed to set region {} role state to {:?}", region_id, region_role_state);
1402 let _ = sender.send(SetRegionRoleStateResponse::invalid_transition(
1403 BoxedError::new(e),
1404 ));
1405 }
1406 }
1407 });
1408 } else {
1409 let _ = sender.send(SetRegionRoleStateResponse::NotFound);
1410 }
1411 }
1412}
1413
1414impl<S> RegionWorkerLoop<S> {
1415 async fn clean(&self) {
1417 let regions = self.regions.list_regions();
1419 for region in regions {
1420 region.stop().await;
1421 }
1422
1423 self.regions.clear();
1424 }
1425
1426 fn notify_group(&mut self) {
1429 let _ = self.flush_sender.send(());
1431 self.flush_receiver.borrow_and_update();
1433 }
1434}
1435
1436#[derive(Default, Clone)]
1438pub(crate) struct WorkerListener {
1439 #[cfg(any(test, feature = "test"))]
1440 listener: Option<crate::engine::listener::EventListenerRef>,
1441}
1442
1443impl WorkerListener {
1444 #[cfg(any(test, feature = "test"))]
1445 pub(crate) fn new(
1446 listener: Option<crate::engine::listener::EventListenerRef>,
1447 ) -> WorkerListener {
1448 WorkerListener { listener }
1449 }
1450
1451 pub(crate) fn on_flush_success(&self, region_id: RegionId) {
1453 #[cfg(any(test, feature = "test"))]
1454 if let Some(listener) = &self.listener {
1455 listener.on_flush_success(region_id);
1456 }
1457 let _ = region_id;
1459 }
1460
1461 pub(crate) fn on_write_stall(&self) {
1463 #[cfg(any(test, feature = "test"))]
1464 if let Some(listener) = &self.listener {
1465 listener.on_write_stall();
1466 }
1467 }
1468
1469 pub(crate) async fn on_flush_begin(&self, region_id: RegionId) {
1470 #[cfg(any(test, feature = "test"))]
1471 if let Some(listener) = &self.listener {
1472 listener.on_flush_begin(region_id).await;
1473 }
1474 let _ = region_id;
1476 }
1477
1478 pub(crate) async fn on_flush_commit_begin(&self, _region_id: RegionId) {
1479 #[cfg(any(test, feature = "test"))]
1480 if let Some(listener) = &self.listener {
1481 listener.on_flush_commit_begin(_region_id).await;
1482 }
1483 }
1484
1485 pub(crate) fn on_flush_cancel_requested(&self, _region_id: RegionId) {
1486 #[cfg(any(test, feature = "test"))]
1487 if let Some(listener) = &self.listener {
1488 listener.on_flush_cancel_requested(_region_id);
1489 }
1490 }
1491
1492 pub(crate) fn on_later_drop_begin(&self, region_id: RegionId) -> Option<Duration> {
1493 #[cfg(any(test, feature = "test"))]
1494 if let Some(listener) = &self.listener {
1495 return listener.on_later_drop_begin(region_id);
1496 }
1497 let _ = region_id;
1499 None
1500 }
1501
1502 pub(crate) fn on_later_drop_end(&self, region_id: RegionId, removed: bool) {
1504 #[cfg(any(test, feature = "test"))]
1505 if let Some(listener) = &self.listener {
1506 listener.on_later_drop_end(region_id, removed);
1507 }
1508 let _ = region_id;
1510 let _ = removed;
1511 }
1512
1513 pub(crate) async fn on_merge_ssts_finished(&self, region_id: RegionId) {
1514 #[cfg(any(test, feature = "test"))]
1515 if let Some(listener) = &self.listener {
1516 listener.on_merge_ssts_finished(region_id).await;
1517 }
1518 let _ = region_id;
1520 }
1521
1522 pub(crate) fn on_recv_requests(&self, request_num: usize) {
1523 #[cfg(any(test, feature = "test"))]
1524 if let Some(listener) = &self.listener {
1525 listener.on_recv_requests(request_num);
1526 }
1527 let _ = request_num;
1529 }
1530
1531 pub(crate) fn on_file_cache_filled(&self, _file_id: FileId) {
1532 #[cfg(any(test, feature = "test"))]
1533 if let Some(listener) = &self.listener {
1534 listener.on_file_cache_filled(_file_id);
1535 }
1536 }
1537
1538 pub(crate) fn on_compaction_scheduled(&self, _region_id: RegionId) {
1539 #[cfg(any(test, feature = "test"))]
1540 if let Some(listener) = &self.listener {
1541 listener.on_compaction_scheduled(_region_id);
1542 }
1543 }
1544
1545 pub(crate) async fn on_compaction_pick_begin(&self, _region_id: RegionId) {
1546 #[cfg(any(test, feature = "test"))]
1547 if let Some(listener) = &self.listener {
1548 listener.on_compaction_pick_begin(_region_id).await;
1549 }
1550 }
1551
1552 pub(crate) async fn on_compaction_commit_begin(&self, _region_id: RegionId) {
1553 #[cfg(any(test, feature = "test"))]
1554 if let Some(listener) = &self.listener {
1555 listener.on_compaction_commit_begin(_region_id).await;
1556 }
1557 }
1558
1559 pub(crate) async fn on_compaction_result_notified(&self, _region_id: RegionId) {
1560 #[cfg(any(test, feature = "test"))]
1561 if let Some(listener) = &self.listener {
1562 listener.on_compaction_result_notified(_region_id).await;
1563 }
1564 }
1565
1566 pub(crate) fn on_compaction_cancel_requested(&self, _region_id: RegionId) {
1567 #[cfg(any(test, feature = "test"))]
1568 if let Some(listener) = &self.listener {
1569 listener.on_compaction_cancel_requested(_region_id);
1570 }
1571 }
1572
1573 pub(crate) async fn on_notify_region_change_result_begin(&self, _region_id: RegionId) {
1574 #[cfg(any(test, feature = "test"))]
1575 if let Some(listener) = &self.listener {
1576 listener
1577 .on_notify_region_change_result_begin(_region_id)
1578 .await;
1579 }
1580 }
1581
1582 pub(crate) async fn on_enter_staging_result_begin(&self, _region_id: RegionId) {
1583 #[cfg(any(test, feature = "test"))]
1584 if let Some(listener) = &self.listener {
1585 listener.on_enter_staging_result_begin(_region_id).await;
1586 }
1587 }
1588
1589 pub(crate) async fn on_index_build_finish(&self, _region_file_id: RegionFileId) {
1590 #[cfg(any(test, feature = "test"))]
1591 if let Some(listener) = &self.listener {
1592 listener.on_index_build_finish(_region_file_id).await;
1593 }
1594 }
1595
1596 pub(crate) async fn on_index_build_begin(&self, _region_file_id: RegionFileId) {
1597 #[cfg(any(test, feature = "test"))]
1598 if let Some(listener) = &self.listener {
1599 listener.on_index_build_begin(_region_file_id).await;
1600 }
1601 }
1602
1603 pub(crate) async fn on_index_build_abort(&self, _region_file_id: RegionFileId) {
1604 #[cfg(any(test, feature = "test"))]
1605 if let Some(listener) = &self.listener {
1606 listener.on_index_build_abort(_region_file_id).await;
1607 }
1608 }
1609
1610 pub(crate) async fn on_index_build_before_manifest_commit(
1611 &self,
1612 _region_file_id: RegionFileId,
1613 ) {
1614 #[cfg(any(test, feature = "test"))]
1615 if let Some(listener) = &self.listener {
1616 listener
1617 .on_index_build_before_manifest_commit(_region_file_id)
1618 .await;
1619 }
1620 }
1621
1622 pub(crate) async fn on_index_build_manifest_committed(&self, _region_file_id: RegionFileId) {
1623 #[cfg(any(test, feature = "test"))]
1624 if let Some(listener) = &self.listener {
1625 listener
1626 .on_index_build_manifest_committed(_region_file_id)
1627 .await;
1628 }
1629 }
1630}
1631
1632#[cfg(test)]
1633mod tests {
1634 use super::*;
1635 use crate::test_util::TestEnv;
1636
1637 #[test]
1638 fn test_region_id_to_index() {
1639 let num_workers = 4;
1640
1641 let region_id = RegionId::new(1, 2);
1642 let index = region_id_to_index(region_id, num_workers);
1643 assert_eq!(index, 3);
1644
1645 let region_id = RegionId::new(2, 3);
1646 let index = region_id_to_index(region_id, num_workers);
1647 assert_eq!(index, 1);
1648 }
1649
1650 #[tokio::test]
1651 async fn test_worker_group_start_stop() {
1652 for (enabled, existing_index) in [(false, false), (true, false), (false, true)] {
1653 let env = TestEnv::with_prefix("group-stop").await;
1655 let root = env.data_home().join("series_index");
1656 if existing_index {
1657 tokio::fs::create_dir_all(&root).await.unwrap();
1658 tokio::fs::write(root.join("retained"), b"index")
1659 .await
1660 .unwrap();
1661 }
1662 let group = env
1663 .create_worker_group(MitoConfig {
1664 num_workers: 4,
1665 experimental_enable_series_index: enabled,
1666 experimental_series_index_maintenance_interval: Duration::from_secs(3600),
1667 ..Default::default()
1668 })
1669 .await;
1670
1671 for worker in &group.workers {
1672 assert_eq!(worker.series_index_task_state.is_some(), enabled);
1673 assert_eq!(worker.series_index_handle.lock().await.is_some(), enabled);
1674 }
1675 assert_eq!(root.is_dir(), enabled || existing_index);
1676 if existing_index {
1677 assert_eq!(
1678 tokio::fs::read(root.join("retained")).await.unwrap(),
1679 b"index"
1680 );
1681 }
1682
1683 tokio::time::timeout(Duration::from_secs(5), group.stop())
1684 .await
1685 .expect("series-index tasks should stop without waiting for their interval")
1686 .unwrap();
1687 }
1688 }
1689}