massa_grpc/stream/send_endorsements.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::endorsement::{EndorsementDeserializer, SecureShareEndorsement};
7use massa_models::secure_share::SecureShareDeserializer;
8use massa_models::timeslots::get_block_slot_timestamp;
9use massa_pos_exports::SelectorController;
10use massa_proto_rs::massa::api::v1 as grpc_api;
11use massa_proto_rs::massa::model::v1 as grpc_model;
12use massa_serialization::{DeserializeError, Deserializer};
13use massa_versioning::consensus_signature::sig_chain_id_for_slot;
14use std::collections::HashMap;
15use std::io::ErrorKind;
16use std::pin::Pin;
17use tracing::{error, warn};
18
19/// Type declaration for SendEndorsements
20pub type SendEndorsementsStreamType = Pin<
21 Box<
22 dyn futures_util::Stream<Item = Result<grpc_api::SendEndorsementsResponse, tonic::Status>>
23 + Send
24 + 'static,
25 >,
26>;
27
28/// This function takes a streaming request of endorsements messages,
29/// verifies, saves and propagates the endorsements received in each message, and sends back a stream of
30/// endorsements ids messages
31pub(crate) async fn send_endorsements(
32 grpc: &MassaPublicGrpc,
33 request: tonic::Request<tonic::Streaming<grpc_api::SendEndorsementsRequest>>,
34) -> Result<SendEndorsementsStreamType, GrpcError> {
35 let mut pool_command_sender = grpc.pool_controller.clone();
36 let protocol_command_sender = grpc.protocol_controller.clone();
37 let selector_controller = grpc.selector_controller.clone();
38 let config = grpc.grpc_config.clone();
39 let storage = grpc.storage.clone_without_refs();
40 let mip_store = grpc.keypair_factory.mip_store.clone();
41
42 // Create a channel to handle communication with the client
43 let (tx, rx) = tokio::sync::mpsc::channel(config.max_channel_size);
44 // Extract the incoming stream of endorsements messages
45 let mut in_stream = request.into_inner();
46
47 // Spawn a task that reads incoming messages and processes the endorsements in each message
48 tokio::spawn(async move {
49 while let Some(result) = in_stream.next().await {
50 match result {
51 Ok(req_content) => {
52 // If the incoming message has no endorsements, send an error message back to the client
53 if req_content.endorsements.is_empty() {
54 report_error(
55 tx.clone(),
56 tonic::Code::InvalidArgument,
57 "the request payload is empty".to_owned(),
58 )
59 .await;
60 } else {
61 // If there are too many endorsements in the incoming message, send an error message back to the client
62 let proto_endorsement = req_content.endorsements;
63 if proto_endorsement.len() as u32 > config.max_endorsements_per_message {
64 report_error(
65 tx.clone(),
66 tonic::Code::InvalidArgument,
67 "too many endorsements per message".to_owned(),
68 )
69 .await;
70 } else {
71 // Deserialize and verify each endorsement in the incoming message
72 let endorsement_deserializer = SecureShareDeserializer::new(
73 EndorsementDeserializer::new(
74 config.thread_count,
75 config.endorsement_count,
76 ),
77 config.chain_id,
78 );
79 let verified_eds_res: Result<HashMap<String, SecureShareEndorsement>, GrpcError> = proto_endorsement
80 .into_iter()
81 .map(|proto_endorsement| {
82 let verified_op = match endorsement_deserializer.deserialize::<DeserializeError>(&proto_endorsement) {
83 Ok(tuple) => {
84 // Deserialize the endorsement and verify its signature
85 let (rest, res_endorsement): (&[u8], SecureShareEndorsement) = tuple;
86 if rest.is_empty() {
87 let slot_ts = get_block_slot_timestamp(
88 config.thread_count,
89 config.t0,
90 config.genesis_timestamp,
91 res_endorsement.content.slot,
92 )?;
93 let sig_chain_id = sig_chain_id_for_slot(
94 &mip_store,
95 config.chain_id,
96 slot_ts,
97 );
98 res_endorsement.verify_signature(sig_chain_id)
99 .map(|_| (res_endorsement.id.to_string(), res_endorsement))
100 .map_err(|e| e.into())
101 } else {
102 Err(GrpcError::InternalServerError(
103 "there is data left after endorsement deserialization".to_owned()
104 ))
105 }
106 }
107 Err(e) => {
108 Err(GrpcError::InternalServerError(format!("failed to deserialize endorsement: {}", e)
109 ))
110 }
111 };
112 verified_op
113 })
114 .collect();
115
116 match verified_eds_res {
117 // If all endorsements in the incoming message are valid, store and propagate them
118 Ok(verified_eds) => {
119 // Check the PoS draws before letting those endorsements in: the
120 // propagation flow trusts its caller and does not check them again.
121 if let Err(error) = check_endorsement_draws(
122 &verified_eds,
123 selector_controller.as_ref(),
124 ) {
125 report_error(
126 tx.clone(),
127 tonic::Code::InvalidArgument,
128 error,
129 )
130 .await;
131 continue;
132 }
133
134 let mut endorsement_storage = storage.clone_without_refs();
135 endorsement_storage.store_endorsements(
136 verified_eds.values().cloned().collect(),
137 );
138 // Add the received endorsements to the endorsements pool
139 pool_command_sender
140 .add_endorsements(endorsement_storage.clone());
141
142 // Propagate the endorsements to the network
143 if let Err(e) = protocol_command_sender
144 .propagate_endorsements(endorsement_storage)
145 {
146 // If propagation failed, send an error message back to the client
147 let error =
148 format!("failed to propagate endorsement: {}", e);
149 report_error(
150 tx.clone(),
151 tonic::Code::Internal,
152 error.to_owned(),
153 )
154 .await;
155 };
156
157 // Build the response message
158 let result = grpc_model::EndorsementIds {
159 endorsement_ids: verified_eds.keys().cloned().collect(),
160 };
161 // Send the response message back to the client
162 if let Err(e) = tx
163 .send(Ok(grpc_api::SendEndorsementsResponse {
164 result: Some(
165 grpc_api::send_endorsements_response::Result::EndorsementIds(
166 result,
167 ),
168 ),
169 }))
170 .await
171 {
172 error!("failed to send back endorsement response: {}", e)
173 };
174 }
175 // If the verification failed, send an error message back to the client
176 Err(e) => {
177 let error = format!("invalid endorsement(s): {}", e);
178 report_error(
179 tx.clone(),
180 tonic::Code::InvalidArgument,
181 error.to_owned(),
182 )
183 .await;
184 }
185 }
186 }
187 }
188 }
189 // Handles errors that occur while sending a response back to a client
190 Err(err) => {
191 // Check if the error matches any IO errors
192 if let Some(io_err) = match_for_io_error(&err) {
193 if io_err.kind() == ErrorKind::BrokenPipe {
194 warn!("client disconnected, broken pipe: {}", io_err);
195 break;
196 }
197 }
198 error!("{}", err);
199 // Send the error response back to the client
200 if let Err(e) = tx.send(Err(err)).await {
201 error!(
202 "failed to send back send_endorsements error response: {}",
203 e
204 );
205 break;
206 }
207 }
208 }
209 }
210 });
211
212 let out_stream = tokio_stream::wrappers::ReceiverStream::new(rx);
213 Ok(Box::pin(out_stream) as SendEndorsementsStreamType)
214}
215
216/// Checks that each endorsement was created by the endorser drawn for its `(slot, index)` pair.
217///
218/// This mirrors the check done on endorsements received from peers (`note_endorsements_from_peer`)
219/// and in the endorsement pool: endorsements accepted here are handed over to the protocol
220/// propagation flow, which trusts its caller and rebroadcasts them without further validation.
221fn check_endorsement_draws(
222 endorsements: &HashMap<String, SecureShareEndorsement>,
223 selector_controller: &dyn SelectorController,
224) -> Result<(), String> {
225 for endorsement in endorsements.values() {
226 let selection = selector_controller
227 .get_selection(endorsement.content.slot)
228 .map_err(|e| {
229 format!(
230 "failed to get the PoS draw at slot {}: {}",
231 endorsement.content.slot, e
232 )
233 })?
234 .endorsements;
235 let Some(address) = selection.get(endorsement.content.index as usize) else {
236 return Err(format!(
237 "no selection at slot {} for index {}",
238 endorsement.content.slot, endorsement.content.index
239 ));
240 };
241 if address != &endorsement.content_creator_address {
242 return Err(format!(
243 "invalid endorsement producer selection at slot {} index {}: expected address {}, got {}",
244 endorsement.content.slot,
245 endorsement.content.index,
246 address,
247 endorsement.content_creator_address
248 ));
249 }
250 }
251 Ok(())
252}
253
254// This function reports an error to the sender by sending a gRPC response message to the client
255async fn report_error(
256 sender: tokio::sync::mpsc::Sender<Result<grpc_api::SendEndorsementsResponse, tonic::Status>>,
257 code: tonic::Code,
258 error: String,
259) {
260 error!("{}", error);
261 // Attempt to send the error response message to the sender
262 if let Err(e) = sender
263 .send(Ok(grpc_api::SendEndorsementsResponse {
264 result: Some(grpc_api::send_endorsements_response::Result::Error(
265 massa_proto_rs::massa::model::v1::Error {
266 code: code.into(),
267 message: error,
268 },
269 )),
270 }))
271 .await
272 {
273 // If sending the message fails, log the error message
274 error!(
275 "failed to send back send_endorsements error response: {}",
276 e
277 );
278 }
279}