massa_grpc/stream/
send_blocks.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::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
20/// Type declaration for SendBlockStream
21pub 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/// This function takes a streaming request of block messages,
30/// verifies, saves and propagates the block received in each message, and sends back a stream of
31/// block id messages
32#[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    // Create a channel to handle communication with the client
44    let (tx, rx) = tokio::sync::mpsc::channel(config.max_channel_size);
45    // Extract the incoming stream of block messages
46    let mut in_stream = request.into_inner();
47
48    // Spawn a task that reads incoming messages and processes the block in each message
49    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                    // Create a block deserializer arguments
64                    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                    // Deserialize and verify received block in the incoming message
74                    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                            // Add the received block to the graph
106                            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                            // Propagate the block(header) to the network
115                            if let Err(e) =
116                                protocol_command_sender.integrated_block(block_id, block_storage)
117                            {
118                                // If propagation failed, send an error message back to the client
119                                report_error(
120                                    tx.clone(),
121                                    tonic::Code::Internal,
122                                    format!("failed to propagate block: {}", e),
123                                )
124                                .await;
125                                continue;
126                            };
127
128                            // Send the response message back to the client
129                            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                        // If the verification failed, send an error message back to the client
141                        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                // Handle any errors that may occur during receiving the data
153                Err(err) => {
154                    // Check if the error matches any IO errors
155                    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                    // Send the error response back to the client
163                    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
176/// This function reports an error to the sender by sending a gRPC response message to the client
177async fn report_error(
178    sender: Sender<Result<grpc_api::SendBlocksResponse, tonic::Status>>,
179    code: tonic::Code,
180    error: String,
181) {
182    error!("{}", error);
183    // Attempt to send the error response message to the sender
184    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        // If sending the message fails, log the error message
196        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    // Block::verify_signature delegates to the nested header signing domain (F87).
206    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}