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}