Skip to main content

meta_srv/
metasrv.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
15pub 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// The datastores that implements metadata kvbackend.
94#[derive(Clone, Debug, PartialEq, Serialize, Default, Deserialize, ValueEnum)]
95#[serde(rename_all = "snake_case")]
96pub enum BackendImpl {
97    // Etcd as metadata storage.
98    #[default]
99    EtcdStore,
100    // In memory metadata storage - mostly used for testing.
101    MemoryStore,
102    #[cfg(feature = "pg_kvbackend")]
103    // Postgres as metadata storage.
104    PostgresStore,
105    #[cfg(feature = "mysql_kvbackend")]
106    // MySql as metadata storage.
107    MysqlStore,
108}
109
110/// Configuration options for the stats persistence.
111#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
112pub struct StatsPersistenceOptions {
113    /// TTL for the stats table that will be used to store the stats.
114    #[serde(with = "humantime_serde")]
115    pub ttl: Duration,
116    /// The interval to persist the stats.
117    #[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/// Heartbeat configuration for a single node type.
131#[derive(Clone, PartialEq, Serialize, Deserialize, Debug)]
132#[serde(default)]
133pub struct HeartbeatOptions {
134    /// Heartbeat interval.
135    #[serde(with = "humantime_serde")]
136    pub interval: Duration,
137    /// Retry interval when heartbeat connection fails.
138    #[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    /// The address the server listens on.
209    #[deprecated(note = "Use grpc.bind_addr instead")]
210    pub bind_addr: String,
211    /// The address the server advertises to the clients.
212    #[deprecated(note = "Use grpc.server_addr instead")]
213    pub server_addr: String,
214    /// The address of the store, e.g., etcd.
215    pub store_addrs: Vec<String>,
216    /// TLS configuration for kv store backend (PostgreSQL/MySQL)
217    /// Only applicable when using PostgreSQL or MySQL as the metadata store
218    #[serde(default)]
219    pub backend_tls: Option<TlsOption>,
220    /// The backend client options.
221    /// Currently, only applicable when using etcd as the metadata store.
222    #[serde(default)]
223    pub backend_client: BackendClientOptions,
224    /// The type of selector.
225    pub selector: SelectorType,
226    /// Whether to enable region failover.
227    pub enable_region_failover: bool,
228    /// The base heartbeat interval.
229    ///
230    /// This value is used to calculate the distributed time constants for components.
231    /// e.g., the region lease time is `heartbeat_interval * 3 + Duration::from_secs(1)`.
232    #[serde(with = "humantime_serde")]
233    pub heartbeat_interval: Duration,
234    /// The delay before starting region failure detection.
235    /// This delay helps prevent Metasrv from triggering unnecessary region failovers before all Datanodes are fully started.
236    /// Especially useful when the cluster is not deployed with GreptimeDB Operator and maintenance mode is not enabled.
237    #[serde(with = "humantime_serde")]
238    pub region_failure_detector_initialization_delay: Duration,
239    /// Whether to allow region failover on local WAL.
240    ///
241    /// If it's true, the region failover will be allowed even if the local WAL is used.
242    /// Note that this option is not recommended to be set to true, because it may lead to data loss during failover.
243    pub allow_region_failover_on_local_wal: bool,
244    pub grpc: GrpcOptions,
245    /// The HTTP server options.
246    pub http: HttpOptions,
247    /// The logging options.
248    pub logging: LoggingOptions,
249    /// The procedure options.
250    pub procedure: ProcedureConfig,
251    /// The failure detector options.
252    pub failure_detector: PhiAccrualFailureDetectorOptions,
253    /// The datanode options.
254    pub datanode: DatanodeClientOptions,
255    /// Whether to enable telemetry.
256    pub enable_telemetry: bool,
257    /// The data home directory.
258    pub data_home: String,
259    /// The WAL options.
260    pub wal: MetasrvWalConfig,
261    /// The store key prefix. If it is not empty, all keys in the store will be prefixed with it.
262    /// This is useful when multiple metasrv clusters share the same store.
263    pub store_key_prefix: String,
264    /// The max operations per txn
265    ///
266    /// This value is usually limited by which store is used for the `KvBackend`.
267    /// For example, if using etcd, this value should ensure that it is less than
268    /// or equal to the `--max-txn-ops` option value of etcd.
269    ///
270    /// TODO(jeremy): Currently, this option only affects the etcd store, but it may
271    /// also affect other stores in the future. In other words, each store needs to
272    /// limit the number of operations in a txn because an infinitely large txn could
273    /// potentially block other operations.
274    pub max_txn_ops: usize,
275    /// The factor that determines how often statistics should be flushed,
276    /// based on the number of received heartbeats. When the number of heartbeats
277    /// reaches this factor, a flush operation is triggered.
278    pub flush_stats_factor: usize,
279    /// The tracing options.
280    pub tracing: TracingOptions,
281    /// The memory options.
282    pub memory: MemoryOptions,
283    /// The datastore for kv metadata.
284    pub backend: BackendImpl,
285    #[cfg(any(feature = "pg_kvbackend", feature = "mysql_kvbackend"))]
286    /// Table name of rds kv backend.
287    pub meta_table_name: String,
288    #[cfg(feature = "pg_kvbackend")]
289    /// Lock id for meta kv election. Only effect when using pg_kvbackend.
290    pub meta_election_lock_id: u64,
291    #[cfg(feature = "pg_kvbackend")]
292    /// Optional PostgreSQL schema for metadata table (defaults to current search_path if empty).
293    pub meta_schema_name: Option<String>,
294    #[cfg(feature = "pg_kvbackend")]
295    /// Automatically create PostgreSQL schema if it doesn't exist (default: true).
296    pub auto_create_schema: bool,
297    #[serde(with = "humantime_serde")]
298    pub node_max_idle_time: Duration,
299    /// The event recorder options.
300    pub event_recorder: EventRecorderOptions,
301    /// The stats persistence options.
302    pub stats_persistence: StatsPersistenceOptions,
303    /// The GC scheduler options.
304    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                // The etcd the maximum size of any request is 1.5 MiB
378                // 1500KiB = 1536KiB (1.5MiB) - 36KiB (reserved size of key)
379                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
493/// Context passed to a selector factory during metasrv bootstrap.
494///
495/// The factory runs after bootstrap has constructed the selector configured by
496/// [`MetasrvOptions::selector`], so plugins can either decorate `base_selector` or
497/// build a completely different selector using bootstrap-only dependencies like
498/// [`MetaPeerClientRef`].
499pub 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
507/// Builds the final datanode selector metasrv should use.
508pub trait SelectorFactory: Send + Sync {
509    fn build(&self, ctx: SelectorFactoryContext) -> SelectorRef;
510}
511
512/// Shared selector factory plugin registered through [`common_base::Plugins`].
513pub 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        // Suspends reporting.
550        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    // It is only valid at the leader node and is used to temporarily
567    // store some data that will not be persisted.
568    in_memory: ResettableKvBackendRef,
569    kv_backend: KvBackendRef,
570    leader_cached_kv_backend: Arc<LeaderCachedKvBackend>,
571    meta_peer_client: MetaPeerClientRef,
572    // The selector is used to select a target datanode.
573    selector: SelectorRef,
574    selector_ctx: SelectorContext,
575    // The flow selector is used to select a target flownode.
576    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        // Creates default schema if not exists
638        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            // Builds leadership change notifier.
655            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            // Register candidate and keep lease in background.
718            {
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            // Campaign
733            {
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            // Always load kv into cached kv store.
757            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    /// Looks up a datanode peer by peer_id, returning it only when it's alive.
819    /// A datanode is considered alive when it's still within the lease period.
820    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}