Skip to main content

mito2/
worker.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Structs and utilities for writing regions.
16
17mod handle_alter;
18mod handle_apply_staging;
19mod handle_bulk_insert;
20mod handle_catchup;
21mod handle_close;
22mod handle_compaction;
23mod handle_copy_region;
24mod handle_create;
25mod handle_drop;
26mod handle_enter_staging;
27mod handle_flush;
28mod handle_manifest;
29mod handle_open;
30mod handle_rebuild_index;
31mod handle_remap;
32mod handle_truncate;
33mod handle_write;
34
35use std::collections::HashMap;
36use std::path::Path;
37use std::sync::Arc;
38use std::sync::atomic::{AtomicBool, Ordering};
39use std::time::Duration;
40
41use common_base::Plugins;
42use common_error::ext::BoxedError;
43use common_meta::key::SchemaMetadataManagerRef;
44use common_runtime::JoinHandle;
45use common_stat::get_total_memory_bytes;
46use common_telemetry::{error, info, warn};
47use futures::future::try_join_all;
48use object_store::manager::ObjectStoreManagerRef;
49use prometheus::{Histogram, IntGauge};
50use rand::{Rng, rng};
51use snafu::{ResultExt, ensure};
52use store_api::logstore::LogStore;
53use store_api::region_engine::{
54    SetRegionRoleStateResponse, SetRegionRoleStateSuccess, SettableRegionRoleState,
55};
56use store_api::storage::{FileId, RegionId};
57use tokio::sync::mpsc::{Receiver, Sender};
58use tokio::sync::{Mutex, Semaphore, mpsc, oneshot, watch};
59
60use crate::cache::write_cache::{WriteCache, WriteCacheRef};
61use crate::cache::{CacheManager, CacheManagerRef};
62use crate::compaction::CompactionScheduler;
63use crate::compaction::memory_manager::{CompactionMemoryManager, new_compaction_memory_manager};
64use crate::config::MitoConfig;
65use crate::error::{self, CreateDirSnafu, JoinSnafu, Result, WorkerStoppedSnafu};
66use crate::flush::{FlushScheduler, WriteBufferManagerImpl, WriteBufferManagerRef};
67use crate::gc::{GcLimiter, GcLimiterRef};
68use crate::memtable::MemtableBuilderProvider;
69use crate::metrics::{REGION_COUNT, REQUEST_WAIT_TIME, WRITE_STALLING};
70use crate::region::opener::PartitionExprFetcherRef;
71use crate::region::{
72    CatchupRegions, CatchupRegionsRef, MitoRegionRef, OpeningRegions, OpeningRegionsRef, RegionMap,
73    RegionMapRef,
74};
75use crate::request::{
76    BackgroundNotify, BulkInsertRequest, DdlRequest, OptionOutputTx, SenderBulkRequest,
77    SenderDdlRequest, SenderWriteRequest, WorkerRequest, WorkerRequestWithTime,
78};
79use crate::schedule::scheduler::{LocalScheduler, SchedulerRef};
80use crate::sst::file::RegionFileId;
81use crate::sst::file_ref::FileReferenceManagerRef;
82use crate::sst::index::IndexBuildScheduler;
83use crate::sst::index::intermediate::IntermediateManager;
84use crate::sst::index::puffin_manager::PuffinManagerFactory;
85use crate::time_provider::{StdTimeProvider, TimeProviderRef};
86use crate::wal::Wal;
87use crate::worker::handle_manifest::RegionEditQueues;
88
89/// Identifier for a worker.
90pub(crate) type WorkerId = u32;
91
92pub(crate) const DROPPING_MARKER_FILE: &str = ".dropping";
93
94/// Interval to check whether regions should flush.
95pub(crate) const CHECK_REGION_INTERVAL: Duration = Duration::from_secs(60);
96/// Max delay to check region periodical tasks.
97pub(crate) const MAX_INITIAL_CHECK_DELAY_SECS: u64 = 60 * 3;
98
99#[cfg_attr(doc, aquamarine::aquamarine)]
100/// A fixed size group of [RegionWorkers](RegionWorker).
101///
102/// A worker group binds each region to a specific [RegionWorker] and sends
103/// requests to region's dedicated worker.
104///
105/// ```mermaid
106/// graph LR
107///
108/// RegionRequest -- Route by region id --> Worker0 & Worker1
109///
110/// subgraph MitoEngine
111///     subgraph WorkerGroup
112///         Worker0["RegionWorker 0"]
113///         Worker1["RegionWorker 1"]
114///     end
115/// end
116///
117/// Chan0[" Request channel 0"]
118/// Chan1[" Request channel 1"]
119/// WorkerThread1["RegionWorkerLoop 1"]
120///
121/// subgraph WorkerThread0["RegionWorkerLoop 0"]
122///     subgraph RegionMap["RegionMap (regions bound to worker 0)"]
123///         Region0["Region 0"]
124///         Region2["Region 2"]
125///     end
126///     Buffer0["RequestBuffer"]
127///
128///     Buffer0 -- modify regions --> RegionMap
129/// end
130///
131/// Worker0 --> Chan0
132/// Worker1 --> Chan1
133/// Chan0 --> Buffer0
134/// Chan1 --> WorkerThread1
135/// ```
136pub(crate) struct WorkerGroup {
137    /// Workers of the group.
138    workers: Vec<RegionWorker>,
139    /// Flush background job pool.
140    flush_job_pool: SchedulerRef,
141    /// Compaction background job pool.
142    compact_job_pool: SchedulerRef,
143    /// Scheduler for index build jobs.
144    index_build_job_pool: SchedulerRef,
145    /// Scheduler for file purgers.
146    purge_scheduler: SchedulerRef,
147    /// Cache.
148    cache_manager: CacheManagerRef,
149    /// File reference manager.
150    file_ref_manager: FileReferenceManagerRef,
151    /// Gc limiter to limit concurrent gc jobs.
152    gc_limiter: GcLimiterRef,
153    /// Object store manager.
154    object_store_manager: ObjectStoreManagerRef,
155    /// Puffin manager factory.
156    puffin_manager_factory: PuffinManagerFactory,
157    /// Intermediate manager.
158    intermediate_manager: IntermediateManager,
159    /// Schema metadata manager.
160    schema_metadata_manager: SchemaMetadataManagerRef,
161}
162
163impl WorkerGroup {
164    /// Starts a worker group.
165    ///
166    /// The number of workers should be power of two.
167    pub(crate) async fn start<S: LogStore>(
168        config: Arc<MitoConfig>,
169        log_store: Arc<S>,
170        object_store_manager: ObjectStoreManagerRef,
171        schema_metadata_manager: SchemaMetadataManagerRef,
172        file_ref_manager: FileReferenceManagerRef,
173        partition_expr_fetcher: PartitionExprFetcherRef,
174        plugins: Plugins,
175    ) -> Result<WorkerGroup> {
176        let (flush_sender, flush_receiver) = watch::channel(());
177        let write_buffer_manager = Arc::new(
178            WriteBufferManagerImpl::new(config.global_write_buffer_size.as_bytes() as usize)
179                .with_notifier(flush_sender.clone()),
180        );
181        let puffin_manager_factory = PuffinManagerFactory::new(
182            &config.index.aux_path,
183            config.index.staging_size.as_bytes(),
184            Some(config.index.write_buffer_size.as_bytes() as _),
185            config.index.staging_ttl,
186        )
187        .await?;
188        let intermediate_manager = IntermediateManager::init_fs(&config.index.aux_path)
189            .await?
190            .with_buffer_size(Some(config.index.write_buffer_size.as_bytes() as _));
191        let index_build_job_pool =
192            Arc::new(LocalScheduler::new(config.max_background_index_builds));
193        let flush_job_pool = Arc::new(LocalScheduler::new(config.max_background_flushes));
194        let compact_job_pool = Arc::new(LocalScheduler::new(config.max_background_compactions));
195        let flush_semaphore = Arc::new(Semaphore::new(config.max_background_flushes));
196        // We use another scheduler to avoid purge jobs blocking other jobs.
197        let purge_scheduler = Arc::new(LocalScheduler::new(config.max_background_purges));
198        let write_cache = write_cache_from_config(
199            &config,
200            puffin_manager_factory.clone(),
201            intermediate_manager.clone(),
202        )
203        .await?;
204        let cache_manager = Arc::new(
205            CacheManager::builder()
206                .sst_meta_cache_size(config.sst_meta_cache_size.as_bytes())
207                .vector_cache_size(config.vector_cache_size.as_bytes())
208                .page_cache_size(config.page_cache_size.as_bytes())
209                .selector_result_cache_size(config.selector_result_cache_size.as_bytes())
210                .range_result_cache_size(config.range_result_cache_size.as_bytes())
211                .prefilter_result_cache_size(config.prefilter_result_cache_size.as_bytes())
212                .index_metadata_size(config.index.metadata_cache_size.as_bytes())
213                .index_content_size(config.index.content_cache_size.as_bytes())
214                .index_content_page_size(config.index.content_cache_page_size.as_bytes())
215                .index_result_cache_size(config.index.result_cache_size.as_bytes())
216                .puffin_metadata_size(config.index.metadata_cache_size.as_bytes())
217                .write_cache(write_cache)
218                .build(),
219        );
220        let time_provider = Arc::new(StdTimeProvider);
221        let total_memory = get_total_memory_bytes();
222        let total_memory = if total_memory > 0 {
223            total_memory as u64
224        } else {
225            0
226        };
227        let compaction_limit_bytes = config
228            .experimental_compaction_memory_limit
229            .resolve(total_memory);
230        let compaction_memory_manager =
231            Arc::new(new_compaction_memory_manager(compaction_limit_bytes));
232        let gc_limiter = Arc::new(GcLimiter::new(config.gc.max_concurrent_gc_job));
233
234        let workers = (0..config.num_workers)
235            .map(|id| {
236                WorkerStarter {
237                    id: id as WorkerId,
238                    config: config.clone(),
239                    log_store: log_store.clone(),
240                    object_store_manager: object_store_manager.clone(),
241                    write_buffer_manager: write_buffer_manager.clone(),
242                    index_build_job_pool: index_build_job_pool.clone(),
243                    flush_job_pool: flush_job_pool.clone(),
244                    compact_job_pool: compact_job_pool.clone(),
245                    purge_scheduler: purge_scheduler.clone(),
246                    listener: WorkerListener::default(),
247                    cache_manager: cache_manager.clone(),
248                    compaction_memory_manager: compaction_memory_manager.clone(),
249                    puffin_manager_factory: puffin_manager_factory.clone(),
250                    intermediate_manager: intermediate_manager.clone(),
251                    time_provider: time_provider.clone(),
252                    flush_sender: flush_sender.clone(),
253                    flush_receiver: flush_receiver.clone(),
254                    plugins: plugins.clone(),
255                    schema_metadata_manager: schema_metadata_manager.clone(),
256                    file_ref_manager: file_ref_manager.clone(),
257                    partition_expr_fetcher: partition_expr_fetcher.clone(),
258                    flush_semaphore: flush_semaphore.clone(),
259                }
260                .start()
261            })
262            .collect::<Result<Vec<_>>>()?;
263
264        Ok(WorkerGroup {
265            workers,
266            flush_job_pool,
267            compact_job_pool,
268            index_build_job_pool,
269            purge_scheduler,
270            cache_manager,
271            file_ref_manager,
272            gc_limiter,
273            object_store_manager,
274            puffin_manager_factory,
275            intermediate_manager,
276            schema_metadata_manager,
277        })
278    }
279
280    /// Stops the worker group.
281    pub(crate) async fn stop(&self) -> Result<()> {
282        info!("Stop region worker group");
283
284        // TODO(yingwen): Do we need to stop gracefully?
285        // Stops the scheduler gracefully.
286        self.compact_job_pool.stop(true).await?;
287        // Stops the scheduler gracefully.
288        self.flush_job_pool.stop(true).await?;
289        // Stops the purge scheduler gracefully.
290        self.purge_scheduler.stop(true).await?;
291        // Stops the index build job pool gracefully.
292        self.index_build_job_pool.stop(true).await?;
293
294        try_join_all(self.workers.iter().map(|worker| worker.stop())).await?;
295
296        Ok(())
297    }
298
299    /// Submits a request to a worker in the group.
300    pub(crate) async fn submit_to_worker(
301        &self,
302        region_id: RegionId,
303        request: WorkerRequest,
304    ) -> Result<()> {
305        self.worker(region_id).submit_request(request).await
306    }
307
308    /// Returns true if the specific region exists.
309    pub(crate) fn is_region_exists(&self, region_id: RegionId) -> bool {
310        self.worker(region_id).is_region_exists(region_id)
311    }
312
313    /// Returns true if the specific region is opening.
314    pub(crate) fn is_region_opening(&self, region_id: RegionId) -> bool {
315        self.worker(region_id).is_region_opening(region_id)
316    }
317
318    /// Returns true if the specific region is catching up.
319    pub(crate) fn is_region_catching_up(&self, region_id: RegionId) -> bool {
320        self.worker(region_id).is_region_catching_up(region_id)
321    }
322
323    /// Returns region of specific `region_id`.
324    ///
325    /// This method should not be public.
326    pub(crate) fn get_region(&self, region_id: RegionId) -> Option<MitoRegionRef> {
327        self.worker(region_id).get_region(region_id)
328    }
329
330    /// Returns cache of the group.
331    pub(crate) fn cache_manager(&self) -> CacheManagerRef {
332        self.cache_manager.clone()
333    }
334
335    pub(crate) fn file_ref_manager(&self) -> FileReferenceManagerRef {
336        self.file_ref_manager.clone()
337    }
338
339    pub(crate) fn gc_limiter(&self) -> GcLimiterRef {
340        self.gc_limiter.clone()
341    }
342
343    /// Get worker for specific `region_id`.
344    pub(crate) fn worker(&self, region_id: RegionId) -> &RegionWorker {
345        let index = region_id_to_index(region_id, self.workers.len());
346
347        &self.workers[index]
348    }
349
350    pub(crate) fn all_regions(&self) -> impl Iterator<Item = MitoRegionRef> + use<'_> {
351        self.workers
352            .iter()
353            .flat_map(|worker| worker.regions.list_regions())
354    }
355
356    pub(crate) fn object_store_manager(&self) -> &ObjectStoreManagerRef {
357        &self.object_store_manager
358    }
359
360    pub(crate) fn puffin_manager_factory(&self) -> &PuffinManagerFactory {
361        &self.puffin_manager_factory
362    }
363
364    pub(crate) fn intermediate_manager(&self) -> &IntermediateManager {
365        &self.intermediate_manager
366    }
367
368    pub(crate) fn schema_metadata_manager(&self) -> &SchemaMetadataManagerRef {
369        &self.schema_metadata_manager
370    }
371}
372
373// Tests methods.
374#[cfg(any(test, feature = "test"))]
375impl WorkerGroup {
376    /// Starts a worker group with `write_buffer_manager` and `listener` for tests.
377    ///
378    /// The number of workers should be power of two.
379    #[allow(clippy::too_many_arguments)]
380    pub(crate) async fn start_for_test<S: LogStore>(
381        config: Arc<MitoConfig>,
382        log_store: Arc<S>,
383        object_store_manager: ObjectStoreManagerRef,
384        write_buffer_manager: Option<WriteBufferManagerRef>,
385        listener: Option<crate::engine::listener::EventListenerRef>,
386        schema_metadata_manager: SchemaMetadataManagerRef,
387        file_ref_manager: FileReferenceManagerRef,
388        time_provider: TimeProviderRef,
389        partition_expr_fetcher: PartitionExprFetcherRef,
390    ) -> Result<WorkerGroup> {
391        let (flush_sender, flush_receiver) = watch::channel(());
392        let write_buffer_manager = write_buffer_manager.unwrap_or_else(|| {
393            Arc::new(
394                WriteBufferManagerImpl::new(config.global_write_buffer_size.as_bytes() as usize)
395                    .with_notifier(flush_sender.clone()),
396            )
397        });
398        let index_build_job_pool =
399            Arc::new(LocalScheduler::new(config.max_background_index_builds));
400        let flush_job_pool = Arc::new(LocalScheduler::new(config.max_background_flushes));
401        let compact_job_pool = Arc::new(LocalScheduler::new(config.max_background_compactions));
402        let flush_semaphore = Arc::new(Semaphore::new(config.max_background_flushes));
403        let purge_scheduler = Arc::new(LocalScheduler::new(config.max_background_flushes));
404        let puffin_manager_factory = PuffinManagerFactory::new(
405            &config.index.aux_path,
406            config.index.staging_size.as_bytes(),
407            Some(config.index.write_buffer_size.as_bytes() as _),
408            config.index.staging_ttl,
409        )
410        .await?;
411        let intermediate_manager = IntermediateManager::init_fs(&config.index.aux_path)
412            .await?
413            .with_buffer_size(Some(config.index.write_buffer_size.as_bytes() as _));
414        let write_cache = write_cache_from_config(
415            &config,
416            puffin_manager_factory.clone(),
417            intermediate_manager.clone(),
418        )
419        .await?;
420        let cache_manager = Arc::new(
421            CacheManager::builder()
422                .sst_meta_cache_size(config.sst_meta_cache_size.as_bytes())
423                .vector_cache_size(config.vector_cache_size.as_bytes())
424                .page_cache_size(config.page_cache_size.as_bytes())
425                .selector_result_cache_size(config.selector_result_cache_size.as_bytes())
426                .range_result_cache_size(config.range_result_cache_size.as_bytes())
427                .prefilter_result_cache_size(config.prefilter_result_cache_size.as_bytes())
428                .write_cache(write_cache)
429                .build(),
430        );
431        let total_memory = get_total_memory_bytes();
432        let total_memory = if total_memory > 0 {
433            total_memory as u64
434        } else {
435            0
436        };
437        let compaction_limit_bytes = config
438            .experimental_compaction_memory_limit
439            .resolve(total_memory);
440        let compaction_memory_manager =
441            Arc::new(new_compaction_memory_manager(compaction_limit_bytes));
442        let gc_limiter = Arc::new(GcLimiter::new(config.gc.max_concurrent_gc_job));
443        let workers = (0..config.num_workers)
444            .map(|id| {
445                WorkerStarter {
446                    id: id as WorkerId,
447                    config: config.clone(),
448                    log_store: log_store.clone(),
449                    object_store_manager: object_store_manager.clone(),
450                    write_buffer_manager: write_buffer_manager.clone(),
451                    index_build_job_pool: index_build_job_pool.clone(),
452                    flush_job_pool: flush_job_pool.clone(),
453                    compact_job_pool: compact_job_pool.clone(),
454                    purge_scheduler: purge_scheduler.clone(),
455                    listener: WorkerListener::new(listener.clone()),
456                    cache_manager: cache_manager.clone(),
457                    compaction_memory_manager: compaction_memory_manager.clone(),
458                    puffin_manager_factory: puffin_manager_factory.clone(),
459                    intermediate_manager: intermediate_manager.clone(),
460                    time_provider: time_provider.clone(),
461                    flush_sender: flush_sender.clone(),
462                    flush_receiver: flush_receiver.clone(),
463                    plugins: Plugins::new(),
464                    schema_metadata_manager: schema_metadata_manager.clone(),
465                    file_ref_manager: file_ref_manager.clone(),
466                    partition_expr_fetcher: partition_expr_fetcher.clone(),
467                    flush_semaphore: flush_semaphore.clone(),
468                }
469                .start()
470            })
471            .collect::<Result<Vec<_>>>()?;
472
473        Ok(WorkerGroup {
474            workers,
475            flush_job_pool,
476            compact_job_pool,
477            index_build_job_pool,
478            purge_scheduler,
479            cache_manager,
480            file_ref_manager,
481            gc_limiter,
482            object_store_manager,
483            puffin_manager_factory,
484            intermediate_manager,
485            schema_metadata_manager,
486        })
487    }
488
489    /// Returns the purge scheduler.
490    pub(crate) fn purge_scheduler(&self) -> &SchedulerRef {
491        &self.purge_scheduler
492    }
493}
494
495fn region_id_to_index(id: RegionId, num_workers: usize) -> usize {
496    ((id.table_id() as usize % num_workers) + (id.region_number() as usize % num_workers))
497        % num_workers
498}
499
500pub async fn write_cache_from_config(
501    config: &MitoConfig,
502    puffin_manager_factory: PuffinManagerFactory,
503    intermediate_manager: IntermediateManager,
504) -> Result<Option<WriteCacheRef>> {
505    if !config.enable_write_cache {
506        return Ok(None);
507    }
508
509    tokio::fs::create_dir_all(Path::new(&config.write_cache_path))
510        .await
511        .context(CreateDirSnafu {
512            dir: &config.write_cache_path,
513        })?;
514
515    let cache = WriteCache::new_fs(
516        &config.write_cache_path,
517        config.write_cache_size,
518        config.write_cache_ttl,
519        Some(config.index_cache_percent),
520        config.enable_refill_cache_on_read,
521        puffin_manager_factory,
522        intermediate_manager,
523        config.manifest_cache_size,
524    )
525    .await?;
526    Ok(Some(Arc::new(cache)))
527}
528
529/// Computes a initial check delay for a worker.
530pub(crate) fn worker_init_check_delay() -> Duration {
531    let init_check_delay = rng().random_range(0..MAX_INITIAL_CHECK_DELAY_SECS);
532    Duration::from_secs(init_check_delay)
533}
534
535/// Worker start config.
536struct WorkerStarter<S> {
537    id: WorkerId,
538    config: Arc<MitoConfig>,
539    log_store: Arc<S>,
540    object_store_manager: ObjectStoreManagerRef,
541    write_buffer_manager: WriteBufferManagerRef,
542    compact_job_pool: SchedulerRef,
543    index_build_job_pool: SchedulerRef,
544    flush_job_pool: SchedulerRef,
545    purge_scheduler: SchedulerRef,
546    listener: WorkerListener,
547    cache_manager: CacheManagerRef,
548    compaction_memory_manager: Arc<CompactionMemoryManager>,
549    puffin_manager_factory: PuffinManagerFactory,
550    intermediate_manager: IntermediateManager,
551    time_provider: TimeProviderRef,
552    /// Watch channel sender to notify workers to handle stalled requests.
553    flush_sender: watch::Sender<()>,
554    /// Watch channel receiver to wait for background flush job.
555    flush_receiver: watch::Receiver<()>,
556    plugins: Plugins,
557    schema_metadata_manager: SchemaMetadataManagerRef,
558    file_ref_manager: FileReferenceManagerRef,
559    partition_expr_fetcher: PartitionExprFetcherRef,
560    flush_semaphore: Arc<Semaphore>,
561}
562
563impl<S: LogStore> WorkerStarter<S> {
564    /// Starts a region worker and its background thread.
565    fn start(self) -> Result<RegionWorker> {
566        let regions = Arc::new(RegionMap::default());
567        let opening_regions = Arc::new(OpeningRegions::default());
568        let catchup_regions = Arc::new(CatchupRegions::default());
569        let (sender, receiver) = mpsc::channel(self.config.worker_channel_size);
570
571        let running = Arc::new(AtomicBool::new(true));
572        let now = self.time_provider.current_time_millis();
573        let id_string = self.id.to_string();
574        let mut worker_thread = RegionWorkerLoop {
575            id: self.id,
576            config: self.config.clone(),
577            regions: regions.clone(),
578            catchup_regions: catchup_regions.clone(),
579            dropping_regions: Arc::new(RegionMap::default()),
580            opening_regions: opening_regions.clone(),
581            sender: sender.clone(),
582            receiver,
583            wal: Wal::new(self.log_store),
584            object_store_manager: self.object_store_manager.clone(),
585            running: running.clone(),
586            memtable_builder_provider: MemtableBuilderProvider::new(
587                Some(self.write_buffer_manager.clone()),
588                self.config.clone(),
589            ),
590            purge_scheduler: self.purge_scheduler.clone(),
591            write_buffer_manager: self.write_buffer_manager,
592            index_build_scheduler: IndexBuildScheduler::new(
593                self.index_build_job_pool,
594                self.config.max_background_index_builds,
595            ),
596            flush_scheduler: FlushScheduler::new(self.flush_job_pool),
597            compaction_scheduler: CompactionScheduler::new(
598                self.compact_job_pool,
599                sender.clone(),
600                self.cache_manager.clone(),
601                self.config.clone(),
602                self.listener.clone(),
603                self.plugins.clone(),
604                self.compaction_memory_manager.clone(),
605                self.config.experimental_compaction_on_exhausted,
606            ),
607            stalled_requests: StalledRequests::default(),
608            listener: self.listener,
609            cache_manager: self.cache_manager,
610            puffin_manager_factory: self.puffin_manager_factory,
611            intermediate_manager: self.intermediate_manager,
612            time_provider: self.time_provider,
613            last_periodical_check_millis: now,
614            flush_sender: self.flush_sender,
615            flush_receiver: self.flush_receiver,
616            stalling_count: WRITE_STALLING.with_label_values(&[&id_string]),
617            region_count: REGION_COUNT.with_label_values(&[&id_string]),
618            request_wait_time: REQUEST_WAIT_TIME.with_label_values(&[&id_string]),
619            region_edit_queues: RegionEditQueues::default(),
620            schema_metadata_manager: self.schema_metadata_manager,
621            file_ref_manager: self.file_ref_manager.clone(),
622            partition_expr_fetcher: self.partition_expr_fetcher,
623            flush_semaphore: self.flush_semaphore,
624            plugins: self.plugins,
625        };
626        let handle = common_runtime::spawn_global(async move {
627            worker_thread.run().await;
628        });
629
630        Ok(RegionWorker {
631            id: self.id,
632            regions,
633            opening_regions,
634            catchup_regions,
635            sender,
636            handle: Mutex::new(Some(handle)),
637            running,
638        })
639    }
640}
641
642/// Worker to write and alter regions bound to it.
643pub(crate) struct RegionWorker {
644    /// Id of the worker.
645    id: WorkerId,
646    /// Regions bound to the worker.
647    regions: RegionMapRef,
648    /// The opening regions.
649    opening_regions: OpeningRegionsRef,
650    /// The catching up regions.
651    catchup_regions: CatchupRegionsRef,
652    /// Request sender.
653    sender: Sender<WorkerRequestWithTime>,
654    /// Handle to the worker thread.
655    handle: Mutex<Option<JoinHandle<()>>>,
656    /// Whether to run the worker thread.
657    running: Arc<AtomicBool>,
658}
659
660impl RegionWorker {
661    /// Submits request to background worker thread.
662    async fn submit_request(&self, request: WorkerRequest) -> Result<()> {
663        ensure!(self.is_running(), WorkerStoppedSnafu { id: self.id });
664        let request_with_time = WorkerRequestWithTime::new(request);
665        if self.sender.send(request_with_time).await.is_err() {
666            warn!(
667                "Worker {} is already exited but the running flag is still true",
668                self.id
669            );
670            // Manually set the running flag to false to avoid printing more warning logs.
671            self.set_running(false);
672            return WorkerStoppedSnafu { id: self.id }.fail();
673        }
674
675        Ok(())
676    }
677
678    /// Stop the worker.
679    ///
680    /// This method waits until the worker thread exists.
681    async fn stop(&self) -> Result<()> {
682        let handle = self.handle.lock().await.take();
683        if let Some(handle) = handle {
684            info!("Stop region worker {}", self.id);
685
686            self.set_running(false);
687            if self
688                .sender
689                .send(WorkerRequestWithTime::new(WorkerRequest::Stop))
690                .await
691                .is_err()
692            {
693                warn!("Worker {} is already exited before stop", self.id);
694            }
695
696            handle.await.context(JoinSnafu)?;
697        }
698
699        Ok(())
700    }
701
702    /// Returns true if the worker is still running.
703    fn is_running(&self) -> bool {
704        self.running.load(Ordering::Relaxed)
705    }
706
707    /// Sets whether the worker is still running.
708    fn set_running(&self, value: bool) {
709        self.running.store(value, Ordering::Relaxed)
710    }
711
712    /// Returns true if the worker contains specific region.
713    fn is_region_exists(&self, region_id: RegionId) -> bool {
714        self.regions.is_region_exists(region_id)
715    }
716
717    /// Returns true if the region is opening.
718    fn is_region_opening(&self, region_id: RegionId) -> bool {
719        self.opening_regions.is_region_exists(region_id)
720    }
721
722    /// Returns true if the region is catching up.
723    fn is_region_catching_up(&self, region_id: RegionId) -> bool {
724        self.catchup_regions.is_region_exists(region_id)
725    }
726
727    /// Returns region of specific `region_id`.
728    fn get_region(&self, region_id: RegionId) -> Option<MitoRegionRef> {
729        self.regions.get_region(region_id)
730    }
731
732    #[cfg(test)]
733    /// Returns the [OpeningRegionsRef].
734    pub(crate) fn opening_regions(&self) -> &OpeningRegionsRef {
735        &self.opening_regions
736    }
737
738    #[cfg(test)]
739    /// Returns the [CatchupRegionsRef].
740    pub(crate) fn catchup_regions(&self) -> &CatchupRegionsRef {
741        &self.catchup_regions
742    }
743}
744
745impl Drop for RegionWorker {
746    fn drop(&mut self) {
747        if self.is_running() {
748            self.set_running(false);
749            // Once we drop the sender, the worker thread will receive a disconnected error.
750        }
751    }
752}
753
754type RequestBuffer = Vec<WorkerRequest>;
755
756/// Buffer for stalled write requests.
757///
758/// Maintains stalled write requests and their estimated size.
759#[derive(Default)]
760pub(crate) struct StalledRequests {
761    /// Stalled requests.
762    /// Remember to use `StalledRequests::stalled_count()` to get the total number of stalled requests
763    /// instead of `StalledRequests::requests.len()`.
764    ///
765    /// Key: RegionId
766    /// Value: (estimated size, stalled requests)
767    pub(crate) requests:
768        HashMap<RegionId, (usize, Vec<SenderWriteRequest>, Vec<SenderBulkRequest>)>,
769    /// Estimated size of all stalled requests.
770    pub(crate) estimated_size: usize,
771}
772
773impl StalledRequests {
774    /// Returns the estimated size of stalled requests for a region.
775    pub(crate) fn estimated_size(&self, region_id: &RegionId) -> usize {
776        self.requests
777            .get(region_id)
778            .map(|(size, _, _)| *size)
779            .unwrap_or_default()
780    }
781
782    /// Appends stalled requests.
783    pub(crate) fn append(
784        &mut self,
785        requests: &mut Vec<SenderWriteRequest>,
786        bulk_requests: &mut Vec<SenderBulkRequest>,
787    ) {
788        for req in requests.drain(..) {
789            self.push(req);
790        }
791        for req in bulk_requests.drain(..) {
792            self.push_bulk(req);
793        }
794    }
795
796    /// Pushes a stalled request to the buffer.
797    pub(crate) fn push(&mut self, req: SenderWriteRequest) {
798        let (size, requests, _) = self.requests.entry(req.request.region_id).or_default();
799        let req_size = req.request.estimated_size();
800        *size += req_size;
801        self.estimated_size += req_size;
802        requests.push(req);
803    }
804
805    pub(crate) fn push_bulk(&mut self, req: SenderBulkRequest) {
806        let region_id = req.region_id;
807        let (size, _, requests) = self.requests.entry(region_id).or_default();
808        let req_size = req.request.estimated_size();
809        *size += req_size;
810        self.estimated_size += req_size;
811        requests.push(req);
812    }
813
814    /// Removes stalled requests of specific region.
815    pub(crate) fn remove(
816        &mut self,
817        region_id: &RegionId,
818    ) -> (Vec<SenderWriteRequest>, Vec<SenderBulkRequest>) {
819        if let Some((size, write_reqs, bulk_reqs)) = self.requests.remove(region_id) {
820            self.estimated_size -= size;
821            (write_reqs, bulk_reqs)
822        } else {
823            (vec![], vec![])
824        }
825    }
826
827    /// Returns the total number of all stalled requests.
828    pub(crate) fn stalled_count(&self) -> usize {
829        self.requests
830            .values()
831            .map(|(_, reqs, bulk_reqs)| reqs.len() + bulk_reqs.len())
832            .sum()
833    }
834}
835
836/// Background worker loop to handle requests.
837struct RegionWorkerLoop<S> {
838    /// Id of the worker.
839    id: WorkerId,
840    /// Engine config.
841    config: Arc<MitoConfig>,
842    /// Regions bound to the worker.
843    regions: RegionMapRef,
844    /// Regions that are not yet fully dropped.
845    dropping_regions: RegionMapRef,
846    /// Regions that are opening.
847    opening_regions: OpeningRegionsRef,
848    /// Regions that are catching up.
849    catchup_regions: CatchupRegionsRef,
850    /// Request sender.
851    sender: Sender<WorkerRequestWithTime>,
852    /// Request receiver.
853    receiver: Receiver<WorkerRequestWithTime>,
854    /// WAL of the engine.
855    wal: Wal<S>,
856    /// Manages object stores for manifest and SSTs.
857    object_store_manager: ObjectStoreManagerRef,
858    /// Whether the worker thread is still running.
859    running: Arc<AtomicBool>,
860    /// Memtable builder provider for each region.
861    memtable_builder_provider: MemtableBuilderProvider,
862    /// Background purge job scheduler.
863    purge_scheduler: SchedulerRef,
864    /// Engine write buffer manager.
865    write_buffer_manager: WriteBufferManagerRef,
866    /// Scheduler for index build task.
867    index_build_scheduler: IndexBuildScheduler,
868    /// Schedules background flush requests.
869    flush_scheduler: FlushScheduler,
870    /// Scheduler for compaction tasks.
871    compaction_scheduler: CompactionScheduler,
872    /// Stalled write requests.
873    stalled_requests: StalledRequests,
874    /// Event listener for tests.
875    listener: WorkerListener,
876    /// Cache.
877    cache_manager: CacheManagerRef,
878    /// Puffin manager factory for index.
879    puffin_manager_factory: PuffinManagerFactory,
880    /// Intermediate manager for inverted index.
881    intermediate_manager: IntermediateManager,
882    /// Provider to get current time.
883    time_provider: TimeProviderRef,
884    /// Last time to check regions periodically.
885    last_periodical_check_millis: i64,
886    /// Watch channel sender to notify workers to handle stalled requests.
887    flush_sender: watch::Sender<()>,
888    /// Watch channel receiver to wait for background flush job.
889    flush_receiver: watch::Receiver<()>,
890    /// Gauge of stalling request count.
891    stalling_count: IntGauge,
892    /// Gauge of regions in the worker.
893    region_count: IntGauge,
894    /// Histogram of request wait time for this worker.
895    request_wait_time: Histogram,
896    /// Queues for region edit requests.
897    region_edit_queues: RegionEditQueues,
898    /// Database level metadata manager.
899    schema_metadata_manager: SchemaMetadataManagerRef,
900    /// Datanode level file references manager.
901    file_ref_manager: FileReferenceManagerRef,
902    /// Partition expr fetcher used to backfill partition expr on open for compatibility.
903    partition_expr_fetcher: PartitionExprFetcherRef,
904    /// Semaphore to control flush concurrency.
905    flush_semaphore: Arc<Semaphore>,
906    /// Plugins for flush hooks.
907    plugins: Plugins,
908}
909
910impl<S: LogStore> RegionWorkerLoop<S> {
911    /// Starts the worker loop.
912    async fn run(&mut self) {
913        let init_check_delay = worker_init_check_delay();
914        info!(
915            "Start region worker thread {}, init_check_delay: {:?}",
916            self.id, init_check_delay
917        );
918        self.last_periodical_check_millis += init_check_delay.as_millis() as i64;
919
920        // Buffer to retrieve requests from receiver.
921        let mut write_req_buffer: Vec<SenderWriteRequest> =
922            Vec::with_capacity(self.config.worker_request_batch_size);
923        let mut bulk_req_buffer: Vec<SenderBulkRequest> =
924            Vec::with_capacity(self.config.worker_request_batch_size);
925        let mut ddl_req_buffer: Vec<SenderDdlRequest> =
926            Vec::with_capacity(self.config.worker_request_batch_size);
927        let mut general_req_buffer: Vec<WorkerRequest> =
928            RequestBuffer::with_capacity(self.config.worker_request_batch_size);
929
930        while self.running.load(Ordering::Relaxed) {
931            // Clear the buffer before handling next batch of requests.
932            write_req_buffer.clear();
933            ddl_req_buffer.clear();
934            general_req_buffer.clear();
935            let mut bulk_insert_req_num = 0;
936
937            let max_wait_time = self.time_provider.wait_duration(CHECK_REGION_INTERVAL);
938            let sleep = tokio::time::sleep(max_wait_time);
939            tokio::pin!(sleep);
940
941            tokio::select! {
942                request_opt = self.receiver.recv() => {
943                    match request_opt {
944                        Some(request_with_time) => {
945                            // Observe the wait time
946                            let wait_time = request_with_time.created_at.elapsed();
947                            self.request_wait_time.observe(wait_time.as_secs_f64());
948
949                            match request_with_time.request {
950                                WorkerRequest::Write(sender_req) => write_req_buffer.push(sender_req),
951                                WorkerRequest::Ddl(sender_req) => ddl_req_buffer.push(sender_req),
952                                WorkerRequest::BulkInserts(bulk_insert) => {
953                                    bulk_insert_req_num += 1;
954                                    self.buffer_bulk_insert_request(bulk_insert, &mut bulk_req_buffer)
955                                        .await;
956                                }
957                                req => general_req_buffer.push(req),
958                            }
959                        },
960                        // The channel is disconnected.
961                        None => break,
962                    }
963                }
964                recv_res = self.flush_receiver.changed() => {
965                    if recv_res.is_err() {
966                        // The channel is disconnected.
967                        break;
968                    } else {
969                        // Also flush this worker if other workers trigger flush as this worker may have
970                        // a large memtable to flush. We may not have chance to flush that memtable if we
971                        // never write to this worker. So only flushing other workers may not release enough
972                        // memory.
973                        self.maybe_flush_worker();
974                        // A flush job is finished, handles stalled requests.
975                        self.handle_stalled_requests().await;
976                        continue;
977                    }
978                }
979                _ = &mut sleep => {
980                    // Timeout. Checks periodical tasks.
981                    self.handle_periodical_tasks();
982                    continue;
983                }
984            }
985
986            if self.flush_receiver.has_changed().unwrap_or(false) {
987                // Always checks whether we could process stalled requests to avoid a request
988                // hangs too long.
989                // If the channel is closed, do nothing.
990                self.handle_stalled_requests().await;
991            }
992
993            // Try to recv more requests from the channel.
994            for _ in 1..self.config.worker_request_batch_size {
995                // We have received one request so we start from 1.
996                match self.receiver.try_recv() {
997                    Ok(request_with_time) => {
998                        // Observe the wait time
999                        let wait_time = request_with_time.created_at.elapsed();
1000                        self.request_wait_time.observe(wait_time.as_secs_f64());
1001
1002                        match request_with_time.request {
1003                            WorkerRequest::Write(sender_req) => write_req_buffer.push(sender_req),
1004                            WorkerRequest::Ddl(sender_req) => ddl_req_buffer.push(sender_req),
1005                            WorkerRequest::BulkInserts(bulk_insert) => {
1006                                bulk_insert_req_num += 1;
1007                                self.buffer_bulk_insert_request(bulk_insert, &mut bulk_req_buffer)
1008                                    .await
1009                            }
1010                            req => general_req_buffer.push(req),
1011                        }
1012                    }
1013                    // We still need to handle remaining requests.
1014                    Err(_) => break,
1015                }
1016            }
1017
1018            self.listener.on_recv_requests(
1019                write_req_buffer.len()
1020                    + ddl_req_buffer.len()
1021                    + general_req_buffer.len()
1022                    + bulk_insert_req_num,
1023            );
1024
1025            self.handle_requests(
1026                &mut write_req_buffer,
1027                &mut ddl_req_buffer,
1028                &mut general_req_buffer,
1029                &mut bulk_req_buffer,
1030            )
1031            .await;
1032
1033            self.handle_periodical_tasks();
1034        }
1035
1036        self.clean().await;
1037
1038        info!("Exit region worker thread {}", self.id);
1039    }
1040
1041    async fn buffer_bulk_insert_request(
1042        &mut self,
1043        bulk_insert: BulkInsertRequest,
1044        bulk_requests: &mut Vec<SenderBulkRequest>,
1045    ) {
1046        let BulkInsertRequest {
1047            metadata,
1048            request,
1049            sender,
1050        } = bulk_insert;
1051
1052        if let Some(region_metadata) = metadata {
1053            self.handle_bulk_insert_batch(region_metadata, request, bulk_requests, sender)
1054                .await;
1055        } else {
1056            error!("Cannot find region metadata for {}", request.region_id);
1057            sender.send(
1058                error::RegionNotFoundSnafu {
1059                    region_id: request.region_id,
1060                }
1061                .fail(),
1062            );
1063        }
1064    }
1065
1066    /// Dispatches and processes requests.
1067    ///
1068    /// `general_requests` should not contain categorized write, ddl, or bulk insert requests.
1069    async fn handle_requests(
1070        &mut self,
1071        write_requests: &mut Vec<SenderWriteRequest>,
1072        ddl_requests: &mut Vec<SenderDdlRequest>,
1073        general_requests: &mut Vec<WorkerRequest>,
1074        bulk_requests: &mut Vec<SenderBulkRequest>,
1075    ) {
1076        for worker_req in general_requests.drain(..) {
1077            match worker_req {
1078                WorkerRequest::Write(_) | WorkerRequest::Ddl(_) => {
1079                    // These requests are categorized before dispatching general requests.
1080                    continue;
1081                }
1082                WorkerRequest::BulkInserts(_) => unreachable!("bulk inserts are buffered"),
1083                WorkerRequest::Background { region_id, notify } => {
1084                    if matches!(
1085                        &notify,
1086                        BackgroundNotify::RegionEdit(edit_result)
1087                            if edit_result.update_region_state
1088                    ) {
1089                        // Region state must be Editing when reach here.
1090                        // This call only moves write/bulk write request into stall queue. When region edit result
1091                        // is processed inside handle_background_notify and region state is switched back to Writable,
1092                        // stalled request will be processed before the next region edit is dequeued from
1093                        // RegionEditQueue immediately in handle_region_edit_result. It not only ensured pending writes
1094                        // are processed in time, but also prevents them from starvation.
1095                        // TODO(hl): maybe we need to merge those queues for pending requests like pending_ddl,
1096                        // region edits and stalled request, so we can simplify the coordination between these queues.
1097                        self.handle_buffered_region_write_requests(
1098                            &region_id,
1099                            write_requests,
1100                            bulk_requests,
1101                        )
1102                        .await;
1103                    }
1104                    // For background notify, we handle it directly.
1105                    self.handle_background_notify(region_id, notify).await;
1106                }
1107                WorkerRequest::SetRegionRoleStateGracefully {
1108                    region_id,
1109                    region_role_state,
1110                    sender,
1111                } => {
1112                    self.set_role_state_gracefully(region_id, region_role_state, sender)
1113                        .await;
1114                }
1115                WorkerRequest::EditRegion(request) => {
1116                    self.handle_region_edit(request);
1117                }
1118                WorkerRequest::Stop => {
1119                    debug_assert!(!self.running.load(Ordering::Relaxed));
1120                }
1121                WorkerRequest::SyncRegion(req) => {
1122                    self.handle_region_sync(req).await;
1123                }
1124                WorkerRequest::RemapManifests(req) => {
1125                    self.handle_remap_manifests_request(req);
1126                }
1127                WorkerRequest::CopyRegionFrom(req) => {
1128                    self.handle_copy_region_from_request(req);
1129                }
1130            }
1131        }
1132
1133        // Handles all write requests first. So we can alter regions without
1134        // considering existing write requests.
1135        self.handle_write_requests(write_requests, bulk_requests, true)
1136            .await;
1137
1138        self.handle_ddl_requests(ddl_requests).await;
1139    }
1140
1141    /// Takes and handles all ddl requests.
1142    async fn handle_ddl_requests(&mut self, ddl_requests: &mut Vec<SenderDdlRequest>) {
1143        if ddl_requests.is_empty() {
1144            return;
1145        }
1146
1147        for ddl in ddl_requests.drain(..) {
1148            let res = match ddl.request {
1149                DdlRequest::Create(req) => self.handle_create_request(ddl.region_id, req).await,
1150                DdlRequest::Drop(req) => {
1151                    self.handle_drop_request(ddl.region_id, req, ddl.sender)
1152                        .await;
1153                    continue;
1154                }
1155                DdlRequest::Open((req, wal_entry_receiver)) => {
1156                    self.handle_open_request(ddl.region_id, req, wal_entry_receiver, ddl.sender)
1157                        .await;
1158                    continue;
1159                }
1160                DdlRequest::OfflineCleanup(req) => {
1161                    self.handle_offline_cleanup_request(ddl.region_id, req)
1162                        .await
1163                }
1164                DdlRequest::Close(req) => {
1165                    self.handle_close_request(ddl.region_id, req, ddl.sender)
1166                        .await;
1167                    continue;
1168                }
1169                DdlRequest::Alter(req) => {
1170                    self.handle_alter_request(ddl.region_id, req, ddl.sender)
1171                        .await;
1172                    continue;
1173                }
1174                DdlRequest::Flush(req) => {
1175                    self.handle_flush_request(ddl.region_id, req, ddl.sender);
1176                    continue;
1177                }
1178                DdlRequest::Compact(req) => {
1179                    self.handle_compaction_request(ddl.region_id, req, ddl.sender)
1180                        .await;
1181                    continue;
1182                }
1183                DdlRequest::BuildIndex(req) => {
1184                    self.handle_build_index_request(ddl.region_id, req, ddl.sender)
1185                        .await;
1186                    continue;
1187                }
1188                DdlRequest::Truncate(req) => {
1189                    self.handle_truncate_request(ddl.region_id, req, ddl.sender)
1190                        .await;
1191                    continue;
1192                }
1193                DdlRequest::Catchup((req, wal_entry_receiver)) => {
1194                    self.handle_catchup_request(ddl.region_id, req, wal_entry_receiver, ddl.sender)
1195                        .await;
1196                    continue;
1197                }
1198                DdlRequest::EnterStaging(req) => {
1199                    self.handle_enter_staging_request(
1200                        ddl.region_id,
1201                        req.partition_directive,
1202                        ddl.sender,
1203                    )
1204                    .await;
1205                    continue;
1206                }
1207                DdlRequest::ApplyStagingManifest(req) => {
1208                    self.handle_apply_staging_manifest_request(ddl.region_id, req, ddl.sender)
1209                        .await;
1210                    continue;
1211                }
1212            };
1213
1214            ddl.sender.send(res);
1215        }
1216    }
1217
1218    /// Handle periodical tasks such as region auto flush.
1219    fn handle_periodical_tasks(&mut self) {
1220        let interval = CHECK_REGION_INTERVAL.as_millis() as i64;
1221        if self
1222            .time_provider
1223            .elapsed_since(self.last_periodical_check_millis)
1224            < interval
1225        {
1226            return;
1227        }
1228
1229        self.last_periodical_check_millis = self.time_provider.current_time_millis();
1230
1231        if let Err(e) = self.flush_periodically() {
1232            error!(e; "Failed to flush regions periodically");
1233        }
1234    }
1235
1236    /// Handles region background request
1237    async fn handle_background_notify(&mut self, region_id: RegionId, notify: BackgroundNotify) {
1238        match notify {
1239            BackgroundNotify::CompactionPickFinished(req) => {
1240                self.handle_compaction_pick_finished(region_id, req).await
1241            }
1242            BackgroundNotify::FlushFinished(req) => {
1243                self.handle_flush_finished(region_id, req).await
1244            }
1245            BackgroundNotify::FlushFailed(req) => self.handle_flush_failed(region_id, req).await,
1246            BackgroundNotify::IndexBuildFinished(req) => {
1247                self.handle_index_build_finished(region_id, req).await
1248            }
1249            BackgroundNotify::IndexBuildStopped(req) => {
1250                self.handle_index_build_stopped(region_id, req).await
1251            }
1252            BackgroundNotify::IndexBuildFailed(req) => {
1253                self.handle_index_build_failed(region_id, req).await
1254            }
1255            BackgroundNotify::IndexBuildRetry(req) => {
1256                self.handle_rebuild_index(req, OptionOutputTx::new(None))
1257                    .await
1258            }
1259            BackgroundNotify::CompactionFinished(req) => {
1260                self.handle_compaction_finished(region_id, req).await
1261            }
1262            BackgroundNotify::CompactionCancelled(req) => {
1263                self.handle_compaction_cancelled(region_id, req).await
1264            }
1265            BackgroundNotify::CompactionFailed(req) => self.handle_compaction_failure(req).await,
1266            BackgroundNotify::Truncate(req) => self.handle_truncate_result(req).await,
1267            BackgroundNotify::DiscardUnflushed(req) => {
1268                self.handle_discard_unflushed_result(req).await
1269            }
1270            BackgroundNotify::RegionChange(req) => {
1271                self.handle_manifest_region_change_result(req).await
1272            }
1273            BackgroundNotify::EnterStaging(req) => self.handle_enter_staging_result(req).await,
1274            BackgroundNotify::RegionEdit(req) => self.handle_region_edit_result(req).await,
1275            BackgroundNotify::CopyRegionFromFinished(req) => {
1276                self.handle_copy_region_from_finished(req)
1277            }
1278        }
1279    }
1280
1281    /// Handles `set_region_role_gracefully`.
1282    async fn set_role_state_gracefully(
1283        &mut self,
1284        region_id: RegionId,
1285        region_role_state: SettableRegionRoleState,
1286        sender: oneshot::Sender<SetRegionRoleStateResponse>,
1287    ) {
1288        if let Some(region) = self.regions.get_region(region_id) {
1289            // We need to do this in background as we need the manifest lock.
1290            common_runtime::spawn_global(async move {
1291                match region.set_role_state_gracefully(region_role_state).await {
1292                    Ok(()) => {
1293                        let last_entry_id = region.version_control.current().last_entry_id;
1294                        let _ = sender.send(SetRegionRoleStateResponse::success(
1295                            SetRegionRoleStateSuccess::mito(last_entry_id),
1296                        ));
1297                    }
1298                    Err(e) => {
1299                        error!(e; "Failed to set region {} role state to {:?}", region_id, region_role_state);
1300                        let _ = sender.send(SetRegionRoleStateResponse::invalid_transition(
1301                            BoxedError::new(e),
1302                        ));
1303                    }
1304                }
1305            });
1306        } else {
1307            let _ = sender.send(SetRegionRoleStateResponse::NotFound);
1308        }
1309    }
1310}
1311
1312impl<S> RegionWorkerLoop<S> {
1313    /// Cleans up the worker.
1314    async fn clean(&self) {
1315        // Closes remaining regions.
1316        let regions = self.regions.list_regions();
1317        for region in regions {
1318            region.stop().await;
1319        }
1320
1321        self.regions.clear();
1322    }
1323
1324    /// Notifies the whole group that a flush job is finished so other
1325    /// workers can handle stalled requests.
1326    fn notify_group(&mut self) {
1327        // Notifies all receivers.
1328        let _ = self.flush_sender.send(());
1329        // Marks the receiver in current worker as seen so the loop won't be waked up immediately.
1330        self.flush_receiver.borrow_and_update();
1331    }
1332}
1333
1334/// Wrapper that only calls event listener in tests.
1335#[derive(Default, Clone)]
1336pub(crate) struct WorkerListener {
1337    #[cfg(any(test, feature = "test"))]
1338    listener: Option<crate::engine::listener::EventListenerRef>,
1339}
1340
1341impl WorkerListener {
1342    #[cfg(any(test, feature = "test"))]
1343    pub(crate) fn new(
1344        listener: Option<crate::engine::listener::EventListenerRef>,
1345    ) -> WorkerListener {
1346        WorkerListener { listener }
1347    }
1348
1349    /// Flush is finished successfully.
1350    pub(crate) fn on_flush_success(&self, region_id: RegionId) {
1351        #[cfg(any(test, feature = "test"))]
1352        if let Some(listener) = &self.listener {
1353            listener.on_flush_success(region_id);
1354        }
1355        // Avoid compiler warning.
1356        let _ = region_id;
1357    }
1358
1359    /// Engine is stalled.
1360    pub(crate) fn on_write_stall(&self) {
1361        #[cfg(any(test, feature = "test"))]
1362        if let Some(listener) = &self.listener {
1363            listener.on_write_stall();
1364        }
1365    }
1366
1367    pub(crate) async fn on_flush_begin(&self, region_id: RegionId) {
1368        #[cfg(any(test, feature = "test"))]
1369        if let Some(listener) = &self.listener {
1370            listener.on_flush_begin(region_id).await;
1371        }
1372        // Avoid compiler warning.
1373        let _ = region_id;
1374    }
1375
1376    pub(crate) async fn on_flush_commit_begin(&self, _region_id: RegionId) {
1377        #[cfg(any(test, feature = "test"))]
1378        if let Some(listener) = &self.listener {
1379            listener.on_flush_commit_begin(_region_id).await;
1380        }
1381    }
1382
1383    pub(crate) fn on_flush_cancel_requested(&self, _region_id: RegionId) {
1384        #[cfg(any(test, feature = "test"))]
1385        if let Some(listener) = &self.listener {
1386            listener.on_flush_cancel_requested(_region_id);
1387        }
1388    }
1389
1390    pub(crate) fn on_later_drop_begin(&self, region_id: RegionId) -> Option<Duration> {
1391        #[cfg(any(test, feature = "test"))]
1392        if let Some(listener) = &self.listener {
1393            return listener.on_later_drop_begin(region_id);
1394        }
1395        // Avoid compiler warning.
1396        let _ = region_id;
1397        None
1398    }
1399
1400    /// On later drop task is finished.
1401    pub(crate) fn on_later_drop_end(&self, region_id: RegionId, removed: bool) {
1402        #[cfg(any(test, feature = "test"))]
1403        if let Some(listener) = &self.listener {
1404            listener.on_later_drop_end(region_id, removed);
1405        }
1406        // Avoid compiler warning.
1407        let _ = region_id;
1408        let _ = removed;
1409    }
1410
1411    pub(crate) async fn on_merge_ssts_finished(&self, region_id: RegionId) {
1412        #[cfg(any(test, feature = "test"))]
1413        if let Some(listener) = &self.listener {
1414            listener.on_merge_ssts_finished(region_id).await;
1415        }
1416        // Avoid compiler warning.
1417        let _ = region_id;
1418    }
1419
1420    pub(crate) fn on_recv_requests(&self, request_num: usize) {
1421        #[cfg(any(test, feature = "test"))]
1422        if let Some(listener) = &self.listener {
1423            listener.on_recv_requests(request_num);
1424        }
1425        // Avoid compiler warning.
1426        let _ = request_num;
1427    }
1428
1429    pub(crate) fn on_file_cache_filled(&self, _file_id: FileId) {
1430        #[cfg(any(test, feature = "test"))]
1431        if let Some(listener) = &self.listener {
1432            listener.on_file_cache_filled(_file_id);
1433        }
1434    }
1435
1436    pub(crate) fn on_compaction_scheduled(&self, _region_id: RegionId) {
1437        #[cfg(any(test, feature = "test"))]
1438        if let Some(listener) = &self.listener {
1439            listener.on_compaction_scheduled(_region_id);
1440        }
1441    }
1442
1443    pub(crate) async fn on_compaction_pick_begin(&self, _region_id: RegionId) {
1444        #[cfg(any(test, feature = "test"))]
1445        if let Some(listener) = &self.listener {
1446            listener.on_compaction_pick_begin(_region_id).await;
1447        }
1448    }
1449
1450    pub(crate) async fn on_compaction_commit_begin(&self, _region_id: RegionId) {
1451        #[cfg(any(test, feature = "test"))]
1452        if let Some(listener) = &self.listener {
1453            listener.on_compaction_commit_begin(_region_id).await;
1454        }
1455    }
1456
1457    pub(crate) async fn on_compaction_result_notified(&self, _region_id: RegionId) {
1458        #[cfg(any(test, feature = "test"))]
1459        if let Some(listener) = &self.listener {
1460            listener.on_compaction_result_notified(_region_id).await;
1461        }
1462    }
1463
1464    pub(crate) fn on_compaction_cancel_requested(&self, _region_id: RegionId) {
1465        #[cfg(any(test, feature = "test"))]
1466        if let Some(listener) = &self.listener {
1467            listener.on_compaction_cancel_requested(_region_id);
1468        }
1469    }
1470
1471    pub(crate) async fn on_notify_region_change_result_begin(&self, _region_id: RegionId) {
1472        #[cfg(any(test, feature = "test"))]
1473        if let Some(listener) = &self.listener {
1474            listener
1475                .on_notify_region_change_result_begin(_region_id)
1476                .await;
1477        }
1478    }
1479
1480    pub(crate) async fn on_enter_staging_result_begin(&self, _region_id: RegionId) {
1481        #[cfg(any(test, feature = "test"))]
1482        if let Some(listener) = &self.listener {
1483            listener.on_enter_staging_result_begin(_region_id).await;
1484        }
1485    }
1486
1487    pub(crate) async fn on_index_build_finish(&self, _region_file_id: RegionFileId) {
1488        #[cfg(any(test, feature = "test"))]
1489        if let Some(listener) = &self.listener {
1490            listener.on_index_build_finish(_region_file_id).await;
1491        }
1492    }
1493
1494    pub(crate) async fn on_index_build_begin(&self, _region_file_id: RegionFileId) {
1495        #[cfg(any(test, feature = "test"))]
1496        if let Some(listener) = &self.listener {
1497            listener.on_index_build_begin(_region_file_id).await;
1498        }
1499    }
1500
1501    pub(crate) async fn on_index_build_abort(&self, _region_file_id: RegionFileId) {
1502        #[cfg(any(test, feature = "test"))]
1503        if let Some(listener) = &self.listener {
1504            listener.on_index_build_abort(_region_file_id).await;
1505        }
1506    }
1507
1508    pub(crate) async fn on_index_build_before_manifest_commit(
1509        &self,
1510        _region_file_id: RegionFileId,
1511    ) {
1512        #[cfg(any(test, feature = "test"))]
1513        if let Some(listener) = &self.listener {
1514            listener
1515                .on_index_build_before_manifest_commit(_region_file_id)
1516                .await;
1517        }
1518    }
1519
1520    pub(crate) async fn on_index_build_manifest_committed(&self, _region_file_id: RegionFileId) {
1521        #[cfg(any(test, feature = "test"))]
1522        if let Some(listener) = &self.listener {
1523            listener
1524                .on_index_build_manifest_committed(_region_file_id)
1525                .await;
1526        }
1527    }
1528}
1529
1530#[cfg(test)]
1531mod tests {
1532    use super::*;
1533    use crate::test_util::TestEnv;
1534
1535    #[test]
1536    fn test_region_id_to_index() {
1537        let num_workers = 4;
1538
1539        let region_id = RegionId::new(1, 2);
1540        let index = region_id_to_index(region_id, num_workers);
1541        assert_eq!(index, 3);
1542
1543        let region_id = RegionId::new(2, 3);
1544        let index = region_id_to_index(region_id, num_workers);
1545        assert_eq!(index, 1);
1546    }
1547
1548    #[tokio::test]
1549    async fn test_worker_group_start_stop() {
1550        let env = TestEnv::with_prefix("group-stop").await;
1551        let group = env
1552            .create_worker_group(MitoConfig {
1553                num_workers: 4,
1554                ..Default::default()
1555            })
1556            .await;
1557
1558        group.stop().await.unwrap();
1559    }
1560}