Skip to main content

datanode/
datanode.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//! Datanode implementation.
16
17use std::path::Path;
18use std::sync::Arc;
19use std::time::{Duration, Instant};
20
21use common_base::Plugins;
22use common_datasource::object_store::LocalFileAccess;
23use common_error::ext::BoxedError;
24use common_greptimedb_telemetry::GreptimeDBTelemetryTask;
25use common_meta::cache::{LayeredCacheRegistry, SchemaCacheRef, TableSchemaCacheRef};
26use common_meta::cache_invalidator::CacheInvalidatorRef;
27use common_meta::datanode::TopicStatsReporter;
28use common_meta::key::runtime_switch::RuntimeSwitchManager;
29use common_meta::key::{SchemaMetadataManager, SchemaMetadataManagerRef};
30use common_meta::kv_backend::KvBackendRef;
31pub use common_procedure::options::ProcedureConfig;
32use common_query::prelude::set_default_prefix;
33use common_stat::ResourceStatImpl;
34use common_telemetry::{error, info, warn};
35use common_wal::config::DatanodeWalConfig;
36use common_wal::config::kafka::DatanodeKafkaConfig;
37use common_wal::config::raft_engine::RaftEngineConfig;
38use file_engine::engine::FileRegionEngine;
39use log_store::kafka::log_store::KafkaLogStore;
40use log_store::kafka::{GlobalIndexCollector, default_index_file};
41use log_store::noop::log_store::NoopLogStore;
42use log_store::raft_engine::log_store::RaftEngineLogStore;
43use meta_client::MetaClientRef;
44use metric_engine::engine::MetricEngine;
45use mito2::config::MitoConfig;
46use mito2::engine::{MitoEngine, MitoEngineBuilder};
47use mito2::region::opener::PartitionExprFetcherRef;
48use mito2::sst::file_ref::{FileReferenceManager, FileReferenceManagerRef};
49use object_store::manager::{ObjectStoreManager, ObjectStoreManagerRef};
50use object_store::util::normalize_dir;
51use query::QueryEngineFactory;
52use query::dummy_catalog::{DummyCatalogManager, TableProviderFactoryRef};
53use servers::server::ServerHandlers;
54use snafu::{OptionExt, ResultExt, ensure};
55use store_api::path_utils::WAL_DIR;
56use store_api::region_engine::{
57    RegionEngineRef, RegionRole, SetRegionRoleStateResponse, SettableRegionRoleState,
58};
59use tokio::fs;
60use tokio::sync::Notify;
61
62use crate::config::{DatanodeOptions, RegionEngineConfig, StorageConfig};
63use crate::error::{
64    self, BuildDatanodeSnafu, BuildMetricEngineSnafu, BuildMitoEngineSnafu, CreateDirSnafu,
65    DataFusionSnafu, GetMetadataSnafu, MissingCacheSnafu, MissingNodeIdSnafu, OpenLogStoreSnafu,
66    Result, ShutdownInstanceSnafu, ShutdownServerSnafu, StartServerSnafu,
67};
68use crate::event_listener::{
69    NoopRegionServerEventListener, RegionServerEventListenerRef, RegionServerEventReceiver,
70    new_region_server_event_channel,
71};
72use crate::greptimedb_telemetry::get_greptimedb_telemetry_task;
73use crate::heartbeat::HeartbeatTask;
74use crate::partition_expr_fetcher::MetaPartitionExprFetcher;
75use crate::region_server::{DummyTableProviderFactory, RegionServer};
76use crate::store::{self, new_object_store_without_cache};
77use crate::utils::{RegionOpenRequests, build_region_open_requests};
78
79/// Datanode service.
80pub struct Datanode {
81    services: ServerHandlers,
82    heartbeat_task: Option<HeartbeatTask>,
83    region_event_receiver: Option<RegionServerEventReceiver>,
84    region_server: RegionServer,
85    greptimedb_telemetry_task: Arc<GreptimeDBTelemetryTask>,
86    leases_notifier: Option<Arc<Notify>>,
87    plugins: Plugins,
88}
89
90impl Datanode {
91    pub async fn start(&mut self) -> Result<()> {
92        info!("Starting datanode instance...");
93
94        self.start_heartbeat().await?;
95        self.wait_coordinated().await;
96
97        self.start_telemetry();
98
99        self.services.start_all().await.context(StartServerSnafu)
100    }
101
102    pub fn server_handlers(&self) -> &ServerHandlers {
103        &self.services
104    }
105
106    pub fn start_telemetry(&self) {
107        if let Err(e) = self.greptimedb_telemetry_task.start() {
108            warn!(e; "Failed to start telemetry task!");
109        }
110    }
111
112    pub async fn start_heartbeat(&mut self) -> Result<()> {
113        if let Some(task) = &self.heartbeat_task {
114            // Safety: The event_receiver must exist.
115            let receiver = self.region_event_receiver.take().unwrap();
116
117            task.start(receiver, self.leases_notifier.clone()).await?;
118        }
119        Ok(())
120    }
121
122    /// If `leases_notifier` exists, it waits until leases have been obtained in all regions.
123    pub async fn wait_coordinated(&mut self) {
124        if let Some(notifier) = self.leases_notifier.take() {
125            notifier.notified().await;
126        }
127    }
128
129    pub fn setup_services(&mut self, services: ServerHandlers) {
130        self.services = services;
131    }
132
133    pub async fn shutdown(&mut self) -> Result<()> {
134        self.services
135            .shutdown_all()
136            .await
137            .context(ShutdownServerSnafu)?;
138
139        let _ = self.greptimedb_telemetry_task.stop().await;
140        if let Some(heartbeat_task) = &self.heartbeat_task {
141            heartbeat_task
142                .close()
143                .map_err(BoxedError::new)
144                .context(ShutdownInstanceSnafu)?;
145        }
146        self.region_server.stop().await?;
147        Ok(())
148    }
149
150    pub fn region_server(&self) -> RegionServer {
151        self.region_server.clone()
152    }
153
154    pub fn plugins(&self) -> Plugins {
155        self.plugins.clone()
156    }
157}
158
159pub struct DatanodeBuilder {
160    opts: DatanodeOptions,
161    table_provider_factory: Option<TableProviderFactoryRef>,
162    plugins: Plugins,
163    meta_client: Option<MetaClientRef>,
164    kv_backend: KvBackendRef,
165    cache_registry: Option<Arc<LayeredCacheRegistry>>,
166    topic_stats_reporter: Option<Box<dyn TopicStatsReporter>>,
167    open_regions_writable_override: Option<bool>,
168    local_file_access: LocalFileAccess,
169    #[cfg(feature = "enterprise")]
170    extension_range_provider_factory: Option<mito2::extension::BoxedExtensionRangeProviderFactory>,
171}
172
173impl DatanodeBuilder {
174    pub fn new(opts: DatanodeOptions, plugins: Plugins, kv_backend: KvBackendRef) -> Self {
175        Self {
176            opts,
177            table_provider_factory: None,
178            plugins,
179            meta_client: None,
180            kv_backend,
181            cache_registry: None,
182            open_regions_writable_override: None,
183            local_file_access: LocalFileAccess::Disabled,
184            #[cfg(feature = "enterprise")]
185            extension_range_provider_factory: None,
186            topic_stats_reporter: None,
187        }
188    }
189
190    pub fn options(&self) -> &DatanodeOptions {
191        &self.opts
192    }
193
194    pub fn with_meta_client(&mut self, client: MetaClientRef) -> &mut Self {
195        self.meta_client = Some(client);
196        self
197    }
198
199    pub fn with_cache_registry(&mut self, registry: Arc<LayeredCacheRegistry>) -> &mut Self {
200        self.cache_registry = Some(registry);
201        self
202    }
203
204    pub fn with_local_file_access(&mut self, local_file_access: LocalFileAccess) -> &mut Self {
205        self.local_file_access = local_file_access;
206        self
207    }
208
209    pub fn kv_backend(&self) -> &KvBackendRef {
210        &self.kv_backend
211    }
212
213    pub fn meta_client(&self) -> Option<&MetaClientRef> {
214        self.meta_client.as_ref()
215    }
216
217    pub fn cache_registry(&self) -> Option<&Arc<LayeredCacheRegistry>> {
218        self.cache_registry.as_ref()
219    }
220
221    pub fn set_plugins(&mut self, plugins: Plugins) {
222        self.plugins = plugins;
223    }
224
225    pub fn with_table_provider_factory(&mut self, factory: TableProviderFactoryRef) -> &mut Self {
226        self.table_provider_factory = Some(factory);
227        self
228    }
229
230    /// Overrides whether regions opened during datanode startup should become writable.
231    ///
232    /// When unset, the builder uses its default writable policy for reopened regions
233    /// (writable only when no metasrv client is configured).
234    ///
235    /// Warning: setting this to `true` on a metasrv-controlled datanode (one built
236    /// with `with_meta_client`) will promote regions to Leader before heartbeat and
237    /// lease coordination begin, bypassing the metasrv safety contract and creating a
238    /// potential split-brain window during startup.
239    pub fn with_open_regions_writable_override(&mut self, writable: bool) -> &mut Self {
240        self.open_regions_writable_override = Some(writable);
241        self
242    }
243
244    #[cfg(feature = "enterprise")]
245    pub fn with_extension_range_provider(
246        &mut self,
247        extension_range_provider_factory: mito2::extension::BoxedExtensionRangeProviderFactory,
248    ) -> &mut Self {
249        self.extension_range_provider_factory = Some(extension_range_provider_factory);
250        self
251    }
252
253    pub async fn build(mut self) -> Result<Datanode> {
254        let node_id = self.opts.node_id.context(MissingNodeIdSnafu)?;
255        set_default_prefix(self.opts.default_column_prefix.as_deref())
256            .map_err(BoxedError::new)
257            .context(BuildDatanodeSnafu)?;
258
259        let meta_client = self.meta_client.take();
260
261        // If metasrv client is provided, we will use it to control the region server.
262        // Otherwise the region server is self-controlled, meaning no heartbeat and immediately
263        // writable upon open.
264        let controlled_by_metasrv = meta_client.is_some();
265
266        // build and initialize region server
267        let (region_event_listener, region_event_receiver) = if controlled_by_metasrv {
268            let (tx, rx) = new_region_server_event_channel();
269            (Box::new(tx) as _, Some(rx))
270        } else {
271            (Box::new(NoopRegionServerEventListener) as _, None)
272        };
273
274        let cache_registry = self.cache_registry.take().context(MissingCacheSnafu)?;
275        let schema_cache: SchemaCacheRef = cache_registry.get().context(MissingCacheSnafu)?;
276        let table_id_schema_cache: TableSchemaCacheRef =
277            cache_registry.get().context(MissingCacheSnafu)?;
278
279        let schema_metadata_manager = Arc::new(SchemaMetadataManager::new(
280            table_id_schema_cache,
281            schema_cache,
282        ));
283
284        let gc_enabled = self.opts.region_engine.iter().any(|engine| {
285            if let RegionEngineConfig::Mito(config) = engine {
286                config.gc.enable
287            } else {
288                false
289            }
290        });
291
292        let file_ref_manager = Arc::new(FileReferenceManager::with_gc_enabled(
293            Some(node_id),
294            gc_enabled,
295        ));
296        let region_server = self
297            .new_region_server(
298                schema_metadata_manager,
299                region_event_listener,
300                file_ref_manager,
301            )
302            .await?;
303
304        // TODO(weny): Considering introducing a readonly kv_backend trait.
305        let runtime_switch_manager = RuntimeSwitchManager::new(self.kv_backend.clone());
306        let is_recovery_mode = runtime_switch_manager
307            .recovery_mode()
308            .await
309            .context(GetMetadataSnafu)?;
310
311        let region_open_requests =
312            build_region_open_requests(node_id, self.kv_backend.clone()).await?;
313        let open_with_writable = self
314            .open_regions_writable_override
315            .unwrap_or(!controlled_by_metasrv);
316        let open_all_regions = open_all_regions(
317            region_server.clone(),
318            region_open_requests,
319            open_with_writable,
320            self.opts.init_regions_parallelism,
321            // Ignore nonexistent regions in recovery mode.
322            is_recovery_mode,
323        );
324
325        if self.opts.init_regions_in_background {
326            // Opens regions in background.
327            common_runtime::spawn_global(async move {
328                if let Err(err) = open_all_regions.await {
329                    error!(err; "Failed to open regions during the startup.");
330                }
331            });
332        } else {
333            open_all_regions.await?;
334        }
335
336        let heartbeat_task = if let Some(meta_client) = meta_client {
337            let task = self
338                .create_heartbeat_task(&region_server, meta_client, cache_registry)
339                .await?;
340            Some(task)
341        } else {
342            None
343        };
344
345        let is_standalone = heartbeat_task.is_none();
346        let greptimedb_telemetry_task = get_greptimedb_telemetry_task(
347            Some(self.opts.storage.data_home.clone()),
348            is_standalone && self.opts.enable_telemetry,
349        )
350        .await;
351
352        let leases_notifier = if self.opts.require_lease_before_startup && !is_standalone {
353            Some(Arc::new(Notify::new()))
354        } else {
355            None
356        };
357
358        Ok(Datanode {
359            services: ServerHandlers::default(),
360            heartbeat_task,
361            region_server,
362            greptimedb_telemetry_task,
363            region_event_receiver,
364            leases_notifier,
365            plugins: self.plugins.clone(),
366        })
367    }
368
369    async fn create_heartbeat_task(
370        &self,
371        region_server: &RegionServer,
372        meta_client: MetaClientRef,
373        cache_invalidator: CacheInvalidatorRef,
374    ) -> Result<HeartbeatTask> {
375        let stat = {
376            let mut stat = ResourceStatImpl::default();
377            stat.start_collect_cpu_usage();
378            Arc::new(stat)
379        };
380
381        HeartbeatTask::try_new(
382            &self.opts,
383            region_server.clone(),
384            meta_client,
385            self.kv_backend.clone(),
386            cache_invalidator,
387            self.plugins.clone(),
388            stat,
389        )
390        .await
391    }
392
393    /// Builds [ObjectStoreManager] from [StorageConfig].
394    pub async fn build_object_store_manager(cfg: &StorageConfig) -> Result<ObjectStoreManagerRef> {
395        let object_store = store::new_object_store(cfg.store.clone(), &cfg.data_home).await?;
396        let default_name = cfg.store.config_name();
397        let mut object_store_manager = ObjectStoreManager::new(default_name, object_store);
398        for store in &cfg.providers {
399            object_store_manager.add(
400                store.config_name(),
401                store::new_object_store(store.clone(), &cfg.data_home).await?,
402            );
403        }
404        Ok(Arc::new(object_store_manager))
405    }
406
407    #[cfg(test)]
408    /// Open all regions belong to this datanode.
409    async fn initialize_region_server(
410        &self,
411        region_server: &RegionServer,
412        open_with_writable: bool,
413    ) -> Result<()> {
414        let node_id = self.opts.node_id.context(MissingNodeIdSnafu)?;
415
416        // TODO(weny): Considering introducing a readonly kv_backend trait.
417        let runtime_switch_manager = RuntimeSwitchManager::new(self.kv_backend.clone());
418        let is_recovery_mode = runtime_switch_manager
419            .recovery_mode()
420            .await
421            .context(GetMetadataSnafu)?;
422        let region_open_requests =
423            build_region_open_requests(node_id, self.kv_backend.clone()).await?;
424
425        open_all_regions(
426            region_server.clone(),
427            region_open_requests,
428            open_with_writable,
429            self.opts.init_regions_parallelism,
430            is_recovery_mode,
431        )
432        .await
433    }
434
435    async fn new_region_server(
436        &mut self,
437        schema_metadata_manager: SchemaMetadataManagerRef,
438        event_listener: RegionServerEventListenerRef,
439        file_ref_manager: FileReferenceManagerRef,
440    ) -> Result<RegionServer> {
441        let opts: &DatanodeOptions = &self.opts;
442
443        let query_engine_factory = QueryEngineFactory::try_new_with_plugins(
444            // query engine in datanode only executes plan with resolved table source.
445            DummyCatalogManager::arc(),
446            None,
447            None,
448            None,
449            None,
450            None,
451            false,
452            self.plugins.clone(),
453            opts.query.clone(),
454        )
455        .context(DataFusionSnafu)?;
456        let query_engine = query_engine_factory.query_engine();
457
458        let table_provider_factory = self
459            .table_provider_factory
460            .clone()
461            .unwrap_or_else(|| Arc::new(DummyTableProviderFactory));
462        let mut region_server = RegionServer::with_table_provider(
463            query_engine,
464            common_runtime::global_runtime(),
465            event_listener,
466            table_provider_factory,
467            opts.max_concurrent_queries,
468            opts.concurrent_query_limiter_timeout,
469            opts.grpc.flight_compression,
470        );
471        region_server.install_remote_dyn_filter_receiver_injector(&self.plugins);
472
473        let object_store_manager = Self::build_object_store_manager(&opts.storage).await?;
474        let engines = self
475            .build_store_engines(
476                object_store_manager,
477                schema_metadata_manager,
478                file_ref_manager,
479                self.plugins.clone(),
480            )
481            .await?;
482        for engine in engines {
483            region_server.register_engine(engine);
484        }
485        if let Some(topic_stats_reporter) = self.topic_stats_reporter.take() {
486            region_server.set_topic_stats_reporter(topic_stats_reporter);
487        }
488
489        Ok(region_server)
490    }
491
492    // internal utils
493
494    /// Builds [RegionEngineRef] from `store_engine` section in `opts`
495    async fn build_store_engines(
496        &mut self,
497        object_store_manager: ObjectStoreManagerRef,
498        schema_metadata_manager: SchemaMetadataManagerRef,
499        file_ref_manager: FileReferenceManagerRef,
500        plugins: Plugins,
501    ) -> Result<Vec<RegionEngineRef>> {
502        let mut metric_engine_config = metric_engine::config::EngineConfig::default();
503        let mut mito_engine_config = MitoConfig::default();
504        let mut file_engine_config = file_engine::config::EngineConfig::default();
505
506        for engine in &self.opts.region_engine {
507            match engine {
508                RegionEngineConfig::Mito(config) => {
509                    mito_engine_config = config.clone();
510                }
511                RegionEngineConfig::File(config) => {
512                    file_engine_config = config.clone();
513                }
514                RegionEngineConfig::Metric(metric_config) => {
515                    metric_engine_config = metric_config.clone();
516                }
517            }
518        }
519
520        // Build a fetcher to backfill partition_expr on open.
521        let fetcher = Arc::new(MetaPartitionExprFetcher::new(self.kv_backend.clone()));
522        let mito_engine = self
523            .build_mito_engine(
524                object_store_manager.clone(),
525                mito_engine_config,
526                schema_metadata_manager.clone(),
527                file_ref_manager.clone(),
528                fetcher.clone(),
529                plugins.clone(),
530            )
531            .await?;
532
533        let metric_engine = MetricEngine::try_new(mito_engine.clone(), metric_engine_config)
534            .context(BuildMetricEngineSnafu)?;
535
536        let file_engine = FileRegionEngine::new(
537            file_engine_config,
538            object_store_manager.default_object_store().clone(), // TODO: implement custom storage for file engine
539            self.local_file_access.clone(),
540        );
541
542        Ok(vec![
543            Arc::new(mito_engine) as _,
544            Arc::new(metric_engine) as _,
545            Arc::new(file_engine) as _,
546        ])
547    }
548
549    /// Builds [MitoEngine] according to options.
550    async fn build_mito_engine(
551        &mut self,
552        object_store_manager: ObjectStoreManagerRef,
553        mut config: MitoConfig,
554        schema_metadata_manager: SchemaMetadataManagerRef,
555        file_ref_manager: FileReferenceManagerRef,
556        partition_expr_fetcher: PartitionExprFetcherRef,
557        plugins: Plugins,
558    ) -> Result<MitoEngine> {
559        let opts = &self.opts;
560        if opts.storage.is_object_storage() {
561            // Enable the write cache when setting object storage
562            config.enable_write_cache = true;
563            info!("Configured 'enable_write_cache=true' for mito engine.");
564        }
565
566        let mito_engine = match &opts.wal {
567            DatanodeWalConfig::RaftEngine(raft_engine_config) => {
568                let log_store =
569                    Self::build_raft_engine_log_store(&opts.storage.data_home, raft_engine_config)
570                        .await?;
571
572                let builder = MitoEngineBuilder::new(
573                    &opts.storage.data_home,
574                    config,
575                    log_store,
576                    object_store_manager,
577                    schema_metadata_manager,
578                    file_ref_manager,
579                    partition_expr_fetcher.clone(),
580                    plugins,
581                );
582
583                #[cfg(feature = "enterprise")]
584                let builder = builder.with_extension_range_provider_factory(
585                    self.extension_range_provider_factory.take(),
586                );
587
588                builder.try_build().await.context(BuildMitoEngineSnafu)?
589            }
590            DatanodeWalConfig::Kafka(kafka_config) => {
591                if kafka_config.create_index && opts.node_id.is_none() {
592                    warn!("The WAL index creation only available in distributed mode.")
593                }
594                let global_index_collector = if kafka_config.create_index
595                    && let Some(node_id) = opts.node_id
596                {
597                    let operator = new_object_store_without_cache(
598                        &opts.storage.store,
599                        &opts.storage.data_home,
600                    )
601                    .await?;
602                    let path = default_index_file(node_id);
603                    Some(Self::build_global_index_collector(
604                        kafka_config.dump_index_interval,
605                        operator,
606                        path,
607                    ))
608                } else {
609                    None
610                };
611
612                let log_store =
613                    Self::build_kafka_log_store(kafka_config, global_index_collector).await?;
614                self.topic_stats_reporter = Some(log_store.topic_stats_reporter());
615                let builder = MitoEngineBuilder::new(
616                    &opts.storage.data_home,
617                    config,
618                    log_store,
619                    object_store_manager,
620                    schema_metadata_manager,
621                    file_ref_manager,
622                    partition_expr_fetcher,
623                    plugins,
624                );
625
626                #[cfg(feature = "enterprise")]
627                let builder = builder.with_extension_range_provider_factory(
628                    self.extension_range_provider_factory.take(),
629                );
630
631                builder.try_build().await.context(BuildMitoEngineSnafu)?
632            }
633            DatanodeWalConfig::Noop => {
634                let log_store = Arc::new(NoopLogStore);
635
636                let builder = MitoEngineBuilder::new(
637                    &opts.storage.data_home,
638                    config,
639                    log_store,
640                    object_store_manager,
641                    schema_metadata_manager,
642                    file_ref_manager,
643                    partition_expr_fetcher.clone(),
644                    plugins,
645                );
646
647                #[cfg(feature = "enterprise")]
648                let builder = builder.with_extension_range_provider_factory(
649                    self.extension_range_provider_factory.take(),
650                );
651
652                builder.try_build().await.context(BuildMitoEngineSnafu)?
653            }
654        };
655        Ok(mito_engine)
656    }
657
658    /// Builds [RaftEngineLogStore].
659    async fn build_raft_engine_log_store(
660        data_home: &str,
661        config: &RaftEngineConfig,
662    ) -> Result<Arc<RaftEngineLogStore>> {
663        let data_home = normalize_dir(data_home);
664        let wal_dir = match &config.dir {
665            Some(dir) => dir.clone(),
666            None => format!("{}{WAL_DIR}", data_home),
667        };
668
669        // create WAL directory
670        fs::create_dir_all(Path::new(&wal_dir))
671            .await
672            .context(CreateDirSnafu { dir: &wal_dir })?;
673        info!(
674            "Creating raft-engine logstore with config: {:?} and storage path: {}",
675            config, &wal_dir
676        );
677        let logstore = RaftEngineLogStore::try_new(wal_dir, config)
678            .await
679            .map_err(Box::new)
680            .context(OpenLogStoreSnafu)?;
681
682        Ok(Arc::new(logstore))
683    }
684
685    /// Builds [`KafkaLogStore`].
686    async fn build_kafka_log_store(
687        config: &DatanodeKafkaConfig,
688        global_index_collector: Option<GlobalIndexCollector>,
689    ) -> Result<Arc<KafkaLogStore>> {
690        KafkaLogStore::try_new(config, global_index_collector)
691            .await
692            .map_err(Box::new)
693            .context(OpenLogStoreSnafu)
694            .map(Arc::new)
695    }
696
697    /// Builds [`GlobalIndexCollector`]
698    fn build_global_index_collector(
699        dump_index_interval: Duration,
700        operator: object_store::ObjectStore,
701        path: String,
702    ) -> GlobalIndexCollector {
703        GlobalIndexCollector::new(dump_index_interval, operator, path)
704    }
705}
706
707/// Open all regions belong to this datanode.
708async fn open_all_regions(
709    region_server: RegionServer,
710    region_open_requests: RegionOpenRequests,
711    open_with_writable: bool,
712    init_regions_parallelism: usize,
713    ignore_nonexistent_region: bool,
714) -> Result<()> {
715    let RegionOpenRequests {
716        leader_regions,
717        #[cfg(feature = "enterprise")]
718        follower_regions,
719    } = region_open_requests;
720
721    let leader_region_num = leader_regions.len();
722    info!("going to open {} region(s)", leader_region_num);
723    let now = Instant::now();
724    let open_regions = region_server
725        .handle_batch_open_requests(
726            init_regions_parallelism,
727            leader_regions,
728            ignore_nonexistent_region,
729        )
730        .await?;
731    info!(
732        "Opened {} regions in {:?}",
733        open_regions.len(),
734        now.elapsed()
735    );
736    if !ignore_nonexistent_region {
737        ensure!(
738            open_regions.len() == leader_region_num,
739            error::UnexpectedSnafu {
740                violated: format!(
741                    "Expected to open {} of regions, only {} of regions has opened",
742                    leader_region_num,
743                    open_regions.len()
744                )
745            }
746        );
747    } else if open_regions.len() != leader_region_num {
748        warn!(
749            "ignore nonexistent region, expected to open {} of regions, only {} of regions has opened",
750            leader_region_num,
751            open_regions.len()
752        );
753    }
754
755    for region_id in open_regions {
756        if open_with_writable {
757            let res = region_server.set_region_role(region_id, RegionRole::Leader);
758            match res {
759                Ok(_) => {
760                    // Finalize leadership: persist backfilled metadata.
761                    if let SetRegionRoleStateResponse::InvalidTransition(err) = region_server
762                        .set_region_role_state_gracefully(
763                            region_id,
764                            SettableRegionRoleState::Leader,
765                        )
766                        .await?
767                    {
768                        error!(err; "failed to convert region {region_id} to leader");
769                    }
770                }
771                Err(e) => {
772                    error!(e; "failed to convert region {region_id} to leader");
773                }
774            }
775        }
776    }
777
778    #[cfg(feature = "enterprise")]
779    if !follower_regions.is_empty() {
780        use tokio::time::Instant;
781
782        let follower_region_num = follower_regions.len();
783        info!("going to open {} follower region(s)", follower_region_num);
784
785        let now = Instant::now();
786        let open_regions = region_server
787            .handle_batch_open_requests(
788                init_regions_parallelism,
789                follower_regions,
790                ignore_nonexistent_region,
791            )
792            .await?;
793        info!(
794            "Opened {} follower regions in {:?}",
795            open_regions.len(),
796            now.elapsed()
797        );
798
799        if !ignore_nonexistent_region {
800            ensure!(
801                open_regions.len() == follower_region_num,
802                error::UnexpectedSnafu {
803                    violated: format!(
804                        "Expected to open {} of follower regions, only {} of regions has opened",
805                        follower_region_num,
806                        open_regions.len()
807                    )
808                }
809            );
810        } else if open_regions.len() != follower_region_num {
811            warn!(
812                "ignore nonexistent region, expected to open {} of follower regions, only {} of regions has opened",
813                follower_region_num,
814                open_regions.len()
815            );
816        }
817    }
818
819    info!("all regions are opened");
820
821    Ok(())
822}
823
824#[cfg(test)]
825mod tests {
826    use std::assert_matches;
827    use std::collections::{BTreeMap, HashMap};
828    use std::sync::Arc;
829
830    use cache::build_datanode_cache_registry;
831    use common_base::Plugins;
832    use common_meta::cache::LayeredCacheRegistryBuilder;
833    use common_meta::key::RegionRoleSet;
834    use common_meta::key::datanode_table::DatanodeTableManager;
835    use common_meta::kv_backend::KvBackendRef;
836    use common_meta::kv_backend::memory::MemoryKvBackend;
837    use mito2::engine::MITO_ENGINE_NAME;
838    use store_api::region_request::RegionRequest;
839    use store_api::storage::RegionId;
840
841    use crate::config::DatanodeOptions;
842    use crate::datanode::DatanodeBuilder;
843    use crate::tests::{MockRegionEngine, mock_region_server};
844
845    async fn setup_table_datanode(kv: &KvBackendRef) {
846        let mgr = DatanodeTableManager::new(kv.clone());
847        let txn = mgr
848            .build_create_txn(
849                1028,
850                MITO_ENGINE_NAME,
851                "foo/bar/weny",
852                HashMap::from([("foo".to_string(), "bar".to_string())]),
853                HashMap::default(),
854                BTreeMap::from([(0, RegionRoleSet::new(vec![0, 1, 2], vec![]))]),
855            )
856            .unwrap();
857
858        let r = kv.txn(txn).await.unwrap();
859        assert!(r.succeeded);
860    }
861
862    #[tokio::test]
863    async fn test_initialize_region_server() {
864        common_telemetry::init_default_ut_logging();
865        let mut mock_region_server = mock_region_server();
866        let (mock_region, mut mock_region_handler) = MockRegionEngine::new(MITO_ENGINE_NAME);
867
868        mock_region_server.register_engine(mock_region.clone());
869
870        let kv_backend = Arc::new(MemoryKvBackend::new());
871        let layered_cache_registry = Arc::new(
872            LayeredCacheRegistryBuilder::default()
873                .add_cache_registry(build_datanode_cache_registry(kv_backend.clone()))
874                .build(),
875        );
876
877        let mut builder = DatanodeBuilder::new(
878            DatanodeOptions {
879                node_id: Some(0),
880                ..Default::default()
881            },
882            Plugins::default(),
883            kv_backend.clone(),
884        );
885        builder.with_cache_registry(layered_cache_registry);
886        setup_table_datanode(&(kv_backend as _)).await;
887
888        builder
889            .initialize_region_server(&mock_region_server, false)
890            .await
891            .unwrap();
892
893        for i in 0..3 {
894            let (region_id, req) = mock_region_handler.recv().await.unwrap();
895            assert_eq!(region_id, RegionId::new(1028, i));
896            if let RegionRequest::Open(req) = req {
897                assert_eq!(
898                    req.options,
899                    HashMap::from([("foo".to_string(), "bar".to_string())])
900                )
901            } else {
902                unreachable!()
903            }
904        }
905
906        assert_matches!(
907            mock_region_handler.try_recv(),
908            Err(tokio::sync::mpsc::error::TryRecvError::Empty)
909        );
910    }
911}