1use massa_bootstrap::white_black_list::SharedWhiteBlackList;
4use massa_models::node::NodeId;
5use massa_versioning::keypair_factory::KeyPairFactory;
6use massa_versioning::versioning::MipStore;
7use parking_lot::RwLock;
8use std::convert::Infallible;
9use std::path::Path;
10use std::sync::{Arc, Condvar, Mutex};
11
12use crate::config::{GrpcConfig, ServiceName};
13use crate::error::GrpcError;
14use futures_util::FutureExt;
15use hyper::{Method, Request, Response};
16use massa_consensus_exports::{ConsensusBroadcasts, ConsensusController};
17use massa_execution_exports::{ExecutionChannels, ExecutionController};
18use massa_pool_exports::{PoolBroadcasts, PoolController};
19use massa_pos_exports::SelectorController;
20use massa_proto_rs::massa::api::v1::{
21 private_service_server::PrivateServiceServer, public_service_server::PublicServiceServer,
22};
23use massa_proto_rs::massa::api::v1::{FILE_DESCRIPTOR_SET_PRIVATE, FILE_DESCRIPTOR_SET_PUBLIC};
24use massa_protocol_exports::{ProtocolConfig, ProtocolController};
25use massa_sdk::cert_manager::{gen_cert_for_ca, gen_signed_cert};
26use massa_storage::Storage;
27
28use massa_wallet::Wallet;
29
30use tokio::sync::oneshot;
31use tonic::body::BoxBody;
32use tonic::codegen::CompressionEncoding;
33use tonic::server::NamedService;
34use tonic::transport::{Certificate, Identity, ServerTlsConfig};
35use tonic_health::server::HealthReporter;
36use tonic_web::GrpcWebLayer;
37use tower_http::cors::{Any, CorsLayer};
38use tracing::{info, warn};
39
40#[derive(Clone)]
42pub struct MassaPrivateGrpc {
43 pub consensus_controller: Box<dyn ConsensusController>,
45 pub execution_controller: Box<dyn ExecutionController>,
47 pub pool_controller: Box<dyn PoolController>,
49 pub protocol_controller: Box<dyn ProtocolController>,
51 pub stop_cv: Arc<(Mutex<bool>, Condvar)>,
54 pub node_wallet: Arc<RwLock<Wallet>>,
56 pub grpc_config: GrpcConfig,
58 pub protocol_config: ProtocolConfig,
60 pub node_id: NodeId,
62 pub mip_store: MipStore,
64 pub version: massa_models::version::Version,
66 pub bs_white_black_list: Option<SharedWhiteBlackList<'static>>,
68}
69
70impl MassaPrivateGrpc {
71 pub async fn serve(self, config: &GrpcConfig) -> Result<StopHandle, GrpcError> {
73 let mut service = PrivateServiceServer::new(self)
74 .max_decoding_message_size(config.max_decoding_message_size)
75 .max_encoding_message_size(config.max_encoding_message_size);
76
77 if let Some(encoding) = &config.accept_compressed {
78 if encoding.eq_ignore_ascii_case("Gzip") {
79 service = service.accept_compressed(CompressionEncoding::Gzip);
80 };
81 }
82
83 if let Some(encoding) = &config.send_compressed {
84 if encoding.eq_ignore_ascii_case("Gzip") {
85 service = service.send_compressed(CompressionEncoding::Gzip);
86 };
87 }
88
89 serve(service, config).await
90 }
91}
92
93#[derive(Clone)]
95pub struct MassaPublicGrpc {
96 pub consensus_controller: Box<dyn ConsensusController>,
98 pub consensus_broadcasts: ConsensusBroadcasts,
100 pub execution_controller: Box<dyn ExecutionController>,
102 pub execution_channels: ExecutionChannels,
104 pub pool_broadcasts: PoolBroadcasts,
106 pub pool_controller: Box<dyn PoolController>,
108 pub protocol_controller: Box<dyn ProtocolController>,
110 pub selector_controller: Box<dyn SelectorController>,
112 pub storage: Storage,
114 pub grpc_config: GrpcConfig,
116 pub protocol_config: ProtocolConfig,
118 pub node_id: NodeId,
120 pub version: massa_models::version::Version,
122 pub keypair_factory: KeyPairFactory,
124}
125
126impl MassaPublicGrpc {
127 pub async fn serve(self, config: &GrpcConfig) -> Result<StopHandle, GrpcError> {
129 let mut service = PublicServiceServer::new(self)
130 .max_decoding_message_size(config.max_decoding_message_size)
131 .max_encoding_message_size(config.max_encoding_message_size);
132
133 if let Some(encoding) = &config.accept_compressed {
134 if encoding.eq_ignore_ascii_case("Gzip") {
135 service = service.accept_compressed(CompressionEncoding::Gzip);
136 };
137 }
138
139 if let Some(encoding) = &config.send_compressed {
140 if encoding.eq_ignore_ascii_case("Gzip") {
141 service = service.send_compressed(CompressionEncoding::Gzip);
142 };
143 }
144 serve(service, config).await
145 }
146}
147
148pub struct StopHandle {
150 stop_cmd_sender: oneshot::Sender<()>,
151}
152
153impl StopHandle {
154 pub fn stop(self) {
156 if let Err(e) = self.stop_cmd_sender.send(()) {
157 warn!("gRPC API thread panicked: {:?}", e);
158 } else {
159 info!("gRPC API stop signal sent successfully");
160 }
161 }
162}
163
164async fn massa_service_status(mut reporter: HealthReporter) {
166 reporter
168 .set_serving::<PublicServiceServer<MassaPublicGrpc>>()
169 .await;
170}
171
172pub(crate) fn check_mtls_requires_tls(config: &GrpcConfig) -> Result<(), GrpcError> {
178 massa_sdk::check_mtls_requires_tls(config.enable_tls, config.enable_mtls)
179 .map_err(|err| GrpcError::ConfigError(format!("gRPC {:?} API: {}", config.name, err)))
180}
181
182async fn serve<S>(service: S, config: &GrpcConfig) -> Result<StopHandle, GrpcError>
184where
185 S: tower_service::Service<Request<BoxBody>, Response = Response<BoxBody>, Error = Infallible>
186 + NamedService
187 + Clone
188 + Send
189 + 'static,
190 S::Future: Send + 'static,
191{
192 check_mtls_requires_tls(config)?;
193
194 let (shutdown_send, shutdown_recv) = oneshot::channel::<()>();
195
196 let mut server_builder = tonic::transport::Server::builder()
197 .concurrency_limit_per_connection(config.concurrency_limit_per_connection)
198 .timeout(config.timeout)
199 .initial_stream_window_size(config.initial_stream_window_size)
200 .initial_connection_window_size(config.initial_connection_window_size)
201 .max_concurrent_streams(config.max_concurrent_streams)
202 .tcp_keepalive(config.tcp_keepalive)
203 .tcp_nodelay(config.tcp_nodelay)
204 .http2_keepalive_interval(config.http2_keepalive_interval)
205 .http2_keepalive_timeout(config.http2_keepalive_timeout)
206 .http2_adaptive_window(config.http2_adaptive_window)
207 .max_frame_size(config.max_frame_size);
208
209 if config.enable_tls {
210 if config.generate_self_signed_certificates {
211 if Path::new(&config.certificate_authority_root_path).exists() {
212 warn!("Certificate authority root already exists, remove the file if you want to generate new certificates. Skipping self signed certificates generation.");
213 } else {
214 info!("Generating self signed certificates");
215 generate_self_signed_certificates(config);
216 }
217 }
218
219 let cert = std::fs::read_to_string(config.server_certificate_path.clone())
220 .expect("error, failed to read server certificate");
221 let key = std::fs::read_to_string(config.server_private_key_path.clone())
222 .expect("error, failed to read server private key");
223
224 let server_identity = Identity::from_pem(cert, key);
225 let tls = ServerTlsConfig::new().identity(server_identity);
226
227 if config.enable_mtls {
228 let client_ca_cert =
229 std::fs::read_to_string(config.client_certificate_authority_root_path.clone())
230 .expect("error, failed to read client certificate authority root");
231 let client_ca_cert = Certificate::from_pem(client_ca_cert);
232
233 let tls = tls.client_ca_root(client_ca_cert);
234
235 server_builder = server_builder
236 .tls_config(tls)
237 .expect("error, failed to setup mTLS");
238
239 info!("gRPC mTLS enabled");
240 } else {
241 server_builder = server_builder
242 .tls_config(tls)
243 .expect("error, failed to setup TLS");
244 info!("gRPC TLS enabled");
245 }
246 }
247
248 let reflection_service_opt_alpha = if config.enable_reflection {
249 let file_descriptor_set = match config.name {
250 ServiceName::Public => FILE_DESCRIPTOR_SET_PUBLIC,
251 ServiceName::Private => FILE_DESCRIPTOR_SET_PRIVATE,
252 };
253 let reflection_service = tonic_reflection::server::Builder::configure()
254 .register_encoded_file_descriptor_set(file_descriptor_set)
255 .build_v1alpha()?;
256
257 Some(reflection_service)
258 } else {
259 None
260 };
261
262 let reflection_service_opt = if config.enable_reflection {
263 let file_descriptor_set = match config.name {
264 ServiceName::Public => FILE_DESCRIPTOR_SET_PUBLIC,
265 ServiceName::Private => FILE_DESCRIPTOR_SET_PRIVATE,
266 };
267 let reflection_service = tonic_reflection::server::Builder::configure()
268 .register_encoded_file_descriptor_set(file_descriptor_set)
269 .build_v1()?;
270
271 Some(reflection_service)
272 } else {
273 None
274 };
275
276 let health_service_opt = if config.enable_health {
277 let (mut health_reporter, health_service) = tonic_health::server::health_reporter();
278 health_reporter
279 .set_serving::<PublicServiceServer<MassaPublicGrpc>>()
280 .await;
281 tokio::spawn(massa_service_status(health_reporter.clone()));
282 info!("gRPC health service enabled");
283 Some(health_service)
284 } else {
285 None
286 };
287
288 if config.accept_http1 {
289 if config.enable_cors {
290 let cors = CorsLayer::new()
291 .allow_methods([Method::GET, Method::POST, Method::OPTIONS])
293 .allow_origin(Any)
295 .allow_headers(Any);
296
297 let router_with_http1 = server_builder
298 .accept_http1(true)
299 .layer(cors)
300 .layer(GrpcWebLayer::new())
301 .add_optional_service(reflection_service_opt)
302 .add_optional_service(reflection_service_opt_alpha)
303 .add_optional_service(health_service_opt)
304 .add_service(service);
305
306 tokio::spawn(
307 router_with_http1.serve_with_shutdown(config.bind, shutdown_recv.map(drop)),
308 );
309 } else {
310 let router_with_http1 = server_builder
311 .accept_http1(true)
312 .layer(GrpcWebLayer::new())
313 .add_optional_service(reflection_service_opt)
314 .add_optional_service(reflection_service_opt_alpha)
315 .add_optional_service(health_service_opt)
316 .add_service(service);
317
318 tokio::spawn(
319 router_with_http1.serve_with_shutdown(config.bind, shutdown_recv.map(drop)),
320 );
321 }
322 } else {
323 let router = server_builder
324 .add_optional_service(reflection_service_opt)
325 .add_optional_service(reflection_service_opt_alpha)
326 .add_optional_service(health_service_opt)
327 .add_service(service);
328
329 tokio::spawn(router.serve_with_shutdown(config.bind, shutdown_recv.map(drop)));
330 }
331
332 Ok(StopHandle {
333 stop_cmd_sender: shutdown_send,
334 })
335}
336
337fn generate_self_signed_certificates(config: &GrpcConfig) {
339 let ca_cert = gen_cert_for_ca().expect("error, failed to generate CA cert");
340 let ca_cert_pem = ca_cert
341 .serialize_pem()
342 .expect("error: failed to convert certificate authority to UTF-8");
343
344 if config.enable_mtls {
345 std::fs::write(
346 config.client_certificate_authority_root_path.clone(),
347 ca_cert_pem.clone(),
348 )
349 .expect("error, failed to write client certificate authority root");
350
351 let (client_cert_pem, client_private_key_pem) =
352 gen_signed_cert(&ca_cert, config.subject_alt_names.clone())
353 .expect("error, failed to generate cert");
354 std::fs::write(config.client_certificate_path.clone(), client_cert_pem)
355 .expect("error, failed to write client certificate");
356 std::fs::write(
357 config.client_private_key_path.clone(),
358 client_private_key_pem,
359 )
360 .expect("error, failed to write client private key");
361 }
362
363 std::fs::write(config.certificate_authority_root_path.clone(), ca_cert_pem)
364 .expect("error, failed to write certificate authority root");
365
366 let (cert_pem, server_private_key_pem) =
367 gen_signed_cert(&ca_cert, config.subject_alt_names.clone())
368 .expect("error, failed to generate server certificate");
369 std::fs::write(config.server_certificate_path.clone(), cert_pem)
370 .expect("error, failed to write server certificate");
371 std::fs::write(
372 config.server_private_key_path.clone(),
373 server_private_key_pem,
374 )
375 .expect("error, failed to write server private key");
376}