massa_grpc/stream/
new_blocks.rs

1// Copyright (c) 2025 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, FilterNewBlocks};
15
16/// Type declaration for NewBlocks
17pub type NewBlocksStreamType = Pin<
18    Box<
19        dyn futures_util::Stream<Item = Result<grpc_api::NewBlocksResponse, tonic::Status>>
20            + Send
21            + 'static,
22    >,
23>;
24
25/// Type declaration for NewBlocksServer
26pub type NewBlocksServerStreamType = Pin<
27    Box<
28        dyn futures_util::Stream<Item = Result<grpc_api::NewBlocksServerResponse, tonic::Status>>
29            + Send
30            + 'static,
31    >,
32>;
33
34/// Creates a new stream of new produced and received blocks
35pub(crate) async fn new_blocks(
36    grpc: &MassaPublicGrpc,
37    request: Request<Streaming<grpc_api::NewBlocksRequest>>,
38) -> Result<NewBlocksStreamType, GrpcError> {
39    // Create a channel to handle communication with the client
40    let (tx, rx) = tokio::sync::mpsc::channel(grpc.grpc_config.max_channel_size);
41    // Get the inner stream from the request
42    let mut in_stream = request.into_inner();
43    // Subscribe to the new blocks channel
44    let mut subscriber = grpc.consensus_broadcasts.block_sender.subscribe();
45    // Clone grpc to be able to use it in the spawned task
46    let grpc_config = grpc.grpc_config.clone();
47
48    tokio::spawn(async move {
49        if let Some(Ok(request)) = in_stream.next().await {
50            let mut filters =
51                match FilterNewBlocks::build_from_request(request.filters, &grpc_config) {
52                    Ok(filter) => filter,
53                    Err(err) => {
54                        error!("failed to get filter: {}", err);
55                        // Send the error response back to the client
56                        if let Err(e) = tx.send(Err(err.into())).await {
57                            error!("failed to send back NewBlocks error response: {}", e);
58                        }
59                        return;
60                    }
61                };
62
63            loop {
64                select! {
65                    // Receive a new block from the subscriber
66                    event = subscriber.recv() => {
67                        match event {
68                            Ok(massa_block) => {
69                                // Check if the block should be sent
70                                if let Some(data) = filters.filter_output(massa_block, &grpc_config) {
71                                    // Send the new block through the channel
72                                    if let Err(e) = tx.send(Ok(grpc_api::NewBlocksResponse {
73                                        signed_block: Some(data.into())
74                                    })).await {
75                                        error!("failed to send new block : {}", e);
76                                        break;
77                                    }
78                                }
79
80                            },
81                            Err(e) => {
82                                // fail closed like the websocket path: notify the client
83                                // that blocks were missed so it can reconnect and resync
84                                error!("error on receive new block : {}", e);
85                                if let Err(e) = tx.send(Err(tonic::Status::data_loss(format!("error on receive new block: {}", e)))).await {
86                                    error!("failed to send back NewBlocks error response: {}", e);
87                                }
88                                break;
89                            }
90                        }
91                    },
92                    res = in_stream.next() => {
93                        match res {
94                            Some(res) => {
95                                match res {
96                                    Ok(message) => {
97                                        // Update current filter
98                                        filters = match FilterNewBlocks::build_from_request(message.filters, &grpc_config) {
99                                            Ok(filter) => filter,
100                                            Err(err) => {
101                                                error!("failed to get filter: {}", err);
102                                                // Send the error response back to the client
103                                                if let Err(e) = tx.send(Err(err.into())).await {
104                                                    error!("failed to send back NewBlocks error response: {}", e);
105                                                }
106                                                return;
107                                            }
108                                        };
109                                    },
110                                    Err(err) => {
111                                        // Check if the error matches any IO errors
112                                        if let Some(io_err) = match_for_io_error(&err) {
113                                            if io_err.kind() == ErrorKind::BrokenPipe {
114                                                warn!("client disconnected, broken pipe: {}", io_err);
115                                                break;
116                                            }
117                                        }
118                                        error!("{}", err);
119                                        // Send the error response back to the client
120                                        if let Err(e) = tx.send(Err(err)).await {
121                                            error!("failed to send back NewBlocks error response: {}", e);
122                                            break;
123                                        }
124                                    }
125                            }
126                        },
127                            None => {
128                                // The client has disconnected
129                                break;
130                            },
131                        }
132                    }
133                }
134            }
135        } else {
136            error!("empty request");
137        }
138    });
139
140    // Create a new stream from the received channel
141    let out_stream = tokio_stream::wrappers::ReceiverStream::new(rx);
142
143    // Return the new stream of blocks
144    Ok(Box::pin(out_stream) as NewBlocksStreamType)
145}
146
147/// Creates a new stream of new produced and received blocks
148/// uni-directional streaming
149pub(crate) async fn new_blocks_server(
150    grpc: &MassaPublicGrpc,
151    request: Request<grpc_api::NewBlocksServerRequest>,
152) -> Result<NewBlocksServerStreamType, GrpcError> {
153    // Create a channel to handle communication with the client
154    let (tx, rx) = tokio::sync::mpsc::channel(grpc.grpc_config.max_channel_size);
155    // Get the inner the request
156    let request = request.into_inner();
157    // Subscribe to the new blocks channel
158    let mut subscriber = grpc.consensus_broadcasts.block_sender.subscribe();
159    // Clone grpc to be able to use it in the spawned task
160    let grpc_config = grpc.grpc_config.clone();
161
162    tokio::spawn(async move {
163        // let filters = match get_filter_new_blocks(request, &grpc_config) {
164        let filter = match FilterNewBlocks::build_from_request(request.filters, &grpc_config) {
165            Ok(filter) => filter,
166            Err(err) => {
167                error!("failed to get filter: {}", err);
168                // Send the error response back to the client
169                if let Err(e) = tx.send(Err(err.into())).await {
170                    error!("failed to send back NewBlocks error response: {}", e);
171                }
172                return;
173            }
174        };
175
176        // Create a timer that ticks every 10 seconds to check if the client is still connected
177        let mut interval = time::interval(Duration::from_secs(
178            grpc_config.unidirectional_stream_interval_check,
179        ));
180
181        // Continuously loop until the stream ends or an error occurs
182        loop {
183            select! {
184                // Receive a new filled block from the subscriber
185                event = subscriber.recv() => {
186                    match event {
187                       Ok(massa_block) => {
188                           // Check if the block should be sent
189                           if let Some(data) = filter.filter_output(massa_block, &grpc_config) {
190                               // Send the new block through the channel
191                               if let Err(e) = tx
192                                   .send(Ok(grpc_api::NewBlocksServerResponse {
193                                       signed_block: Some(data.into()),
194                                   }))
195                                   .await
196                               {
197                                   error!("failed to send new block : {}", e);
198                                   break;
199                               }
200                           }
201                       },
202                       Err(e) => {
203                           // fail closed like the websocket path: notify the client
204                           // that blocks were missed so it can reconnect and resync
205                           error!("error on receive new block : {}", e);
206                           if let Err(e) = tx.send(Err(tonic::Status::data_loss(format!("error on receive new block: {}", e)))).await {
207                               error!("failed to send back NewBlocks error response: {}", e);
208                           }
209                           break;
210                       }
211                    }
212                },
213                // Execute the code block whenever the timer ticks
214                _ = interval.tick() => {
215                    if tx.is_closed() {
216                        // Client disconnected
217                        break;
218                    }
219                }
220            }
221        }
222    });
223
224    // Create a new stream from the received channel
225    let out_stream = tokio_stream::wrappers::ReceiverStream::new(rx);
226
227    // Return the new stream of blocks
228    Ok(Box::pin(out_stream) as NewBlocksServerStreamType)
229}