massa_grpc/stream/
new_operations.rs

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