1use std::net::SocketAddr;
16use std::sync::Arc;
17
18use common_config::Configurable;
19use servers::grpc::GrpcServer;
20use servers::grpc::builder::GrpcServerBuilder;
21use servers::http::HttpServerBuilder;
22use servers::metrics_handler::MetricsHandler;
23use servers::server::{ServerHandler, ServerHandlers};
24use snafu::ResultExt;
25
26use crate::config::DatanodeOptions;
27use crate::error::{ParseAddrSnafu, Result, TomlFormatSnafu};
28use crate::region_server::RegionServer;
29
30pub struct DatanodeServiceBuilder<'a> {
31 opts: &'a DatanodeOptions,
32 grpc_server: Option<GrpcServer>,
33 enable_http_service: bool,
34}
35
36impl<'a> DatanodeServiceBuilder<'a> {
37 pub fn new(opts: &'a DatanodeOptions) -> Self {
38 Self {
39 opts,
40 grpc_server: None,
41 enable_http_service: false,
42 }
43 }
44
45 pub fn with_grpc_server(self, grpc_server: GrpcServer) -> Self {
46 Self {
47 grpc_server: Some(grpc_server),
48 ..self
49 }
50 }
51
52 pub fn with_default_grpc_server(mut self, region_server: &RegionServer) -> Self {
53 let grpc_server = Self::grpc_server_builder(self.opts, region_server).build();
54 self.grpc_server = Some(grpc_server);
55 self
56 }
57
58 pub fn enable_http_service(self) -> Self {
59 Self {
60 enable_http_service: true,
61 ..self
62 }
63 }
64
65 pub fn build(mut self) -> Result<ServerHandlers> {
66 let handlers = ServerHandlers::default();
67
68 if let Some(grpc_server) = self.grpc_server.take() {
69 let addr: SocketAddr = self.opts.grpc.bind_addr.parse().context(ParseAddrSnafu {
70 addr: &self.opts.grpc.bind_addr,
71 })?;
72 let handler: ServerHandler = (Box::new(grpc_server), addr);
73 handlers.insert(handler);
74 }
75
76 if self.enable_http_service {
77 let http_server = HttpServerBuilder::new(self.opts.http.clone())
78 .with_metrics_handler(MetricsHandler)
79 .with_greptime_config_options(self.opts.to_toml().context(TomlFormatSnafu)?)
80 .build();
81 let addr: SocketAddr = self.opts.http.addr.parse().context(ParseAddrSnafu {
82 addr: &self.opts.http.addr,
83 })?;
84 let handler: ServerHandler = (Box::new(http_server), addr);
85 handlers.insert(handler);
86 }
87
88 Ok(handlers)
89 }
90
91 pub fn grpc_server_builder(
92 opts: &DatanodeOptions,
93 region_server: &RegionServer,
94 ) -> GrpcServerBuilder {
95 GrpcServerBuilder::new(opts.grpc.as_config(), region_server.runtime())
96 .flight_handler(Arc::new(region_server.clone()))
97 .region_server_handler(Arc::new(region_server.clone()))
98 }
99}
100
101#[cfg(test)]
102mod tests {
103 use api::v1::HealthCheckRequest;
104 use api::v1::health_check_client::HealthCheckClient;
105 use servers::grpc::GRPC_SERVER;
106
107 use super::DatanodeServiceBuilder;
108 use crate::config::DatanodeOptions;
109 use crate::tests::mock_region_server;
110
111 #[tokio::test]
112 async fn test_default_grpc_server_health_check_is_reachable() {
113 let opts = DatanodeOptions {
115 grpc: servers::grpc::GrpcOptions::default().with_bind_addr("127.0.0.1:0"),
116 ..Default::default()
117 };
118 let region_server = mock_region_server();
119 let mut services = DatanodeServiceBuilder::new(&opts)
120 .with_default_grpc_server(®ion_server)
121 .build()
122 .unwrap();
123
124 services.start_all().await.unwrap();
126 let addr = services.addr(GRPC_SERVER).unwrap();
127 let health_check = HealthCheckClient::connect(format!("http://{addr}"))
128 .await
129 .unwrap()
130 .health_check(HealthCheckRequest {})
131 .await;
132 services.shutdown_all().await.unwrap();
133
134 assert!(health_check.is_ok());
136 }
137
138 #[tokio::test]
139 async fn test_service_builder_without_grpc_server_does_not_expose_grpc_address() {
140 let opts = DatanodeOptions::default();
142 let mut services = DatanodeServiceBuilder::new(&opts).build().unwrap();
143
144 services.start_all().await.unwrap();
146 let addr = services.addr(GRPC_SERVER);
147 services.shutdown_all().await.unwrap();
148
149 assert!(addr.is_none());
151 }
152}