Skip to main content

datanode/
service.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 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        // Arrange
114        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(&region_server)
121            .build()
122            .unwrap();
123
124        // Act
125        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
135        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        // Arrange
141        let opts = DatanodeOptions::default();
142        let mut services = DatanodeServiceBuilder::new(&opts).build().unwrap();
143
144        // Act
145        services.start_all().await.unwrap();
146        let addr = services.addr(GRPC_SERVER);
147        services.shutdown_all().await.unwrap();
148
149        // Assert
150        assert!(addr.is_none());
151    }
152}