massa_grpc/stream/send_operations.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_models::operation::{OperationDeserializer, SecureShareOperation};
7use massa_models::secure_share::SecureShareDeserializer;
8use massa_models::timeslots::get_latest_block_slot_at_timestamp;
9use massa_proto_rs::massa::api::v1 as grpc_api;
10use massa_proto_rs::massa::model::v1 as grpc_model;
11use massa_serialization::{DeserializeError, Deserializer};
12use massa_time::MassaTime;
13use std::collections::HashMap;
14use std::io::ErrorKind;
15use std::pin::Pin;
16use tracing::{error, warn};
17
18/// Type declaration for SendOperations
19pub type SendOperationsStreamType = Pin<
20 Box<
21 dyn futures_util::Stream<Item = Result<grpc_api::SendOperationsResponse, tonic::Status>>
22 + Send
23 + 'static,
24 >,
25>;
26
27/// This function takes a streaming request of operations messages,
28/// verifies, saves and propagates the operations received in each message, and sends back a stream of
29/// operations ids messages
30pub(crate) async fn send_operations(
31 grpc: &MassaPublicGrpc,
32 request: tonic::Request<tonic::Streaming<grpc_api::SendOperationsRequest>>,
33) -> Result<SendOperationsStreamType, GrpcError> {
34 let mut pool_controller = grpc.pool_controller.clone();
35 let protocol_controller = grpc.protocol_controller.clone();
36 let config = grpc.grpc_config.clone();
37 let storage = grpc.storage.clone_without_refs();
38
39 // Create a channel for sending responses to the client
40 let (tx, rx) = tokio::sync::mpsc::channel(config.max_channel_size);
41 // Extract the incoming stream of operations messages
42 let mut in_stream = request.into_inner();
43
44 // Spawn a task that reads incoming messages and processes the operations in each message
45 tokio::spawn(async move {
46 while let Some(result) = in_stream.next().await {
47 match result {
48 Ok(req_content) => {
49 // If the incoming message has no operations, send an error message back to the client
50 if req_content.operations.is_empty() {
51 report_error(
52 tx.clone(),
53 tonic::Code::InvalidArgument,
54 "the request payload is empty".to_owned(),
55 )
56 .await;
57 } else {
58 let now = MassaTime::now();
59 let Ok(last_slot) = get_latest_block_slot_at_timestamp(
60 config.thread_count,
61 config.t0,
62 config.genesis_timestamp,
63 now,
64 ) else {
65 report_error(
66 tx.clone(),
67 tonic::Code::InvalidArgument,
68 "failed to get current period".to_owned(),
69 )
70 .await;
71 continue;
72 };
73 // If there are too many operations in the incoming message, send an error message back to the client
74 if req_content.operations.len() as u32 > config.max_operations_per_message {
75 report_error(
76 tx.clone(),
77 tonic::Code::InvalidArgument,
78 "too many operations per message".to_owned(),
79 )
80 .await;
81 } else {
82 // Deserialize and verify each operation in the incoming message
83 let operation_deserializer = SecureShareDeserializer::new(
84 OperationDeserializer::new(
85 config.max_bytecode_size,
86 config.max_function_name_length,
87 config.max_parameter_size,
88 config.max_op_datastore_entry_count,
89 config.max_op_datastore_key_length,
90 config.max_op_datastore_value_length,
91 ),
92 config.chain_id,
93 );
94 let verified_ops_res: Result<HashMap<String, SecureShareOperation>, GrpcError> = req_content.operations
95 .into_iter()
96 .map(|proto_operation| {
97 // Deserialize the operation and verify its signature
98 let verified_op_res = match operation_deserializer.deserialize::<DeserializeError>(&proto_operation) {
99 Ok(tuple) => {
100 let (rest, res_operation): (&[u8], SecureShareOperation) = tuple;
101 // an operation bigger than the whole operations payload of a
102 // block can never be propagated nor included in a block
103 let op_size = res_operation.serialized_size();
104 if op_size > config.max_serialized_operation_size {
105 return Err(GrpcError::InvalidArgument(format!(
106 "serialized operation size {} exceeds the maximum authorized size of {} bytes. Your operation will never be included in a block.",
107 op_size, config.max_serialized_operation_size
108 )));
109 }
110 if let Err(err_msg) = res_operation.check_gas_usage(
111 config.max_gas_per_block,
112 config.base_operation_gas_cost,
113 config.sp_compilation_cost,
114 ) {
115 return Err(GrpcError::InvalidArgument(err_msg));
116 }
117 if let Some(slot) = last_slot {
118 if res_operation.content.expire_period < slot.period {
119 return Err(GrpcError::InvalidArgument("Operation expire_period is lower than the current period of this node. Your operation will never be included in a block.".into()));
120 }
121 }
122
123
124 if res_operation.content.fee.checked_sub(config.minimal_fees).is_none() {
125 return Err(GrpcError::InvalidArgument("Operation fee is lower than the minimal fee. Your operation will never be included in a block.".into()));
126 }
127
128 if rest.is_empty() {
129 res_operation.verify_signature(None)
130 .map(|_| (res_operation.id.to_string(), res_operation))
131 .map_err(|e| e.into())
132 } else {
133 Err(GrpcError::InternalServerError(
134 "there is data left after operation deserialization".to_owned()
135 ))
136 }
137 }
138 Err(e) => {
139 Err(GrpcError::InternalServerError(format!("failed to deserialize operation: {}", e)))
140 }
141 };
142 verified_op_res
143 })
144 .collect();
145
146 match verified_ops_res {
147 // If all operations in the incoming message are valid, store and propagate them
148 Ok(verified_ops) => {
149 let mut operation_storage = storage.clone_without_refs();
150 operation_storage
151 .store_operations(verified_ops.values().cloned().collect());
152 // Propagate the operations to the network before admitting
153 // them into the local pool (same order as the JSON-RPC API)
154 if let Err(e) = protocol_controller
155 .propagate_operations(operation_storage.clone())
156 {
157 // If propagation failed, send an error message back to the client
158 let error =
159 format!("failed to propagate operations: {}", e);
160 report_error(
161 tx.clone(),
162 tonic::Code::Internal,
163 error.to_owned(),
164 )
165 .await;
166 continue;
167 };
168
169 // Add the received operations to the operations pool
170 if let Err(e) =
171 pool_controller.add_operations(operation_storage)
172 {
173 // If pool admission failed, send an error message back to the client
174 let error =
175 format!("failed to add operations to pool: {}", e);
176 report_error(
177 tx.clone(),
178 tonic::Code::Internal,
179 error.to_owned(),
180 )
181 .await;
182 continue;
183 };
184
185 // Build the response message
186 let result = grpc_model::OperationIds {
187 operation_ids: verified_ops.keys().cloned().collect(),
188 };
189 // Send the response message back to the client
190 if let Err(e) = tx
191 .send(Ok(grpc_api::SendOperationsResponse {
192 result: Some(
193 grpc_api::send_operations_response::Result::OperationIds(
194 result,
195 ),
196 ),
197 }))
198 .await
199 {
200 error!("failed to send back operations response: {}", e);
201 };
202 }
203 // If the verification failed, send an error message back to the client
204 Err(e) => {
205 let error = format!("invalid operation(s): {}", e);
206 report_error(
207 tx.clone(),
208 tonic::Code::InvalidArgument,
209 error.to_owned(),
210 )
211 .await;
212 }
213 }
214 }
215 }
216 }
217 // Handle any errors that may occur during receiving the data
218 Err(err) => {
219 // Check if the error matches any IO errors
220 if let Some(io_err) = match_for_io_error(&err) {
221 if io_err.kind() == ErrorKind::BrokenPipe {
222 warn!("client disconnected, broken pipe: {}", io_err);
223 break;
224 }
225 }
226 error!("{}", err);
227 // Send the error response back to the client
228 if let Err(e) = tx.send(Err(err)).await {
229 error!("failed to send back send_operations error response: {}", e);
230 break;
231 }
232 }
233 }
234 }
235 });
236
237 let out_stream = tokio_stream::wrappers::ReceiverStream::new(rx);
238 Ok(Box::pin(out_stream) as SendOperationsStreamType)
239}
240
241// This function reports an error to the sender by sending a gRPC response message to the client
242async fn report_error(
243 sender: tokio::sync::mpsc::Sender<Result<grpc_api::SendOperationsResponse, tonic::Status>>,
244 code: tonic::Code,
245 error: String,
246) {
247 error!("{}", error);
248 // Attempt to send the error response message to the sender
249 if let Err(e) = sender
250 .send(Ok(grpc_api::SendOperationsResponse {
251 result: Some(grpc_api::send_operations_response::Result::Error(
252 grpc_model::Error {
253 code: code.into(),
254 message: error,
255 },
256 )),
257 }))
258 .await
259 {
260 // If sending the message fails, log the error message
261 error!("failed to send back send_operations error response: {}", e);
262 }
263}