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