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}