massa_grpc/
server.rs

1// Copyright (c) 2023 MASSA LABS <info@massa.net>
2
3use 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/// gRPC PRIVATE API content
41#[derive(Clone)]
42pub struct MassaPrivateGrpc {
43    /// link to the consensus component
44    pub consensus_controller: Box<dyn ConsensusController>,
45    /// link to the execution component
46    pub execution_controller: Box<dyn ExecutionController>,
47    /// link to the pool component
48    pub pool_controller: Box<dyn PoolController>,
49    /// link to the protocol component
50    pub protocol_controller: Box<dyn ProtocolController>,
51    /// Mechanism by which to gracefully shut down.
52    /// To be a clone of the same pair provided to the ctrlc handler.
53    pub stop_cv: Arc<(Mutex<bool>, Condvar)>,
54    /// User wallet
55    pub node_wallet: Arc<RwLock<Wallet>>,
56    /// gRPC configuration
57    pub grpc_config: GrpcConfig,
58    /// Massa protocol configuration
59    pub protocol_config: ProtocolConfig,
60    /// our node id
61    pub node_id: NodeId,
62    /// database for all MIP info
63    pub mip_store: MipStore,
64    /// node version
65    pub version: massa_models::version::Version,
66    /// white/black list of bootstrap
67    pub bs_white_black_list: Option<SharedWhiteBlackList<'static>>,
68}
69
70impl MassaPrivateGrpc {
71    /// Start the gRPC PRIVATE API
72    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/// gRPC PUBLIC API content
94#[derive(Clone)]
95pub struct MassaPublicGrpc {
96    /// link to the consensus component
97    pub consensus_controller: Box<dyn ConsensusController>,
98    /// Broadcasts made by consensus component
99    pub consensus_broadcasts: ConsensusBroadcasts,
100    /// link to the execution component
101    pub execution_controller: Box<dyn ExecutionController>,
102    /// link(channels) to the execution component
103    pub execution_channels: ExecutionChannels,
104    /// Broadcasts made by pool component
105    pub pool_broadcasts: PoolBroadcasts,
106    /// link to the pool component
107    pub pool_controller: Box<dyn PoolController>,
108    /// link to the protocol component
109    pub protocol_controller: Box<dyn ProtocolController>,
110    /// link to the selector component
111    pub selector_controller: Box<dyn SelectorController>,
112    /// link to the storage component
113    pub storage: Storage,
114    /// gRPC configuration
115    pub grpc_config: GrpcConfig,
116    /// Massa protocol configuration
117    pub protocol_config: ProtocolConfig,
118    /// our node id
119    pub node_id: NodeId,
120    /// node version
121    pub version: massa_models::version::Version,
122    /// keypair factory
123    pub keypair_factory: KeyPairFactory,
124}
125
126impl MassaPublicGrpc {
127    /// Start the gRPC PUBLIC API
128    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
148/// Used to be able to stop the gRPC API
149pub struct StopHandle {
150    stop_cmd_sender: oneshot::Sender<()>,
151}
152
153impl StopHandle {
154    /// stop the gRPC API gracefully
155    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
164/// Massa service health check implementation
165async fn massa_service_status(mut reporter: HealthReporter) {
166    //TODO add a complete health check based on Massa modules health
167    reporter
168        .set_serving::<PublicServiceServer<MassaPublicGrpc>>()
169        .await;
170}
171
172/// mTLS is a strict extension of TLS: reject the combination rather than silently serving in
173/// plaintext a service the operator meant to protect with client certificates.
174///
175/// The invariant itself is shared with the gRPC client (massa-client), see
176/// `massa_sdk::check_mtls_requires_tls`.
177pub(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
182// Configure and start the gRPC API with the given service
183async 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 `GET`, `POST` and `OPTIONS` when accessing the resource
292                .allow_methods([Method::GET, Method::POST, Method::OPTIONS])
293                // Allow requests from any origin
294                .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
337// Generate self signed certificates
338fn 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}