1use 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 serve_state: Arc<Mutex<Option<oneshot::Receiver<Result<()>>>>>,
75
76 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 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 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 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 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 .accept_http1(true)
260 .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 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 .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}