massa_grpc/stream/
new_filled_blocks.rs1use 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
17pub type NewFilledBlocksStreamType = Pin<
19 Box<
20 dyn futures_util::Stream<Item = Result<grpc_api::NewFilledBlocksResponse, tonic::Status>>
21 + Send
22 + 'static,
23 >,
24>;
25
26pub 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
36pub(crate) async fn new_filled_blocks(
38 grpc: &MassaPublicGrpc,
39 request: Request<Streaming<grpc_api::NewFilledBlocksRequest>>,
40) -> Result<NewFilledBlocksStreamType, GrpcError> {
41 let (tx, rx) = tokio::sync::mpsc::channel(grpc.grpc_config.max_channel_size);
43 let mut in_stream = request.into_inner();
45 let mut subscriber = grpc.consensus_broadcasts.filled_block_sender.subscribe();
47 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 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 event = subscriber.recv() => {
69 match event {
70 Ok(massa_filled_block) => {
71 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 res = in_stream.next() => {
87 match res {
88 Some(res) => {
89 match res {
90 Ok(message) => {
91 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 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 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 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 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
138pub(crate) async fn new_filled_blocks_server(
140 grpc: &MassaPublicGrpc,
141 request: Request<grpc_api::NewFilledBlocksServerRequest>,
142) -> Result<NewFilledBlocksServerStreamType, GrpcError> {
143 let (tx, rx) = tokio::sync::mpsc::channel(grpc.grpc_config.max_channel_size);
145 let request = request.into_inner();
147 let mut subscriber = grpc.consensus_broadcasts.filled_block_sender.subscribe();
149 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 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 let mut interval = time::interval(Duration::from_secs(
168 grpc_config.unidirectional_stream_interval_check,
169 ));
170
171 loop {
173 select! {
174 event = subscriber.recv() => {
176 match event {
177 Ok(massa_filled_block) => {
178 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 _ = interval.tick() => {
196 if tx.is_closed() {
197 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}