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}