1use 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
79pub 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 let receiver = self.region_event_receiver.take().unwrap();
116
117 task.start(receiver, self.leases_notifier.clone()).await?;
118 }
119 Ok(())
120 }
121
122 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 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 let controlled_by_metasrv = meta_client.is_some();
265
266 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 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 is_recovery_mode,
323 );
324
325 if self.opts.init_regions_in_background {
326 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(®ion_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 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 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 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 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 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 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(), 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 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 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 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 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 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 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
707async 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 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}