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}