massa_grpc/stream/
new_transfers_info.rs1use massa_proto_rs::massa::api::v1::{self as grpc_api};
2#[cfg(feature = "execution-info")]
3use tonic::Request;
4
5use std::pin::Pin;
6
7#[cfg(feature = "execution-info")]
8use crate::{error::GrpcError, server::MassaPublicGrpc};
9
10#[cfg(feature = "execution-info")]
11use super::trait_filters_impl::{FilterGrpc, NewExecutionInfoFilter};
12
13pub type NewTransferInfoServerStreamType = Pin<
15 Box<
16 dyn futures_util::Stream<
17 Item = Result<grpc_api::NewTransfersInfoServerResponse, tonic::Status>,
18 > + Send
19 + 'static,
20 >,
21>;
22
23#[cfg(feature = "execution-info")]
24pub(crate) async fn new_transfer_info_server(
25 grpc: &MassaPublicGrpc,
26 request: Request<grpc_api::NewTransfersInfoServerRequest>,
27) -> Result<NewTransferInfoServerStreamType, GrpcError> {
28 use std::time::Duration;
29 use tokio::{select, time};
30 use tracing::error;
31
32 let (tx, rx) = tokio::sync::mpsc::channel(grpc.grpc_config.max_channel_size);
34 let request = request.into_inner();
36 let mut subscriber = grpc
38 .execution_channels
39 .slot_execution_info_sender
40 .subscribe();
41 let config = grpc.grpc_config.clone();
43
44 tokio::spawn(async move {
45 let filter = match NewExecutionInfoFilter::build_from_request(request.address, &config) {
46 Ok(filter) => filter,
47 Err(err) => {
48 error!("failed to get filter: {}", err);
49 if let Err(e) = tx.send(Err(err.into())).await {
51 error!("failed to send back NewOperations error response: {}", e);
52 }
53 return;
54 }
55 };
56
57 let mut interval = time::interval(Duration::from_secs(
60 config.unidirectional_stream_interval_check,
61 ));
62
63 loop {
65 select! {
66 event = subscriber.recv() => {
68 match event {
69 Ok(massa_operation) => {
70 if let Some(data) = filter.filter_output(massa_operation, &config) {
72 if let Err(e) = tx
74 .send(Ok(grpc_api::NewTransfersInfoServerResponse::from(data)))
75 .await
76 {
77 error!("failed to send operation : {}", e);
78 break;
79 }
80 }
81 }
82 Err(e) => error!("error on receive new operation: {}", e)
83 }
84 },
85 _ = interval.tick() => {
87 if tx.is_closed() {
88 break;
90 }
91 }
92 }
93 }
94 });
95
96 let out_stream = tokio_stream::wrappers::ReceiverStream::new(rx);
97 Ok(Box::pin(out_stream) as NewTransferInfoServerStreamType)
98}