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}