Skip to main content

meta_srv/
bootstrap.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
15use std::net::SocketAddr;
16use std::sync::Arc;
17
18use api::v1::meta::cluster_server::ClusterServer;
19use api::v1::meta::config_server::ConfigServer;
20use api::v1::meta::heartbeat_server::HeartbeatServer;
21use api::v1::meta::procedure_service_server::ProcedureServiceServer;
22use api::v1::meta::store_server::StoreServer;
23use common_base::Plugins;
24use common_config::Configurable;
25#[cfg(any(feature = "pg_kvbackend", feature = "mysql_kvbackend"))]
26use common_meta::distributed_time_constants::META_LEASE_SECS;
27use common_meta::election::etcd::EtcdElection;
28use common_meta::kv_backend::chroot::ChrootKvBackend;
29use common_meta::kv_backend::etcd::EtcdStore;
30use common_meta::kv_backend::memory::MemoryKvBackend;
31use common_meta::kv_backend::{KvBackendRef, ResettableKvBackendRef};
32use common_telemetry::info;
33use either::Either;
34use servers::configurator::GrpcRouterConfiguratorRef;
35use servers::http::{ExtraHttpRouterProviders, HttpServer, HttpServerBuilder};
36use servers::metrics_handler::MetricsHandler;
37use servers::server::Server;
38use snafu::ResultExt;
39use tokio::net::TcpListener;
40use tokio::sync::mpsc::{self, Receiver, Sender};
41use tokio::sync::{Mutex, oneshot};
42use tonic::codec::CompressionEncoding;
43use tonic::transport::server::{Router, TcpIncoming};
44
45use crate::cluster::{MetaPeerClientBuilder, MetaPeerClientRef};
46use crate::error::OtherSnafu;
47use crate::metasrv::builder::MetasrvBuilder;
48use crate::metasrv::{
49    BackendImpl, ElectionRef, Metasrv, MetasrvOptions, SelectTarget, SelectorFactoryContext,
50    SelectorFactoryRef, SelectorRef,
51};
52use crate::selector::lease_based::LeaseBasedSelector;
53use crate::selector::load_based::LoadBasedSelector;
54use crate::selector::round_robin::RoundRobinSelector;
55use crate::selector::weight_compute::RegionNumsBasedWeightCompute;
56use crate::selector::{Selector, SelectorType};
57use crate::service::admin;
58use crate::service::admin::admin_axum_router;
59use crate::utils::etcd::create_etcd_client_with_tls;
60use crate::{Result, error};
61
62pub struct MetasrvInstance {
63    metasrv: Arc<Metasrv>,
64
65    http_server: Either<Option<HttpServerBuilder>, HttpServer>,
66
67    opts: MetasrvOptions,
68
69    signal_sender: Option<Sender<()>>,
70
71    plugins: Plugins,
72
73    /// gRPC serving state receiver. Only present if the gRPC server is started.
74    serve_state: Arc<Mutex<Option<oneshot::Receiver<Result<()>>>>>,
75
76    /// gRPC bind addr
77    bind_addr: Option<SocketAddr>,
78}
79
80impl MetasrvInstance {
81    pub async fn new(metasrv: Metasrv) -> Result<MetasrvInstance> {
82        let opts = metasrv.options().clone();
83        let plugins = metasrv.plugins().clone();
84        let metasrv = Arc::new(metasrv);
85
86        // Wire up the admin_axum_router as an extra router
87        let extra_routers = admin_axum_router(metasrv.clone());
88
89        let mut builder = HttpServerBuilder::new(opts.http.clone())
90            .with_metrics_handler(MetricsHandler)
91            .with_greptime_config_options(opts.to_toml().context(error::TomlFormatSnafu)?);
92        builder = builder.with_extra_router(extra_routers);
93
94        // put metasrv into plugins for later use
95        plugins.insert::<Arc<Metasrv>>(metasrv.clone());
96        Ok(MetasrvInstance {
97            metasrv,
98            http_server: Either::Left(Some(builder)),
99            opts,
100            signal_sender: None,
101            plugins,
102            serve_state: Default::default(),
103            bind_addr: None,
104        })
105    }
106
107    pub async fn start(&mut self) -> Result<()> {
108        self.metasrv.ensure_repartition_gc_enabled().await?;
109
110        if let Some(builder) = self.http_server.as_mut().left()
111            && let Some(mut builder) = builder.take()
112        {
113            if let Some(providers) = self.plugins.get::<ExtraHttpRouterProviders>() {
114                for provider in providers.iter() {
115                    builder = builder.with_extra_router(provider.router());
116                }
117            }
118            let mut server = builder.build();
119
120            let addr = self.opts.http.addr.parse().context(error::ParseAddrSnafu {
121                addr: &self.opts.http.addr,
122            })?;
123            info!("starting http server at {}", addr);
124            server.start(addr).await.context(error::StartHttpSnafu)?;
125
126            self.http_server = Either::Right(server);
127        } else {
128            // If the http server builder is not present, the Metasrv has to be called "start"
129            // already, regardless of the startup was successful or not. Return an `Ok` here for
130            // simplicity.
131            return Ok(());
132        };
133
134        self.metasrv.try_start_after_gc_check().await?;
135
136        let (tx, rx) = mpsc::channel::<()>(1);
137
138        self.signal_sender = Some(tx);
139
140        // Start gRPC server with admin services for backward compatibility
141        let mut router = router(self.metasrv.clone());
142        if let Some(configurator) = self
143            .metasrv
144            .plugins()
145            .get::<GrpcRouterConfiguratorRef<()>>()
146        {
147            router = configurator
148                .configure_grpc_router(router, ())
149                .await
150                .context(OtherSnafu)?;
151        }
152
153        let (serve_state_tx, serve_state_rx) = oneshot::channel();
154
155        let socket_addr =
156            bootstrap_metasrv_with_router(&self.opts.grpc.bind_addr, router, serve_state_tx, rx)
157                .await?;
158        self.bind_addr = Some(socket_addr);
159
160        *self.serve_state.lock().await = Some(serve_state_rx);
161        Ok(())
162    }
163
164    pub async fn shutdown(&self) -> Result<()> {
165        if let Some(mut rx) = self.serve_state.lock().await.take()
166            && let Ok(Err(err)) = rx.try_recv()
167        {
168            common_telemetry::error!(err; "Metasrv start failed")
169        }
170        if let Some(signal) = &self.signal_sender {
171            signal
172                .send(())
173                .await
174                .context(error::SendShutdownSignalSnafu)?;
175        }
176        self.metasrv.shutdown().await?;
177
178        if let Some(http_server) = self.http_server.as_ref().right() {
179            http_server
180                .shutdown()
181                .await
182                .context(error::ShutdownServerSnafu {
183                    server: http_server.name(),
184                })?;
185        }
186        Ok(())
187    }
188
189    pub fn plugins(&self) -> Plugins {
190        self.plugins.clone()
191    }
192
193    pub fn get_inner(&self) -> &Metasrv {
194        &self.metasrv
195    }
196    pub fn bind_addr(&self) -> &Option<SocketAddr> {
197        &self.bind_addr
198    }
199
200    pub fn mut_http_server(&mut self) -> &mut Either<Option<HttpServerBuilder>, HttpServer> {
201        &mut self.http_server
202    }
203
204    pub fn http_server(&self) -> Option<&HttpServer> {
205        self.http_server.as_ref().right()
206    }
207}
208
209pub async fn bootstrap_metasrv_with_router(
210    bind_addr: &str,
211    router: Router,
212    serve_state_tx: oneshot::Sender<Result<()>>,
213    mut shutdown_rx: Receiver<()>,
214) -> Result<SocketAddr> {
215    let listener = TcpListener::bind(bind_addr)
216        .await
217        .context(error::TcpBindSnafu { addr: bind_addr })?;
218
219    let real_bind_addr = listener
220        .local_addr()
221        .context(error::TcpBindSnafu { addr: bind_addr })?;
222
223    info!("gRPC server is bound to: {}", real_bind_addr);
224
225    let incoming = TcpIncoming::from(listener).with_nodelay(Some(true));
226
227    let _handle = common_runtime::spawn_global(async move {
228        let result = router
229            .serve_with_incoming_shutdown(incoming, async {
230                let _ = shutdown_rx.recv().await;
231            })
232            .await
233            .inspect_err(|err| common_telemetry::error!(err;"Failed to start metasrv"))
234            .context(error::StartGrpcSnafu);
235        let _ = serve_state_tx.send(result);
236    });
237
238    Ok(real_bind_addr)
239}
240
241#[macro_export]
242macro_rules! add_compressed_service {
243    ($builder:expr, $server:expr, $grpc_config:expr) => {
244        $builder.add_service(
245            $server
246                .accept_compressed(CompressionEncoding::Gzip)
247                .accept_compressed(CompressionEncoding::Zstd)
248                .send_compressed(CompressionEncoding::Gzip)
249                .send_compressed(CompressionEncoding::Zstd)
250                .max_decoding_message_size($grpc_config.max_recv_message_size)
251                .max_encoding_message_size($grpc_config.max_send_message_size),
252        )
253    };
254}
255
256pub fn router(metasrv: Arc<Metasrv>) -> Router {
257    let mut router = tonic::transport::Server::builder()
258        // for admin services
259        .accept_http1(true)
260        // For quick network failures detection.
261        .http2_keepalive_interval(Some(metasrv.options().grpc.http2_keep_alive_interval))
262        .http2_keepalive_timeout(Some(metasrv.options().grpc.http2_keep_alive_timeout));
263    let grpc_config = metasrv.options().grpc.as_config();
264    let router = add_compressed_service!(
265        router,
266        HeartbeatServer::from_arc(metasrv.clone()),
267        grpc_config
268    );
269    let router =
270        add_compressed_service!(router, StoreServer::from_arc(metasrv.clone()), grpc_config);
271    let router = add_compressed_service!(
272        router,
273        ClusterServer::from_arc(metasrv.clone()),
274        grpc_config
275    );
276    let router = add_compressed_service!(
277        router,
278        ProcedureServiceServer::from_arc(metasrv.clone()),
279        grpc_config
280    );
281    let router =
282        add_compressed_service!(router, ConfigServer::from_arc(metasrv.clone()), grpc_config);
283    router.add_service(admin::make_admin_service(metasrv))
284}
285
286pub async fn metasrv_builder(
287    opts: &MetasrvOptions,
288    plugins: &Plugins,
289    kv_backend: Option<KvBackendRef>,
290) -> Result<MetasrvBuilder> {
291    let (mut kv_backend, election) = match (kv_backend, &opts.backend) {
292        (Some(kv_backend), _) => (kv_backend, None),
293        (None, BackendImpl::MemoryStore) => (Arc::new(MemoryKvBackend::new()) as _, None),
294        (None, BackendImpl::EtcdStore) => {
295            let etcd_client = create_etcd_client_with_tls(
296                &opts.store_addrs,
297                &opts.backend_client,
298                opts.backend_tls.as_ref(),
299            )
300            .await?;
301            let kv_backend = EtcdStore::with_etcd_client(etcd_client.clone(), opts.max_txn_ops);
302            let election = EtcdElection::with_etcd_client(
303                &opts.grpc.server_addr,
304                etcd_client,
305                opts.store_key_prefix.clone(),
306            )
307            .await
308            .context(error::KvBackendSnafu)?;
309
310            (kv_backend, Some(election))
311        }
312        #[cfg(feature = "pg_kvbackend")]
313        (None, BackendImpl::PostgresStore) => {
314            use std::time::Duration;
315
316            use common_meta::distributed_time_constants::POSTGRES_KEEP_ALIVE_SECS;
317            use common_meta::election::CANDIDATE_LEASE_SECS;
318            use deadpool_postgres::{Config, ManagerConfig, RecyclingMethod};
319
320            use crate::utils::postgres::{build_postgres_election, build_postgres_kv_backend};
321
322            let candidate_lease_ttl = Duration::from_secs(CANDIDATE_LEASE_SECS);
323            let meta_lease_ttl = Duration::from_secs(META_LEASE_SECS);
324
325            let mut cfg = Config::new();
326            cfg.keepalives = Some(true);
327            cfg.keepalives_idle = Some(Duration::from_secs(POSTGRES_KEEP_ALIVE_SECS));
328            cfg.manager = Some(ManagerConfig {
329                recycling_method: RecyclingMethod::Verified,
330            });
331
332            let election = build_postgres_election(
333                &opts.store_addrs,
334                Some(cfg.clone()),
335                opts.backend_tls.clone(),
336                opts.grpc.server_addr.clone(),
337                opts.store_key_prefix.clone(),
338                candidate_lease_ttl,
339                meta_lease_ttl,
340                opts.meta_schema_name.as_deref(),
341                &opts.meta_table_name,
342                opts.meta_election_lock_id,
343            )
344            .await?;
345
346            let kv_backend = build_postgres_kv_backend(
347                &opts.store_addrs,
348                Some(cfg),
349                opts.backend_tls.clone(),
350                opts.meta_schema_name.as_deref(),
351                &opts.meta_table_name,
352                opts.max_txn_ops,
353                opts.auto_create_schema,
354            )
355            .await?;
356
357            (kv_backend, Some(election))
358        }
359        #[cfg(feature = "mysql_kvbackend")]
360        (None, BackendImpl::MysqlStore) => {
361            use std::time::Duration;
362
363            use common_meta::election::CANDIDATE_LEASE_SECS;
364
365            use crate::utils::mysql::{build_mysql_election, build_mysql_kv_backend};
366
367            let kv_backend = build_mysql_kv_backend(
368                &opts.store_addrs,
369                opts.backend_tls.as_ref(),
370                &opts.meta_table_name,
371                opts.max_txn_ops,
372            )
373            .await?;
374            // Since election will acquire a lock of the table, we need a separate table for election.
375            let election_table_name = opts.meta_table_name.clone() + "_election";
376            let innode_lock_wait_timeout = Duration::from_secs(META_LEASE_SECS / 2);
377            let meta_lease_ttl = Duration::from_secs(META_LEASE_SECS);
378            let candidate_lease_ttl = Duration::from_secs(CANDIDATE_LEASE_SECS);
379
380            let election = build_mysql_election(
381                &opts.store_addrs,
382                opts.backend_tls.as_ref(),
383                opts.grpc.server_addr.clone(),
384                opts.store_key_prefix.clone(),
385                candidate_lease_ttl,
386                meta_lease_ttl,
387                &election_table_name,
388                innode_lock_wait_timeout,
389            )
390            .await?;
391            (kv_backend, Some(election))
392        }
393    };
394
395    if !opts.store_key_prefix.is_empty() {
396        info!(
397            "using chroot kv backend with prefix: {prefix}",
398            prefix = opts.store_key_prefix
399        );
400        kv_backend = Arc::new(ChrootKvBackend::new(
401            opts.store_key_prefix.clone().into_bytes(),
402            kv_backend,
403        ))
404    }
405
406    let in_memory = Arc::new(MemoryKvBackend::new()) as ResettableKvBackendRef;
407    let meta_peer_client = build_default_meta_peer_client(&election, &in_memory);
408
409    let base_selector: Arc<
410        dyn Selector<
411                Context = crate::metasrv::SelectorContext,
412                Output = Vec<common_meta::peer::Peer>,
413            >,
414    > = match opts.selector {
415        SelectorType::LoadBased => Arc::new(LoadBasedSelector::new(
416            RegionNumsBasedWeightCompute,
417            meta_peer_client.clone(),
418        )) as SelectorRef,
419        SelectorType::LeaseBased => Arc::new(LeaseBasedSelector) as SelectorRef,
420        SelectorType::RoundRobin => {
421            Arc::new(RoundRobinSelector::new(SelectTarget::Datanode)) as SelectorRef
422        }
423    };
424    info!(
425        "Using selector from options, selector type: {}",
426        opts.selector.as_ref()
427    );
428
429    let selector = if let Some(factory) = plugins.get::<SelectorFactoryRef>() {
430        info!("Building selector from plugin factory");
431        factory.build(SelectorFactoryContext {
432            metasrv_options: opts.clone(),
433            meta_peer_client: meta_peer_client.clone(),
434            in_memory: in_memory.clone(),
435            election: election.clone(),
436            base_selector,
437        })
438    } else {
439        base_selector
440    };
441
442    Ok(MetasrvBuilder::new()
443        .options(opts.clone())
444        .kv_backend(kv_backend)
445        .in_memory(in_memory)
446        .selector(selector)
447        .election(election)
448        .meta_peer_client(meta_peer_client))
449}
450
451pub(crate) fn build_default_meta_peer_client(
452    election: &Option<ElectionRef>,
453    in_memory: &ResettableKvBackendRef,
454) -> MetaPeerClientRef {
455    MetaPeerClientBuilder::default()
456        .election(election.clone())
457        .in_memory(in_memory.clone())
458        .build()
459        .map(Arc::new)
460        // Safety: all required fields set at initialization
461        .unwrap()
462}
463
464#[cfg(test)]
465mod tests {
466    use std::sync::atomic::{AtomicBool, Ordering};
467
468    use axum::body::Body;
469    use axum::http::{Request, StatusCode};
470    use axum::routing::get;
471    use common_meta::kv_backend::memory::MemoryKvBackend;
472    use servers::http::{ExtraHttpRouterProvider, ExtraHttpRouterProviderRef};
473    use tower::ServiceExt;
474
475    use super::*;
476    use crate::metasrv::{SelectorFactory, SelectorFactoryContext};
477    use crate::procedure::repartition::gc_requirement::RepartitionGcRequirementManager;
478
479    struct RecordingSelectorFactory {
480        called: Arc<AtomicBool>,
481    }
482
483    impl SelectorFactory for RecordingSelectorFactory {
484        fn build(&self, ctx: SelectorFactoryContext) -> SelectorRef {
485            self.called.store(true, Ordering::Relaxed);
486            ctx.base_selector
487        }
488    }
489
490    #[tokio::test]
491    async fn metasrv_builder_builds_load_based_selector_from_plugin_factory() {
492        let called = Arc::new(AtomicBool::new(false));
493        let plugins = Plugins::new();
494        plugins.insert(Arc::new(RecordingSelectorFactory {
495            called: called.clone(),
496        }) as SelectorFactoryRef);
497        let opts = MetasrvOptions {
498            selector: SelectorType::LoadBased,
499            ..Default::default()
500        };
501
502        metasrv_builder(
503            &opts,
504            &plugins,
505            Some(Arc::new(MemoryKvBackend::new()) as KvBackendRef),
506        )
507        .await
508        .unwrap();
509
510        assert!(called.load(Ordering::Relaxed));
511    }
512
513    #[tokio::test]
514    async fn gc_requirement_is_checked_before_http_server_start() {
515        let kv_backend: KvBackendRef = Arc::new(MemoryKvBackend::new());
516        RepartitionGcRequirementManager::new(kv_backend.clone())
517            .require_gc()
518            .await
519            .unwrap();
520
521        let opts = MetasrvOptions {
522            enable_telemetry: false,
523            ..Default::default()
524        };
525        let metasrv = MetasrvBuilder::new()
526            .options(opts)
527            .kv_backend(kv_backend)
528            .build()
529            .await
530            .unwrap();
531        let mut instance = MetasrvInstance::new(metasrv).await.unwrap();
532
533        let err = instance.start().await.unwrap_err();
534        assert!(matches!(err, error::Error::RepartitionGcRequired { .. }));
535        assert!(instance.http_server().is_none());
536        assert!(matches!(instance.mut_http_server(), Either::Left(Some(_))));
537    }
538
539    struct TestExtraHttpRouterProvider;
540
541    impl ExtraHttpRouterProvider for TestExtraHttpRouterProvider {
542        fn router(&self) -> axum::Router {
543            axum::Router::new().route(
544                "/test-extra-http-router",
545                get(|| async { StatusCode::NO_CONTENT }),
546            )
547        }
548    }
549
550    #[tokio::test]
551    async fn test_metasrv_add_extra_http_router() -> Result<()> {
552        let mut opts = MetasrvOptions::default();
553        opts.grpc.bind_addr = "127.0.0.1:0".to_string();
554        opts.http.addr = "127.0.0.1:0".to_string();
555        opts.enable_telemetry = false;
556
557        let plugins = Plugins::new();
558        let mut providers = ExtraHttpRouterProviders::new();
559        providers.add(Arc::new(TestExtraHttpRouterProvider) as ExtraHttpRouterProviderRef);
560        plugins.insert(providers);
561
562        let metasrv = MetasrvBuilder::new()
563            .options(opts)
564            .plugins(plugins)
565            .build()
566            .await?;
567        let mut instance = MetasrvInstance::new(metasrv).await?;
568        instance.start().await?;
569
570        let response = instance
571            .http_server()
572            .unwrap()
573            .make_app()
574            .oneshot(
575                Request::get("/test-extra-http-router")
576                    .body(Body::empty())
577                    .unwrap(),
578            )
579            .await
580            .unwrap();
581
582        assert_eq!(StatusCode::NO_CONTENT, response.status());
583        instance.shutdown().await?;
584        Ok(())
585    }
586}