massa_grpc/stream/new_endorsements.rs
1// Copyright (c) 2023 MASSA LABS <info@massa.net>
2
3use crate::error::{match_for_io_error, GrpcError};
4use crate::server::MassaPublicGrpc;
5use futures_util::StreamExt;
6use massa_proto_rs::massa::api::v1::{self as grpc_api};
7use std::io::ErrorKind;
8use std::pin::Pin;
9use std::time::Duration;
10use tokio::{select, time};
11use tonic::{Request, Streaming};
12use tracing::{error, warn};
13
14use super::trait_filters_impl::{FilterGrpc, NewEndorsementsFilter};
15
16/// Type declaration for NewEndorsements
17pub type NewEndorsementsStreamType = Pin<
18 Box<
19 dyn futures_util::Stream<Item = Result<grpc_api::NewEndorsementsResponse, tonic::Status>>
20 + Send
21 + 'static,
22 >,
23>;
24
25/// Type declaration for NewEndorsementsServer
26pub type NewEndorsementsServerStreamType = Pin<
27 Box<
28 dyn futures_util::Stream<
29 Item = Result<grpc_api::NewEndorsementsServerResponse, tonic::Status>,
30 > + Send
31 + 'static,
32 >,
33>;
34
35/// Creates a new stream of new produced and received endorsements
36pub(crate) async fn new_endorsements(
37 grpc: &MassaPublicGrpc,
38 request: Request<Streaming<grpc_api::NewEndorsementsRequest>>,
39) -> Result<NewEndorsementsStreamType, GrpcError> {
40 // Create a channel to handle communication with the client
41 let (tx, rx) = tokio::sync::mpsc::channel(grpc.grpc_config.max_channel_size);
42 // Get the inner stream from the request
43 let mut in_stream = request.into_inner();
44 // Subscribe to the new endorsements channel
45 let mut subscriber = grpc.pool_broadcasts.endorsement_sender.subscribe();
46 // Clone grpc to be able to use it in the spawned task
47 let grpc_config = grpc.grpc_config.clone();
48
49 tokio::spawn(async move {
50 if let Some(Ok(request)) = in_stream.next().await {
51 let mut filter =
52 match NewEndorsementsFilter::build_from_request(request.filters, &grpc_config) {
53 Ok(filter) => filter,
54 Err(err) => {
55 error!("failed to get filter: {}", err);
56 // Send the error response back to the client
57 if let Err(e) = tx.send(Err(err.into())).await {
58 error!("failed to send back NewEndorsements error response: {}", e);
59 }
60 return;
61 }
62 };
63
64 loop {
65 select! {
66 // Receive a new endorsement from the subscriber
67 event = subscriber.recv() => {
68 match event {
69 Ok(massa_endorsement) => {
70 // Check if the endorsement should be sent
71 if let Some(data) = filter.filter_output(massa_endorsement, &grpc_config) {
72 // Send the new endorsement through the channel
73 if let Err(e) = tx.send(Ok(grpc_api::NewEndorsementsResponse {
74 signed_endorsement: Some(data.into())
75 })).await {
76 error!("failed to send new endorsement : {}", e);
77 break;
78 }
79 }
80 },
81 Err(e) => error!("error on receive new endorsement : {}", e)
82 }
83 },
84 // Receive a new message from the in_stream
85 res = in_stream.next() => {
86 match res {
87 Some(res) => {
88 match res {
89 Ok(message) => {
90 // Update current filter
91 filter = match NewEndorsementsFilter::build_from_request(message.filters, &grpc_config) {
92 Ok(filter) => filter,
93 Err(err) => {
94 error!("failed to get filter: {}", err);
95 // Send the error response back to the client
96 if let Err(e) = tx.send(Err(err.into())).await {
97 error!("failed to send back NewEndorsements error response: {}", e);
98 }
99 return;
100 }
101 };
102 },
103 Err(err) => {
104 // Check if the error matches any IO errors
105 if let Some(io_err) = match_for_io_error(&err) {
106 if io_err.kind() == ErrorKind::BrokenPipe {
107 warn!("client disconnected, broken pipe: {}", io_err);
108 break;
109 }
110 }
111 error!("{}", err);
112 // Send the error response back to the client
113 if let Err(e) = tx.send(Err(err)).await {
114 error!("failed to send back NewEndorsements error response: {}", e);
115 break;
116 }
117 }
118 }
119 },
120 None => {
121 // The client has disconnected
122 break;
123 },
124 }
125 }
126 }
127 }
128 } else {
129 error!("empty request");
130 }
131 });
132
133 // Create a new stream from the received channel
134 let out_stream = tokio_stream::wrappers::ReceiverStream::new(rx);
135
136 // Return the new stream of endorsements
137 Ok(Box::pin(out_stream) as NewEndorsementsStreamType)
138}
139
140/// Creates a new unidirectional stream of new produced and received endorsements
141pub(crate) async fn new_endorsements_server(
142 grpc: &MassaPublicGrpc,
143 request: Request<grpc_api::NewEndorsementsServerRequest>,
144) -> Result<NewEndorsementsServerStreamType, GrpcError> {
145 // Create a channel to handle communication with the client
146 let (tx, rx) = tokio::sync::mpsc::channel(grpc.grpc_config.max_channel_size);
147 // Get the inner request
148 let request = request.into_inner();
149 // Subscribe to the new endorsements channel
150 let mut subscriber = grpc.pool_broadcasts.endorsement_sender.subscribe();
151 // Clone grpc to be able to use it in the spawned task
152 let grpc_config = grpc.grpc_config.clone();
153
154 tokio::spawn(async move {
155 let filter = match NewEndorsementsFilter::build_from_request(request.filters, &grpc_config)
156 {
157 Ok(filter) => filter,
158 Err(err) => {
159 error!("failed to get filter: {}", err);
160 // Send the error response back to the client
161 if let Err(e) = tx.send(Err(err.into())).await {
162 error!("failed to send back new endorsement error response: {}", e);
163 }
164 return;
165 }
166 };
167
168 // Create a timer that ticks every 10 seconds to check if the client is still connected
169 let mut interval = time::interval(Duration::from_secs(
170 grpc_config.unidirectional_stream_interval_check,
171 ));
172
173 loop {
174 select! {
175 // Receive a new endorsement from the subscriber
176 event = subscriber.recv() => {
177 match event {
178 Ok(massa_endorsement) => {
179 // Check if the endorsement should be sent
180 if let Some(data) = filter.filter_output(massa_endorsement, &grpc_config) {
181 // Send the new endorsement through the channel
182 if let Err(e) = tx.send(Ok(grpc_api::NewEndorsementsServerResponse {
183 signed_endorsement: Some(data.into())
184 })).await {
185 error!("failed to send new endorsement : {}", e);
186 break;
187 }
188 }
189 },
190 Err(e) => error!("error on receive new endorsement : {}", e)
191 }
192 },
193 _ = interval.tick() => {
194 if tx.is_closed() {
195 // Client disconnected
196 break;
197 }
198 }
199 }
200 }
201 });
202
203 // Create a new stream from the received channel
204 let out_stream = tokio_stream::wrappers::ReceiverStream::new(rx);
205
206 // Return the new stream of endorsements
207 Ok(Box::pin(out_stream) as NewEndorsementsServerStreamType)
208}