1pub mod builder;
16
17use std::fmt::{self, Display};
18use std::sync::atomic::{AtomicBool, Ordering};
19use std::sync::{Arc, Mutex, RwLock};
20use std::time::Duration;
21
22use api::v1::meta::{HeartbeatConfig, Role};
23use clap::ValueEnum;
24use common_base::Plugins;
25use common_base::readable_size::ReadableSize;
26use common_config::{Configurable, DEFAULT_DATA_HOME};
27use common_event_recorder::EventRecorderOptions;
28use common_greptimedb_telemetry::GreptimeDBTelemetryTask;
29use common_meta::cache_invalidator::CacheInvalidatorRef;
30use common_meta::ddl::allocator::resource_id::ResourceIdAllocatorRef;
31use common_meta::ddl_manager::DdlManagerRef;
32use common_meta::distributed_time_constants::{
33 self, BASE_HEARTBEAT_INTERVAL, default_distributed_time_constants, frontend_heartbeat_interval,
34};
35use common_meta::election::LeaderChangeMessage;
36pub use common_meta::election::{ElectionRef, MetasrvNodeInfo};
37use common_meta::key::TableMetadataManagerRef;
38use common_meta::key::runtime_switch::RuntimeSwitchManagerRef;
39use common_meta::kv_backend::{KvBackendRef, ResettableKvBackend, ResettableKvBackendRef};
40use common_meta::leadership_notifier::{
41 LeadershipChangeNotifier, LeadershipChangeNotifierCustomizerRef,
42};
43use common_meta::node_expiry_listener::NodeExpiryListener;
44use common_meta::peer::{Peer, PeerDiscoveryRef};
45use common_meta::reconciliation::manager::ReconciliationManagerRef;
46use common_meta::region_keeper::MemoryRegionKeeperRef;
47use common_meta::region_registry::LeaderRegionRegistryRef;
48use common_meta::stats::topic::TopicStatsRegistryRef;
49use common_meta::wal_provider::WalProviderRef;
50use common_options::datanode::DatanodeClientOptions;
51use common_options::memory::MemoryOptions;
52use common_procedure::ProcedureManagerRef;
53use common_procedure::options::ProcedureConfig;
54use common_stat::ResourceStatRef;
55use common_telemetry::logging::{LoggingOptions, TracingOptions};
56use common_telemetry::{error, info, warn};
57use common_time::util::DefaultSystemTimer;
58use common_wal::config::MetasrvWalConfig;
59use serde::{Deserialize, Serialize};
60use servers::grpc::GrpcOptions;
61use servers::http::HttpOptions;
62use servers::tls::TlsOption;
63use snafu::{OptionExt, ResultExt};
64use store_api::storage::RegionId;
65use tokio::sync::broadcast::error::RecvError;
66
67use crate::cluster::MetaPeerClientRef;
68use crate::discovery;
69use crate::error::{
70 self, InitMetadataSnafu, KvBackendSnafu, Result, StartProcedureManagerSnafu,
71 StartTelemetryTaskSnafu, StopProcedureManagerSnafu,
72};
73use crate::failure_detector::PhiAccrualFailureDetectorOptions;
74use crate::gc::{GcSchedulerOptions, GcTickerRef};
75use crate::handler::{HeartbeatHandlerGroupBuilder, HeartbeatHandlerGroupRef};
76use crate::procedure::ProcedureManagerListenerAdapter;
77use crate::procedure::region_migration::manager::RegionMigrationManagerRef;
78use crate::procedure::repartition::gc_requirement::RepartitionGcRequirementManagerRef;
79use crate::procedure::wal_prune::manager::WalPruneTickerRef;
80use crate::pubsub::{PublisherRef, SubscriptionManagerRef};
81use crate::region::flush_trigger::RegionFlushTickerRef;
82use crate::region::supervisor::RegionSupervisorTickerRef;
83use crate::selector::{RegionStatAwareSelector, Selector, SelectorType};
84use crate::service::mailbox::MailboxRef;
85use crate::service::store::cached_kv::LeaderCachedKvBackend;
86use crate::state::{StateRef, become_follower, become_leader};
87use crate::utils::database::DatabaseOperatorRef;
88
89pub const TABLE_ID_SEQ: &str = "table_id";
90pub const FLOW_ID_SEQ: &str = "flow_id";
91pub const METASRV_DATA_DIR: &str = "metasrv";
92
93#[derive(Clone, Debug, PartialEq, Serialize, Default, Deserialize, ValueEnum)]
95#[serde(rename_all = "snake_case")]
96pub enum BackendImpl {
97 #[default]
99 EtcdStore,
100 MemoryStore,
102 #[cfg(feature = "pg_kvbackend")]
103 PostgresStore,
105 #[cfg(feature = "mysql_kvbackend")]
106 MysqlStore,
108}
109
110#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
112pub struct StatsPersistenceOptions {
113 #[serde(with = "humantime_serde")]
115 pub ttl: Duration,
116 #[serde(with = "humantime_serde")]
118 pub interval: Duration,
119}
120
121impl Default for StatsPersistenceOptions {
122 fn default() -> Self {
123 Self {
124 ttl: Duration::ZERO,
125 interval: Duration::from_mins(10),
126 }
127 }
128}
129
130#[derive(Clone, PartialEq, Serialize, Deserialize, Debug)]
132#[serde(default)]
133pub struct HeartbeatOptions {
134 #[serde(with = "humantime_serde")]
136 pub interval: Duration,
137 #[serde(with = "humantime_serde")]
139 pub retry_interval: Duration,
140}
141
142impl Default for HeartbeatOptions {
143 fn default() -> Self {
144 Self {
145 interval: BASE_HEARTBEAT_INTERVAL,
146 retry_interval: BASE_HEARTBEAT_INTERVAL,
147 }
148 }
149}
150
151impl HeartbeatOptions {
152 pub fn datanode_from(base_interval: Duration) -> Self {
153 Self {
154 interval: base_interval,
155 retry_interval: base_interval,
156 }
157 }
158
159 pub fn frontend_from(base_interval: Duration) -> Self {
160 Self {
161 interval: frontend_heartbeat_interval(base_interval),
162 retry_interval: base_interval,
163 }
164 }
165
166 pub fn flownode_from(base_interval: Duration) -> Self {
167 Self {
168 interval: base_interval,
169 retry_interval: base_interval,
170 }
171 }
172}
173
174impl From<HeartbeatOptions> for HeartbeatConfig {
175 fn from(opts: HeartbeatOptions) -> Self {
176 Self {
177 heartbeat_interval_ms: opts.interval.as_millis() as u64,
178 retry_interval_ms: opts.retry_interval.as_millis() as u64,
179 gc_enabled: false,
180 }
181 }
182}
183
184#[derive(Clone, PartialEq, Serialize, Deserialize, Debug)]
185#[serde(default)]
186pub struct BackendClientOptions {
187 #[serde(with = "humantime_serde")]
188 pub keep_alive_timeout: Duration,
189 #[serde(with = "humantime_serde")]
190 pub keep_alive_interval: Duration,
191 #[serde(with = "humantime_serde")]
192 pub connect_timeout: Duration,
193}
194
195impl Default for BackendClientOptions {
196 fn default() -> Self {
197 Self {
198 keep_alive_interval: Duration::from_secs(10),
199 keep_alive_timeout: Duration::from_secs(3),
200 connect_timeout: Duration::from_secs(3),
201 }
202 }
203}
204
205#[derive(Clone, PartialEq, Serialize, Deserialize)]
206#[serde(default)]
207pub struct MetasrvOptions {
208 #[deprecated(note = "Use grpc.bind_addr instead")]
210 pub bind_addr: String,
211 #[deprecated(note = "Use grpc.server_addr instead")]
213 pub server_addr: String,
214 pub store_addrs: Vec<String>,
216 #[serde(default)]
219 pub backend_tls: Option<TlsOption>,
220 #[serde(default)]
223 pub backend_client: BackendClientOptions,
224 pub selector: SelectorType,
226 pub enable_region_failover: bool,
228 #[serde(with = "humantime_serde")]
233 pub heartbeat_interval: Duration,
234 #[serde(with = "humantime_serde")]
238 pub region_failure_detector_initialization_delay: Duration,
239 pub allow_region_failover_on_local_wal: bool,
244 pub grpc: GrpcOptions,
245 pub http: HttpOptions,
247 pub logging: LoggingOptions,
249 pub procedure: ProcedureConfig,
251 pub failure_detector: PhiAccrualFailureDetectorOptions,
253 pub datanode: DatanodeClientOptions,
255 pub enable_telemetry: bool,
257 pub data_home: String,
259 pub wal: MetasrvWalConfig,
261 pub store_key_prefix: String,
264 pub max_txn_ops: usize,
275 pub flush_stats_factor: usize,
279 pub tracing: TracingOptions,
281 pub memory: MemoryOptions,
283 pub backend: BackendImpl,
285 #[cfg(any(feature = "pg_kvbackend", feature = "mysql_kvbackend"))]
286 pub meta_table_name: String,
288 #[cfg(feature = "pg_kvbackend")]
289 pub meta_election_lock_id: u64,
291 #[cfg(feature = "pg_kvbackend")]
292 pub meta_schema_name: Option<String>,
294 #[cfg(feature = "pg_kvbackend")]
295 pub auto_create_schema: bool,
297 #[serde(with = "humantime_serde")]
298 pub node_max_idle_time: Duration,
299 pub event_recorder: EventRecorderOptions,
301 pub stats_persistence: StatsPersistenceOptions,
303 pub gc: GcSchedulerOptions,
305}
306
307impl fmt::Debug for MetasrvOptions {
308 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
309 let mut debug_struct = f.debug_struct("MetasrvOptions");
310 debug_struct
311 .field("store_addrs", &self.sanitize_store_addrs())
312 .field("backend_tls", &self.backend_tls)
313 .field("selector", &self.selector)
314 .field("enable_region_failover", &self.enable_region_failover)
315 .field(
316 "allow_region_failover_on_local_wal",
317 &self.allow_region_failover_on_local_wal,
318 )
319 .field("grpc", &self.grpc)
320 .field("http", &self.http)
321 .field("logging", &self.logging)
322 .field("procedure", &self.procedure)
323 .field("failure_detector", &self.failure_detector)
324 .field("datanode", &self.datanode)
325 .field("enable_telemetry", &self.enable_telemetry)
326 .field("data_home", &self.data_home)
327 .field("wal", &self.wal)
328 .field("store_key_prefix", &self.store_key_prefix)
329 .field("max_txn_ops", &self.max_txn_ops)
330 .field("flush_stats_factor", &self.flush_stats_factor)
331 .field("tracing", &self.tracing)
332 .field("backend", &self.backend)
333 .field("event_recorder", &self.event_recorder)
334 .field("stats_persistence", &self.stats_persistence)
335 .field("heartbeat_interval", &self.heartbeat_interval)
336 .field("backend_client", &self.backend_client);
337
338 #[cfg(any(feature = "pg_kvbackend", feature = "mysql_kvbackend"))]
339 debug_struct.field("meta_table_name", &self.meta_table_name);
340
341 #[cfg(feature = "pg_kvbackend")]
342 debug_struct.field("meta_election_lock_id", &self.meta_election_lock_id);
343 #[cfg(feature = "pg_kvbackend")]
344 debug_struct.field("meta_schema_name", &self.meta_schema_name);
345
346 debug_struct
347 .field("node_max_idle_time", &self.node_max_idle_time)
348 .finish()
349 }
350}
351
352const DEFAULT_METASRV_ADDR_PORT: &str = "3002";
353
354impl Default for MetasrvOptions {
355 fn default() -> Self {
356 Self {
357 #[allow(deprecated)]
358 bind_addr: String::new(),
359 #[allow(deprecated)]
360 server_addr: String::new(),
361 store_addrs: vec!["127.0.0.1:2379".to_string()],
362 backend_tls: Some(TlsOption::prefer()),
363 selector: SelectorType::default(),
364 enable_region_failover: false,
365 heartbeat_interval: distributed_time_constants::BASE_HEARTBEAT_INTERVAL,
366 region_failure_detector_initialization_delay: Duration::from_secs(10 * 60),
367 allow_region_failover_on_local_wal: false,
368 grpc: GrpcOptions {
369 bind_addr: format!("127.0.0.1:{}", DEFAULT_METASRV_ADDR_PORT),
370 ..Default::default()
371 },
372 http: HttpOptions::default(),
373 logging: LoggingOptions::default(),
374 procedure: ProcedureConfig {
375 max_retry_times: 12,
376 retry_delay: Duration::from_millis(500),
377 max_metadata_value_size: Some(ReadableSize::kb(1500)),
380 max_running_procedures: 128,
381 },
382 failure_detector: PhiAccrualFailureDetectorOptions::default(),
383 datanode: DatanodeClientOptions::default(),
384 enable_telemetry: true,
385 data_home: DEFAULT_DATA_HOME.to_string(),
386 wal: MetasrvWalConfig::default(),
387 store_key_prefix: String::new(),
388 max_txn_ops: 128,
389 flush_stats_factor: 3,
390 tracing: TracingOptions::default(),
391 memory: MemoryOptions::default(),
392 backend: BackendImpl::EtcdStore,
393 #[cfg(any(feature = "pg_kvbackend", feature = "mysql_kvbackend"))]
394 meta_table_name: common_meta::kv_backend::DEFAULT_META_TABLE_NAME.to_string(),
395 #[cfg(feature = "pg_kvbackend")]
396 meta_election_lock_id: common_meta::kv_backend::DEFAULT_META_ELECTION_LOCK_ID,
397 #[cfg(feature = "pg_kvbackend")]
398 meta_schema_name: None,
399 #[cfg(feature = "pg_kvbackend")]
400 auto_create_schema: true,
401 node_max_idle_time: Duration::from_secs(24 * 60 * 60),
402 event_recorder: EventRecorderOptions::default(),
403 stats_persistence: StatsPersistenceOptions::default(),
404 gc: GcSchedulerOptions::default(),
405 backend_client: BackendClientOptions::default(),
406 }
407 }
408}
409
410impl Configurable for MetasrvOptions {
411 fn env_list_keys() -> Option<&'static [&'static str]> {
412 Some(&[
413 "wal.broker_endpoints",
414 "store_addrs",
415 "event_recorder.event_types",
416 ])
417 }
418}
419
420impl MetasrvOptions {
421 fn sanitize_store_addrs(&self) -> Vec<String> {
422 self.store_addrs
423 .iter()
424 .map(|addr| common_meta::kv_backend::util::sanitize_connection_string(addr))
425 .collect()
426 }
427}
428
429pub struct MetasrvInfo {
430 pub server_addr: String,
431}
432#[derive(Clone)]
433pub struct Context {
434 pub server_addr: String,
435 pub in_memory: ResettableKvBackendRef,
436 pub kv_backend: KvBackendRef,
437 pub leader_cached_kv_backend: ResettableKvBackendRef,
438 pub meta_peer_client: MetaPeerClientRef,
439 pub mailbox: MailboxRef,
440 pub election: Option<ElectionRef>,
441 pub is_infancy: bool,
442 pub table_metadata_manager: TableMetadataManagerRef,
443 pub cache_invalidator: CacheInvalidatorRef,
444 pub leader_region_registry: LeaderRegionRegistryRef,
445 pub topic_stats_registry: TopicStatsRegistryRef,
446 pub heartbeat_interval: Duration,
447 pub gc_enabled: bool,
448 pub is_handshake: bool,
449}
450
451impl Context {
452 pub fn reset_in_memory(&self) {
453 self.in_memory.reset();
454 self.leader_region_registry.reset();
455 }
456
457 pub fn with_handshake(mut self, is_handshake: bool) -> Self {
458 self.is_handshake = is_handshake;
459 self
460 }
461
462 pub fn heartbeat_options_for(&self, role: Role) -> HeartbeatOptions {
463 match role {
464 Role::Datanode => HeartbeatOptions::datanode_from(self.heartbeat_interval),
465 Role::Frontend => HeartbeatOptions::frontend_from(self.heartbeat_interval),
466 Role::Flownode => HeartbeatOptions::flownode_from(self.heartbeat_interval),
467 }
468 }
469}
470
471#[derive(Clone, Copy)]
472pub enum SelectTarget {
473 Datanode,
474 Flownode,
475}
476
477impl Display for SelectTarget {
478 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
479 match self {
480 SelectTarget::Datanode => write!(f, "datanode"),
481 SelectTarget::Flownode => write!(f, "flownode"),
482 }
483 }
484}
485
486#[derive(Clone)]
487pub struct SelectorContext {
488 pub peer_discovery: PeerDiscoveryRef,
489}
490
491pub type SelectorRef = Arc<dyn Selector<Context = SelectorContext, Output = Vec<Peer>>>;
492
493pub struct SelectorFactoryContext {
500 pub metasrv_options: MetasrvOptions,
501 pub meta_peer_client: MetaPeerClientRef,
502 pub in_memory: ResettableKvBackendRef,
503 pub election: Option<ElectionRef>,
504 pub base_selector: SelectorRef,
505}
506
507pub trait SelectorFactory: Send + Sync {
509 fn build(&self, ctx: SelectorFactoryContext) -> SelectorRef;
510}
511
512pub type SelectorFactoryRef = Arc<dyn SelectorFactory>;
514pub type RegionStatAwareSelectorRef =
515 Arc<dyn RegionStatAwareSelector<Context = SelectorContext, Output = Vec<(RegionId, Peer)>>>;
516
517pub struct MetaStateHandler {
518 subscribe_manager: Option<SubscriptionManagerRef>,
519 greptimedb_telemetry_task: Arc<GreptimeDBTelemetryTask>,
520 leader_cached_kv_backend: Arc<LeaderCachedKvBackend>,
521 leadership_change_notifier: LeadershipChangeNotifier,
522 state: StateRef,
523}
524
525impl MetaStateHandler {
526 pub async fn on_leader_start(&self) {
527 self.state.write().unwrap().next_state(become_leader(false));
528
529 if let Err(e) = self.leader_cached_kv_backend.load().await {
530 error!(e; "Failed to load kv into leader cache kv store");
531 } else {
532 self.state.write().unwrap().next_state(become_leader(true));
533 }
534
535 self.leadership_change_notifier
536 .notify_on_leader_start()
537 .await;
538
539 self.greptimedb_telemetry_task.should_report(true);
540 }
541
542 pub async fn on_leader_stop(&self) {
543 self.state.write().unwrap().next_state(become_follower());
544
545 self.leadership_change_notifier
546 .notify_on_leader_stop()
547 .await;
548
549 self.greptimedb_telemetry_task.should_report(false);
551
552 if let Some(sub_manager) = self.subscribe_manager.clone() {
553 info!("Leader changed, un_subscribe all");
554 if let Err(e) = sub_manager.unsubscribe_all() {
555 error!(e; "Failed to un_subscribe all");
556 }
557 }
558 }
559}
560
561pub struct Metasrv {
562 state: StateRef,
563 started: Arc<AtomicBool>,
564 start_time_ms: u64,
565 options: MetasrvOptions,
566 in_memory: ResettableKvBackendRef,
569 kv_backend: KvBackendRef,
570 leader_cached_kv_backend: Arc<LeaderCachedKvBackend>,
571 meta_peer_client: MetaPeerClientRef,
572 selector: SelectorRef,
574 selector_ctx: SelectorContext,
575 flow_selector: SelectorRef,
577 handler_group: RwLock<Option<HeartbeatHandlerGroupRef>>,
578 handler_group_builder: Mutex<Option<HeartbeatHandlerGroupBuilder>>,
579 election: Option<ElectionRef>,
580 procedure_manager: ProcedureManagerRef,
581 mailbox: MailboxRef,
582 ddl_manager: DdlManagerRef,
583 wal_provider: WalProviderRef,
584 table_metadata_manager: TableMetadataManagerRef,
585 runtime_switch_manager: RuntimeSwitchManagerRef,
586 repartition_gc_requirement_manager: RepartitionGcRequirementManagerRef,
587 memory_region_keeper: MemoryRegionKeeperRef,
588 greptimedb_telemetry_task: Arc<GreptimeDBTelemetryTask>,
589 region_migration_manager: RegionMigrationManagerRef,
590 region_supervisor_ticker: Option<RegionSupervisorTickerRef>,
591 cache_invalidator: CacheInvalidatorRef,
592 leader_region_registry: LeaderRegionRegistryRef,
593 topic_stats_registry: TopicStatsRegistryRef,
594 wal_prune_ticker: Option<WalPruneTickerRef>,
595 region_flush_ticker: Option<RegionFlushTickerRef>,
596 table_id_allocator: ResourceIdAllocatorRef,
597 reconciliation_manager: ReconciliationManagerRef,
598 resource_stat: ResourceStatRef,
599 gc_ticker: Option<GcTickerRef>,
600 database_operator: DatabaseOperatorRef,
601
602 plugins: Plugins,
603}
604
605impl Metasrv {
606 pub(crate) async fn ensure_repartition_gc_enabled(&self) -> Result<()> {
607 self.repartition_gc_requirement_manager
608 .ensure_gc_enabled(self.options.gc.enable, &self.procedure_manager)
609 .await
610 }
611
612 pub async fn try_start(&self) -> Result<()> {
613 self.ensure_repartition_gc_enabled().await?;
614 self.try_start_after_gc_check().await
615 }
616
617 pub(crate) async fn try_start_after_gc_check(&self) -> Result<()> {
618 if self
619 .started
620 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
621 .is_err()
622 {
623 warn!("Metasrv already started");
624 return Ok(());
625 }
626
627 let handler_group_builder =
628 self.handler_group_builder
629 .lock()
630 .unwrap()
631 .take()
632 .context(error::UnexpectedSnafu {
633 violated: "expected heartbeat handler group builder",
634 })?;
635 *self.handler_group.write().unwrap() = Some(Arc::new(handler_group_builder.build()?));
636
637 self.table_metadata_manager
639 .init()
640 .await
641 .context(InitMetadataSnafu)?;
642
643 if let Some(election) = self.election() {
644 let procedure_manager = self.procedure_manager.clone();
645 let in_memory = self.in_memory.clone();
646 let leader_cached_kv_backend = self.leader_cached_kv_backend.clone();
647 let subscribe_manager = self.subscription_manager();
648 let mut rx = election.subscribe_leader_change();
649 let greptimedb_telemetry_task = self.greptimedb_telemetry_task.clone();
650 greptimedb_telemetry_task
651 .start()
652 .context(StartTelemetryTaskSnafu)?;
653
654 let mut leadership_change_notifier = LeadershipChangeNotifier::default();
656 leadership_change_notifier.add_listener(self.wal_provider.clone());
657 leadership_change_notifier
658 .add_listener(Arc::new(ProcedureManagerListenerAdapter(procedure_manager)));
659 leadership_change_notifier.add_listener(Arc::new(NodeExpiryListener::new(
660 self.options.node_max_idle_time,
661 self.in_memory.clone(),
662 )));
663 if let Some(region_supervisor_ticker) = &self.region_supervisor_ticker {
664 leadership_change_notifier.add_listener(region_supervisor_ticker.clone() as _);
665 }
666 if let Some(wal_prune_ticker) = &self.wal_prune_ticker {
667 leadership_change_notifier.add_listener(wal_prune_ticker.clone() as _);
668 }
669 if let Some(region_flush_trigger) = &self.region_flush_ticker {
670 leadership_change_notifier.add_listener(region_flush_trigger.clone() as _);
671 }
672 if let Some(gc_ticker) = &self.gc_ticker {
673 leadership_change_notifier.add_listener(gc_ticker.clone() as _);
674 }
675 if let Some(customizer) = self.plugins.get::<LeadershipChangeNotifierCustomizerRef>() {
676 customizer.customize(&mut leadership_change_notifier);
677 }
678
679 let state_handler = MetaStateHandler {
680 greptimedb_telemetry_task,
681 subscribe_manager,
682 state: self.state.clone(),
683 leader_cached_kv_backend: leader_cached_kv_backend.clone(),
684 leadership_change_notifier,
685 };
686 let _handle = common_runtime::spawn_global(async move {
687 loop {
688 match rx.recv().await {
689 Ok(msg) => {
690 in_memory.reset();
691 leader_cached_kv_backend.reset();
692 info!("Leader's cache has bean cleared on leader change: {msg}");
693 match msg {
694 LeaderChangeMessage::Elected(_) => {
695 state_handler.on_leader_start().await;
696 }
697 LeaderChangeMessage::StepDown(leader) => {
698 error!("Leader :{:?} step down", leader);
699
700 state_handler.on_leader_stop().await;
701 }
702 }
703 }
704 Err(RecvError::Closed) => {
705 error!("Not expected, is leader election loop still running?");
706 break;
707 }
708 Err(RecvError::Lagged(_)) => {
709 break;
710 }
711 }
712 }
713
714 state_handler.on_leader_stop().await;
715 });
716
717 {
719 let election = election.clone();
720 let started = self.started.clone();
721 let node_info = self.node_info();
722 let _handle = common_runtime::spawn_global(async move {
723 while started.load(Ordering::Acquire) {
724 let res = election.register_candidate(&node_info).await;
725 if let Err(e) = res {
726 warn!(e; "Metasrv register candidate error");
727 }
728 }
729 });
730 }
731
732 {
734 let election = election.clone();
735 let started = self.started.clone();
736 let _handle = common_runtime::spawn_global(async move {
737 while started.load(Ordering::Acquire) {
738 let res = election.campaign().await;
739 if let Err(e) = res {
740 warn!(e; "Metasrv election error");
741 }
742 election.reset_campaign().await;
743 info!("Metasrv re-initiate election");
744 }
745 info!("Metasrv stopped");
746 });
747 }
748 } else {
749 warn!(
750 "Ensure only one instance of Metasrv is running, as there is no election service."
751 );
752
753 if let Err(e) = self.wal_provider.start().await {
754 error!(e; "Failed to start wal provider");
755 }
756 self.leader_cached_kv_backend
758 .load()
759 .await
760 .context(KvBackendSnafu)?;
761 self.procedure_manager
762 .start()
763 .await
764 .context(StartProcedureManagerSnafu)?;
765 }
766
767 info!("Metasrv started");
768
769 Ok(())
770 }
771
772 pub async fn shutdown(&self) -> Result<()> {
773 if self
774 .started
775 .compare_exchange(true, false, Ordering::AcqRel, Ordering::Acquire)
776 .is_err()
777 {
778 warn!("Metasrv already stopped");
779 return Ok(());
780 }
781
782 self.procedure_manager
783 .stop()
784 .await
785 .context(StopProcedureManagerSnafu)?;
786
787 info!("Metasrv stopped");
788
789 Ok(())
790 }
791
792 pub fn start_time_ms(&self) -> u64 {
793 self.start_time_ms
794 }
795
796 pub fn resource_stat(&self) -> &ResourceStatRef {
797 &self.resource_stat
798 }
799
800 pub fn node_info(&self) -> MetasrvNodeInfo {
801 let build_info = common_version::build_info();
802 MetasrvNodeInfo {
803 addr: self.options().grpc.server_addr.clone(),
804 version: build_info.version.to_string(),
805 git_commit: build_info.commit_short.to_string(),
806 start_time_ms: self.start_time_ms(),
807 total_cpu_millicores: self.resource_stat.get_total_cpu_millicores(),
808 total_memory_bytes: self.resource_stat.get_total_memory_bytes(),
809 cpu_usage_millicores: self.resource_stat.get_cpu_usage_millicores(),
810 memory_usage_bytes: self.resource_stat.get_memory_usage_bytes(),
811 hostname: hostname::get()
812 .unwrap_or_default()
813 .to_string_lossy()
814 .to_string(),
815 }
816 }
817
818 pub(crate) async fn lookup_datanode_peer(&self, peer_id: u64) -> Result<Option<Peer>> {
821 discovery::utils::alive_datanode(
822 &DefaultSystemTimer,
823 self.meta_peer_client.as_ref(),
824 peer_id,
825 default_distributed_time_constants().datanode_lease,
826 )
827 .await
828 }
829
830 pub fn options(&self) -> &MetasrvOptions {
831 &self.options
832 }
833
834 pub fn in_memory(&self) -> &ResettableKvBackendRef {
835 &self.in_memory
836 }
837
838 pub fn kv_backend(&self) -> &KvBackendRef {
839 &self.kv_backend
840 }
841
842 pub fn meta_peer_client(&self) -> &MetaPeerClientRef {
843 &self.meta_peer_client
844 }
845
846 pub fn selector(&self) -> &SelectorRef {
847 &self.selector
848 }
849
850 pub fn selector_ctx(&self) -> &SelectorContext {
851 &self.selector_ctx
852 }
853
854 pub fn flow_selector(&self) -> &SelectorRef {
855 &self.flow_selector
856 }
857
858 pub fn handler_group(&self) -> Option<HeartbeatHandlerGroupRef> {
859 self.handler_group.read().unwrap().clone()
860 }
861
862 pub fn election(&self) -> Option<&ElectionRef> {
863 self.election.as_ref()
864 }
865
866 pub fn mailbox(&self) -> &MailboxRef {
867 &self.mailbox
868 }
869
870 pub fn ddl_manager(&self) -> &DdlManagerRef {
871 &self.ddl_manager
872 }
873
874 pub fn procedure_manager(&self) -> &ProcedureManagerRef {
875 &self.procedure_manager
876 }
877
878 pub fn table_metadata_manager(&self) -> &TableMetadataManagerRef {
879 &self.table_metadata_manager
880 }
881
882 pub fn runtime_switch_manager(&self) -> &RuntimeSwitchManagerRef {
883 &self.runtime_switch_manager
884 }
885
886 pub fn memory_region_keeper(&self) -> &MemoryRegionKeeperRef {
887 &self.memory_region_keeper
888 }
889
890 pub fn region_migration_manager(&self) -> &RegionMigrationManagerRef {
891 &self.region_migration_manager
892 }
893
894 pub fn publish(&self) -> Option<PublisherRef> {
895 self.plugins.get::<PublisherRef>()
896 }
897
898 pub fn subscription_manager(&self) -> Option<SubscriptionManagerRef> {
899 self.plugins.get::<SubscriptionManagerRef>()
900 }
901
902 pub fn table_id_allocator(&self) -> &ResourceIdAllocatorRef {
903 &self.table_id_allocator
904 }
905
906 pub fn reconciliation_manager(&self) -> &ReconciliationManagerRef {
907 &self.reconciliation_manager
908 }
909
910 pub fn database_operator(&self) -> &DatabaseOperatorRef {
911 &self.database_operator
912 }
913
914 pub fn plugins(&self) -> &Plugins {
915 &self.plugins
916 }
917
918 pub fn started(&self) -> Arc<AtomicBool> {
919 self.started.clone()
920 }
921
922 pub fn gc_ticker(&self) -> Option<GcTickerRef> {
923 self.gc_ticker.as_ref().cloned()
924 }
925
926 #[inline]
927 pub fn new_ctx(&self) -> Context {
928 let server_addr = self.options().grpc.server_addr.clone();
929 let in_memory = self.in_memory.clone();
930 let kv_backend = self.kv_backend.clone();
931 let leader_cached_kv_backend = self.leader_cached_kv_backend.clone();
932 let meta_peer_client = self.meta_peer_client.clone();
933 let mailbox = self.mailbox.clone();
934 let election = self.election.clone();
935 let table_metadata_manager = self.table_metadata_manager.clone();
936 let cache_invalidator = self.cache_invalidator.clone();
937 let leader_region_registry = self.leader_region_registry.clone();
938 let topic_stats_registry = self.topic_stats_registry.clone();
939
940 Context {
941 server_addr,
942 in_memory,
943 kv_backend,
944 leader_cached_kv_backend,
945 meta_peer_client,
946 mailbox,
947 election,
948 is_infancy: false,
949 table_metadata_manager,
950 cache_invalidator,
951 leader_region_registry,
952 topic_stats_registry,
953 heartbeat_interval: self.options().heartbeat_interval,
954 gc_enabled: self.options().gc.enable,
955 is_handshake: false,
956 }
957 }
958}
959
960#[cfg(test)]
961mod tests {
962 use common_event_recorder::EventTypeFilter;
963
964 use super::*;
965 use crate::metasrv::MetasrvNodeInfo;
966
967 #[test]
968 fn test_deserialize_metasrv_node_info() {
969 let str = r#"{"addr":"127.0.0.1:4002","version":"0.1.0","git_commit":"1234567890","start_time_ms":1715145600}"#;
970 let node_info: MetasrvNodeInfo = serde_json::from_str(str).unwrap();
971 assert_eq!(node_info.addr, "127.0.0.1:4002");
972 assert_eq!(node_info.version, "0.1.0");
973 assert_eq!(node_info.git_commit, "1234567890");
974 assert_eq!(node_info.start_time_ms, 1715145600);
975 }
976
977 #[test]
978 fn test_metasrv_event_recorder_options_preserve_event_type_filter_semantics() {
979 let all = MetasrvOptions::default().event_recorder;
980 let none: EventRecorderOptions = toml::from_str("ttl = '90d'\nevent_types = []").unwrap();
981 let selected: EventRecorderOptions =
982 toml::from_str("ttl = '90d'\nevent_types = ['create_database']").unwrap();
983
984 assert!(all.event_types.allows("future_event"));
985 assert_eq!(
986 none.event_types.as_ref(),
987 &EventTypeFilter::Only(Default::default())
988 );
989 assert!(selected.event_types.allows("create_database"));
990 assert!(!selected.event_types.allows("drop_database"));
991 }
992}