massa_grpc/stream/
new_filled_blocks.rs

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