sc_consensus_beefy/communication/request_response/
incoming_requests_handler.rs1use codec::DecodeAll;
21use futures::{channel::oneshot, StreamExt};
22use log::{debug, trace};
23use sc_client_api::BlockBackend;
24use sc_network::{
25 config as netconfig, service::traits::RequestResponseConfig, types::ProtocolName,
26 NetworkBackend, ReputationChange,
27};
28use sc_network_types::PeerId;
29use sp_consensus_beefy::BEEFY_ENGINE_ID;
30use sp_runtime::traits::Block;
31use std::{marker::PhantomData, sync::Arc};
32
33use crate::{
34 communication::{
35 cost,
36 request_response::{
37 on_demand_justifications_protocol_config, Error, JustificationRequest,
38 BEEFY_SYNC_LOG_TARGET,
39 },
40 },
41 metric_inc,
42 metrics::{register_metrics, OnDemandIncomingRequestsMetrics},
43};
44
45#[derive(Debug)]
47pub(crate) struct IncomingRequest<B: Block> {
48 pub peer: PeerId,
50 pub payload: JustificationRequest<B>,
52 pub pending_response: oneshot::Sender<netconfig::OutgoingResponse>,
54}
55
56impl<B: Block> IncomingRequest<B> {
57 pub fn new(
59 peer: PeerId,
60 payload: JustificationRequest<B>,
61 pending_response: oneshot::Sender<netconfig::OutgoingResponse>,
62 ) -> Self {
63 Self { peer, payload, pending_response }
64 }
65
66 pub fn try_from_raw<F>(
75 raw: netconfig::IncomingRequest,
76 reputation_changes_on_err: F,
77 ) -> Result<Self, Error>
78 where
79 F: FnOnce(usize) -> Vec<ReputationChange>,
80 {
81 let netconfig::IncomingRequest { payload, peer, pending_response } = raw;
82 let payload = match JustificationRequest::decode_all(&mut payload.as_ref()) {
83 Ok(payload) => payload,
84 Err(err) => {
85 let response = netconfig::OutgoingResponse {
86 result: Err(()),
87 reputation_changes: reputation_changes_on_err(payload.len()),
88 sent_feedback: None,
89 };
90 if let Err(_) = pending_response.send(response) {
91 return Err(Error::DecodingErrorNoReputationChange(peer, err));
92 }
93 return Err(Error::DecodingError(peer, err));
94 },
95 };
96 Ok(Self::new(peer, payload, pending_response))
97 }
98}
99
100pub(crate) struct IncomingRequestReceiver {
104 raw: async_channel::Receiver<netconfig::IncomingRequest>,
105}
106
107impl IncomingRequestReceiver {
108 pub fn new(inner: async_channel::Receiver<netconfig::IncomingRequest>) -> Self {
109 Self { raw: inner }
110 }
111
112 pub async fn recv<B, F>(&mut self, reputation_changes: F) -> Result<IncomingRequest<B>, Error>
117 where
118 B: Block,
119 F: FnOnce(usize) -> Vec<ReputationChange>,
120 {
121 let req = match self.raw.next().await {
122 None => return Err(Error::RequestChannelExhausted),
123 Some(raw) => IncomingRequest::<B>::try_from_raw(raw, reputation_changes)?,
124 };
125 Ok(req)
126 }
127}
128
129pub struct BeefyJustifsRequestHandler<B, Client> {
131 pub(crate) request_receiver: IncomingRequestReceiver,
132 pub(crate) justif_protocol_name: ProtocolName,
133 pub(crate) client: Arc<Client>,
134 pub(crate) metrics: Option<OnDemandIncomingRequestsMetrics>,
135 pub(crate) _block: PhantomData<B>,
136}
137
138impl<B, Client> BeefyJustifsRequestHandler<B, Client>
139where
140 B: Block,
141 Client: BlockBackend<B> + Send + Sync,
142{
143 pub fn new<Hash: AsRef<[u8]>, Network: NetworkBackend<B, <B as Block>::Hash>>(
145 genesis_hash: Hash,
146 fork_id: Option<&str>,
147 client: Arc<Client>,
148 prometheus_registry: Option<prometheus_endpoint::Registry>,
149 ) -> (Self, Network::RequestResponseProtocolConfig) {
150 let (request_receiver, config): (_, Network::RequestResponseProtocolConfig) =
151 on_demand_justifications_protocol_config::<_, _, Network>(genesis_hash, fork_id);
152 let justif_protocol_name = config.protocol_name().clone();
153 let metrics = register_metrics(prometheus_registry);
154 (
155 Self { request_receiver, justif_protocol_name, client, metrics, _block: PhantomData },
156 config,
157 )
158 }
159
160 pub fn protocol_name(&self) -> ProtocolName {
162 self.justif_protocol_name.clone()
163 }
164
165 fn handle_request(&self, request: IncomingRequest<B>) -> Result<(), Error> {
167 let mut reputation_changes = vec![];
168 let maybe_encoded_proof = self
169 .client
170 .block_hash(request.payload.begin)
171 .ok()
172 .flatten()
173 .and_then(|hash| self.client.justifications(hash).ok().flatten())
174 .and_then(|justifs| justifs.get(BEEFY_ENGINE_ID).cloned())
175 .ok_or_else(|| reputation_changes.push(cost::UNKNOWN_PROOF_REQUEST));
176 request
177 .pending_response
178 .send(netconfig::OutgoingResponse {
179 result: maybe_encoded_proof,
180 reputation_changes,
181 sent_feedback: None,
182 })
183 .map_err(|_| Error::SendResponse)
184 }
185
186 pub async fn run(&mut self) -> Error {
190 trace!(target: BEEFY_SYNC_LOG_TARGET, "🥩 Running BeefyJustifsRequestHandler");
191
192 loop {
193 let request = match self
194 .request_receiver
195 .recv(|bytes| {
196 let bytes = bytes.min(i32::MAX as usize) as i32;
197 vec![ReputationChange::new(
198 bytes.saturating_mul(cost::PER_UNDECODABLE_BYTE),
199 "BEEFY: Bad request payload",
200 )]
201 })
202 .await
203 {
204 Ok(request) => request,
205 Err(
206 e @ (Error::DecodingError(_, _) | Error::DecodingErrorNoReputationChange(_, _)),
207 ) => {
208 metric_inc!(self.metrics, beefy_failed_justification_responses);
211 debug!(
212 target: BEEFY_SYNC_LOG_TARGET,
213 "🥩 Ignoring invalid BEEFY justification request: {}", e,
214 );
215 continue;
216 },
217 Err(e) => return e,
218 };
219
220 let peer = request.peer;
221 match self.handle_request(request) {
222 Ok(()) => {
223 metric_inc!(self.metrics, beefy_successful_justification_responses);
224 debug!(
225 target: BEEFY_SYNC_LOG_TARGET,
226 "🥩 Handled BEEFY justification request from {:?}.", peer
227 )
228 },
229 Err(e) => {
230 metric_inc!(self.metrics, beefy_failed_justification_responses);
232 debug!(
233 target: BEEFY_SYNC_LOG_TARGET,
234 "🥩 Failed to handle BEEFY justification request from {:?}: {}", peer, e,
235 )
236 },
237 }
238 }
239 }
240}
241
242#[cfg(test)]
243mod tests {
244 use super::*;
245 use crate::communication::request_response::JUSTIF_CHANNEL_SIZE;
246 use codec::Encode;
247 use sc_block_builder::BlockBuilderBuilder;
248 use sc_network::config::OutgoingResponse;
249 use sp_blockchain::HeaderBackend;
250 use substrate_test_runtime_client::{
251 runtime::Block as TestBlock, Backend, Client, ClientBlockImportExt, ClientExt,
252 DefaultTestClientBuilderExt, TestClientBuilder, TestClientBuilderExt,
253 };
254
255 type TestClient = Client<Backend>;
256
257 fn test_handler(
258 client: TestClient,
259 ) -> (
260 async_channel::Sender<netconfig::IncomingRequest>,
261 BeefyJustifsRequestHandler<TestBlock, TestClient>,
262 ) {
263 let (tx, rx) = async_channel::bounded(JUSTIF_CHANNEL_SIZE);
264 let handler = BeefyJustifsRequestHandler {
265 request_receiver: IncomingRequestReceiver::new(rx),
266 justif_protocol_name: ProtocolName::Static("/beefy/justifications/1"),
267 client: Arc::new(client),
268 metrics: None,
269 _block: PhantomData,
270 };
271 (tx, handler)
272 }
273
274 async fn send_request(
275 tx: &async_channel::Sender<netconfig::IncomingRequest>,
276 payload: Vec<u8>,
277 ) -> Result<OutgoingResponse, oneshot::Canceled> {
278 let (pending_response, rx) = oneshot::channel();
279 tx.send(netconfig::IncomingRequest { peer: PeerId::random(), payload, pending_response })
280 .await
281 .unwrap();
282 rx.await
283 }
284
285 #[tokio::test]
286 async fn misbehaving_sending_peers_are_penalized() {
287 let (tx, mut handler) = test_handler(TestClientBuilder::new().build());
288 let handler_task = tokio::spawn(async move { handler.run().await });
289
290 let response = send_request(&tx, vec![0xff]).await.expect("handler answers the request");
292 assert_eq!(response.result, Err(()));
293 assert!(!response.reputation_changes.is_empty());
294
295 let mut payload = JustificationRequest::<TestBlock> { begin: 1 }.encode();
297 payload.push(0x00);
298 let response = send_request(&tx, payload).await.expect("handler answers the request");
299 assert_eq!(response.result, Err(()));
300 assert!(!response.reputation_changes.is_empty());
301
302 let payload = JustificationRequest::<TestBlock> { begin: 1 }.encode();
306 let response = send_request(&tx, payload).await.expect("handler answers the request");
307 assert_eq!(response.result, Err(()));
308 assert_eq!(response.reputation_changes, vec![cost::UNKNOWN_PROOF_REQUEST]);
309
310 drop(tx);
312 handler_task.await.unwrap();
313 }
314
315 #[tokio::test]
316 async fn known_justification_is_served() {
317 let client = TestClientBuilder::new().build();
318 let justif = vec![42u8];
319
320 let block = BlockBuilderBuilder::new(&client)
322 .on_parent_block(client.info().genesis_hash)
323 .with_parent_block_number(0)
324 .build()
325 .unwrap()
326 .build()
327 .unwrap()
328 .block;
329 let hash = block.header.hash();
330 client.import(sp_consensus::BlockOrigin::Own, block).await.unwrap();
331 client.finalize_block(hash, Some((BEEFY_ENGINE_ID, justif.clone()))).unwrap();
332
333 let (tx, mut handler) = test_handler(client);
334 let handler_task = tokio::spawn(async move { handler.run().await });
335
336 let payload = JustificationRequest::<TestBlock> { begin: 1 }.encode();
337 let response = send_request(&tx, payload).await.expect("handler answers the request");
338 assert_eq!(response.result, Ok(justif));
339 assert!(response.reputation_changes.is_empty());
340
341 drop(tx);
342 handler_task.await.unwrap();
343 }
344}