massa_grpc/stream/
new_slot_execution_outputs.rs

1// Copyright (c) 2023 MASSA LABS <info@massa.net>
2
3use crate::error::{match_for_io_error, GrpcError};
4use crate::server::MassaPublicGrpc;
5use futures_util::StreamExt;
6use massa_proto_rs::massa::api::v1::{self as grpc_api};
7use std::io::ErrorKind;
8use std::pin::Pin;
9use std::time::Duration;
10use tokio::{select, time};
11use tonic::{Request, Streaming};
12use tracing::{error, warn};
13
14use super::trait_filters_impl::{FilterGrpc, FilterNewSlotExec};
15
16/// Type declaration for NewSlotExecutionOutputs
17pub type NewSlotExecutionOutputsStreamType = Pin<
18    Box<
19        dyn futures_util::Stream<
20                Item = Result<grpc_api::NewSlotExecutionOutputsResponse, tonic::Status>,
21            > + Send
22            + 'static,
23    >,
24>;
25
26/// Type declaration for NewSlotExecutionOutputsServer
27pub type NewSlotExecutionOutputsServerStreamType = Pin<
28    Box<
29        dyn futures_util::Stream<
30                Item = Result<grpc_api::NewSlotExecutionOutputsServerResponse, tonic::Status>,
31            > + Send
32            + 'static,
33    >,
34>;
35
36/// Creates a new stream of new produced and received slot execution outputs
37pub(crate) async fn new_slot_execution_outputs(
38    grpc: &MassaPublicGrpc,
39    request: Request<Streaming<grpc_api::NewSlotExecutionOutputsRequest>>,
40) -> Result<NewSlotExecutionOutputsStreamType, GrpcError> {
41    // Create a channel to handle communication with the client
42    let (tx, rx) = tokio::sync::mpsc::channel(grpc.grpc_config.max_channel_size);
43    // Get the inner stream from the request
44    let mut in_stream = request.into_inner();
45    // Subscribe to the new slot execution events channel
46    let mut subscriber = grpc
47        .execution_channels
48        .slot_execution_output_sender
49        .subscribe();
50    let grpc_config = grpc.grpc_config.clone();
51
52    tokio::spawn(async move {
53        if let Some(Ok(request)) = in_stream.next().await {
54            let mut filters: FilterNewSlotExec = match FilterNewSlotExec::build_from_request(
55                request.clone().filters,
56                &grpc_config,
57            ) {
58                Ok(filter) => filter,
59                Err(err) => {
60                    error!("failed to get filter: {}", err);
61                    // Send the error response back to the client
62                    if let Err(e) = tx.send(Err(err.into())).await {
63                        error!("failed to send back NewBlocks error response: {}", e);
64                    }
65                    return;
66                }
67            };
68
69            loop {
70                select! {
71                    // Receive a new slot execution output from the subscriber
72                    event = subscriber.recv() => {
73                        match event {
74                            Ok(massa_slot_execution_output) => {
75                                if let Some(data) = filters.filter_output(massa_slot_execution_output, &grpc_config) {
76                                    if let Err(e) = tx.send(Ok(grpc_api::NewSlotExecutionOutputsResponse {
77                                        output: Some(data.into()),
78                                    })).await {
79                                        error!("failed to send new slot execution output : {}", e);
80                                        break;
81                                    }
82                                }
83                            },
84
85                            Err(e) => error!("error on receive new slot execution output : {}", e)
86                        }
87                    },
88                    // Receive a new message from the in_stream
89                    res = in_stream.next() => {
90                        match res {
91                            Some(res) => {
92                                match res {
93                                    Ok(message) => {
94                                        // Update current filter
95                                        filters = match FilterNewSlotExec::build_from_request(message.clone().filters, &grpc_config) {
96                                            Ok(filter) => filter,
97                                            Err(err) => {
98                                                error!("failed to get filter: {}", err);
99                                                // Send the error response back to the client
100                                                if let Err(e) = tx.send(Err(err.into())).await {
101                                                    error!("failed to send back NewBlocks error response: {}", e);
102                                                }
103                                                return;
104                                            }
105                                        };
106                                    },
107                                    // Handle any errors that may occur during receiving the data
108                                    Err(err) => {
109                                        // Check if the error matches any IO errors
110                                        if let Some(io_err) = match_for_io_error(&err) {
111                                            if io_err.kind() == ErrorKind::BrokenPipe {
112                                                warn!("client disconnected, broken pipe: {}", io_err);
113                                                break;
114                                            }
115                                        }
116                                        error!("{}", err);
117                                        // Send the error response back to the client
118                                        if let Err(e) = tx.send(Err(err)).await {
119                                            error!("failed to send back new_slot_execution_outputs error response: {}", e);
120                                            break;
121                                        }
122                                    }
123                                }
124                            },
125                            None => {
126                                // The client has disconnected
127                                break;
128                            },
129                        }
130                    }
131                }
132            }
133        } else {
134            error!("empty request");
135        }
136    });
137
138    // Create a new stream from the received channel
139    let out_stream = tokio_stream::wrappers::ReceiverStream::new(rx);
140
141    // Return the new stream of slot execution output
142    Ok(Box::pin(out_stream) as NewSlotExecutionOutputsStreamType)
143}
144
145pub(crate) async fn new_slot_execution_outputs_server(
146    grpc: &MassaPublicGrpc,
147    request: tonic::Request<grpc_api::NewSlotExecutionOutputsServerRequest>,
148) -> Result<NewSlotExecutionOutputsServerStreamType, GrpcError> {
149    // Create a channel to handle communication with the client
150    let (tx, rx) = tokio::sync::mpsc::channel(grpc.grpc_config.max_channel_size);
151    // Subscribe to the new slot execution events channel
152    let mut subscriber = grpc
153        .execution_channels
154        .slot_execution_output_sender
155        .subscribe();
156    let grpc = grpc.clone();
157    let inner_req = request.into_inner();
158    tokio::spawn(async move {
159        let filters: FilterNewSlotExec = match FilterNewSlotExec::build_from_request(
160            inner_req.clone().filters,
161            &grpc.grpc_config,
162        ) {
163            Ok(filter) => filter,
164            Err(err) => {
165                error!("failed to get filter: {}", err);
166                // Send the error response back to the client
167                if let Err(e) = tx.send(Err(err.into())).await {
168                    error!("failed to send back error response: {}", e);
169                }
170                return;
171            }
172        };
173
174        // Create a timer that ticks every 10 seconds to check if the client is still connected
175        let mut interval = time::interval(Duration::from_secs(
176            grpc.grpc_config.unidirectional_stream_interval_check,
177        ));
178
179        // Continuously loop until the stream ends or an error occurs
180        loop {
181            select! {
182                // Receive a new filled block from the subscriber
183                event = subscriber.recv() => {
184                    match event {
185                        Ok(massa_slot_execution_output) => {
186                            // Check if the slot execution output should be sent
187                            if let Some(slot_execution_output) =
188                                filters.filter_output(massa_slot_execution_output, &grpc.grpc_config)
189                            {
190                                // Send the new slot execution output through the channel
191                                if let Err(e) = tx
192                                    .send(Ok(grpc_api::NewSlotExecutionOutputsServerResponse {
193                                        output: Some(slot_execution_output.into()),
194                                    }))
195                                    .await
196                                {
197                                    error!("failed to send new slot execution output : {}", e);
198                                    break;
199                                }
200                            }
201                        }
202                        Err(e) => error!("error on receive new slot execution output : {}", e),
203                    }
204                },
205                // Execute the code block whenever the timer ticks
206                _ = interval.tick() => {
207                    if tx.is_closed() {
208                        // Client disconnected
209                        break;
210                    }
211                }
212            }
213        }
214    });
215    // Create a new stream from the received channel
216    let out_stream = tokio_stream::wrappers::ReceiverStream::new(rx);
217    // Return the new stream of slot execution output
218    Ok(Box::pin(out_stream) as NewSlotExecutionOutputsServerStreamType)
219}