massa_grpc/stream/
new_transfers_info.rs

1use 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
13/// Type declaration for New execution Info server
14pub 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    // Create a channel to handle communication with the client
33    let (tx, rx) = tokio::sync::mpsc::channel(grpc.grpc_config.max_channel_size);
34    // Get the inner request
35    let request = request.into_inner();
36    // Subscribe to the new operations channel
37    let mut subscriber = grpc
38        .execution_channels
39        .slot_execution_info_sender
40        .subscribe();
41    // Clone grpc to be able to use it in the spawned task
42    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                // Send the error response back to the client
50                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        // Create a timer that ticks every 10 seconds to check if the client is still connected
58        // otherwise the server has no way to check if client has disconnected (and can help to save some resources)
59        let mut interval = time::interval(Duration::from_secs(
60            config.unidirectional_stream_interval_check,
61        ));
62
63        // Continuously loop until the stream ends or an error occurs
64        loop {
65            select! {
66                // Receive a new filled block from the subscriber
67                event = subscriber.recv() => {
68                    match event {
69                        Ok(massa_operation) => {
70                            // Check if the operation should be sent
71                            if let Some(data) = filter.filter_output(massa_operation, &config) {
72                                // Send the new operation through the channel
73                                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                // Execute the code block whenever the timer ticks
86                _ = interval.tick() => {
87                    if tx.is_closed() {
88                        // Client disconnected
89                        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}