massa_grpc/stream/
send_blocks.rs1use crate::error::{match_for_io_error, GrpcError};
4use crate::server::MassaPublicGrpc;
5use futures_util::StreamExt;
6use massa_models::block::{BlockDeserializer, BlockDeserializerArgs, SecureShareBlock};
7use massa_models::error::ModelsError;
8use massa_models::secure_share::SecureShareDeserializer;
9use massa_models::timeslots::get_block_slot_timestamp;
10use massa_proto_rs::massa::api::v1 as grpc_api;
11use massa_serialization::{DeserializeError, Deserializer};
12use massa_versioning::consensus_signature::sig_chain_id_for_slot;
13use massa_versioning::versioning::MipStore;
14use std::io::ErrorKind;
15use std::pin::Pin;
16use tokio::sync::mpsc::Sender;
17use tonic::Request;
18use tracing::{error, warn};
19
20pub type SendBlocksStreamType = Pin<
22 Box<
23 dyn futures_util::Stream<Item = Result<grpc_api::SendBlocksResponse, tonic::Status>>
24 + Send
25 + 'static,
26 >,
27>;
28
29#[allow(dead_code)]
33pub(crate) async fn send_blocks(
34 grpc: &MassaPublicGrpc,
35 request: Request<tonic::Streaming<grpc_api::SendBlocksRequest>>,
36) -> Result<SendBlocksStreamType, GrpcError> {
37 let consensus_controller = grpc.consensus_controller.clone();
38 let protocol_command_sender = grpc.protocol_controller.clone();
39 let config = grpc.grpc_config.clone();
40 let storage = grpc.storage.clone_without_refs();
41 let mip_store = grpc.keypair_factory.mip_store.clone();
42
43 let (tx, rx) = tokio::sync::mpsc::channel(config.max_channel_size);
45 let mut in_stream = request.into_inner();
47
48 tokio::spawn(async move {
50 while let Some(result) = in_stream.next().await {
51 match result {
52 Ok(req_content) => {
53 if req_content.block.is_empty() {
54 report_error(
55 tx.clone(),
56 tonic::Code::InvalidArgument,
57 "the request payload is empty".to_owned(),
58 )
59 .await;
60 continue;
61 };
62
63 let args = BlockDeserializerArgs {
65 thread_count: config.thread_count,
66 max_operations_per_block: config.max_operations_per_block,
67 endorsement_count: config.endorsement_count,
68 max_denunciations_per_block_header: config
69 .max_denunciations_per_block_header,
70 last_start_period: Some(config.last_start_period),
71 chain_id: config.chain_id,
72 };
73 match SecureShareDeserializer::new(
75 BlockDeserializer::new(args),
76 config.chain_id,
77 )
78 .deserialize::<DeserializeError>(&req_content.block)
79 {
80 Ok(tuple) => {
81 let (rest, res_block): (&[u8], SecureShareBlock) = tuple;
82 if !rest.is_empty() {
83 report_error(
84 tx.clone(),
85 tonic::Code::InvalidArgument,
86 "the request payload is too large".to_owned(),
87 )
88 .await;
89 continue;
90 }
91 if let Err(e) = verify_received_block(&res_block, &config, &mip_store) {
92 report_error(
93 tx.clone(),
94 tonic::Code::InvalidArgument,
95 format!("wrong signature: {}", e),
96 )
97 .await;
98 continue;
99 }
100
101 let block_id = res_block.id;
102 let slot = res_block.content.header.content.slot;
103 let mut block_storage = storage.clone_without_refs();
104
105 block_storage.store_block(res_block.clone());
107 consensus_controller.register_block(
108 block_id,
109 slot,
110 block_storage.clone(),
111 false,
112 );
113
114 if let Err(e) =
116 protocol_command_sender.integrated_block(block_id, block_storage)
117 {
118 report_error(
120 tx.clone(),
121 tonic::Code::Internal,
122 format!("failed to propagate block: {}", e),
123 )
124 .await;
125 continue;
126 };
127
128 if let Err(e) = tx
130 .send(Ok(grpc_api::SendBlocksResponse {
131 result: Some(grpc_api::send_blocks_response::Result::BlockId(
132 res_block.id.to_string(),
133 )),
134 }))
135 .await
136 {
137 error!("failed to send back block response: {}", e);
138 };
139 }
140 Err(e) => {
142 report_error(
143 tx.clone(),
144 tonic::Code::InvalidArgument,
145 format!("failed to deserialize block: {}", e),
146 )
147 .await;
148 continue;
149 }
150 };
151 }
152 Err(err) => {
154 if let Some(io_err) = match_for_io_error(&err) {
156 if io_err.kind() == ErrorKind::BrokenPipe {
157 warn!("client disconnected, broken pipe: {}", io_err);
158 break;
159 }
160 }
161 error!("{}", err);
162 if let Err(e) = tx.send(Err(err)).await {
164 error!("failed to send back send_blocks error response: {}", e);
165 break;
166 }
167 }
168 }
169 }
170 });
171
172 let out_stream = tokio_stream::wrappers::ReceiverStream::new(rx);
173 Ok(Box::pin(out_stream) as SendBlocksStreamType)
174}
175
176async fn report_error(
178 sender: Sender<Result<grpc_api::SendBlocksResponse, tonic::Status>>,
179 code: tonic::Code,
180 error: String,
181) {
182 error!("{}", error);
183 if let Err(e) = sender
185 .send(Ok(grpc_api::SendBlocksResponse {
186 result: Some(grpc_api::send_blocks_response::Result::Error(
187 massa_proto_rs::massa::model::v1::Error {
188 code: code.into(),
189 message: error,
190 },
191 )),
192 }))
193 .await
194 {
195 error!("failed to send back send_blocks error response: {}", e);
197 }
198}
199
200fn verify_received_block(
201 block: &SecureShareBlock,
202 config: &crate::config::GrpcConfig,
203 mip_store: &MipStore,
204) -> Result<(), ModelsError> {
205 let header = &block.content.header;
207 let header_ts = get_block_slot_timestamp(
208 config.thread_count,
209 config.t0,
210 config.genesis_timestamp,
211 header.content.slot,
212 )?;
213 block.verify_signature(sig_chain_id_for_slot(mip_store, config.chain_id, header_ts))?;
214 for endorsement in header.content.endorsements.iter() {
215 let slot_ts = get_block_slot_timestamp(
216 config.thread_count,
217 config.t0,
218 config.genesis_timestamp,
219 endorsement.content.slot,
220 )?;
221 endorsement.verify_signature(sig_chain_id_for_slot(mip_store, config.chain_id, slot_ts))?;
222 }
223 Ok(())
224}