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