referrerpolicy=no-referrer-when-downgrade

sc_network/bitswap/
mod.rs

1// Copyright (C) Parity Technologies (UK) Ltd.
2// This file is part of Substrate.
3// SPDX-License-Identifier: GPL-3.0-or-later WITH Classpath-exception-2.0
4
5// Substrate is free software: you can redistribute it and/or modify
6// it under the terms of the GNU General Public License as published by
7// the Free Software Foundation, either version 3 of the License, or
8// (at your option) any later version.
9
10// Substrate is distributed in the hope that it will be useful,
11// but WITHOUT ANY WARRANTY; without even the implied warranty of
12// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
13// GNU General Public License for more details.
14
15// You should have received a copy of the GNU General Public License
16// along with Substrate. If not, see <https://www.gnu.org/licenses/>.
17
18//! Bitswap server for Substrate.
19//!
20//! Supports querying indexed transactions by hash over the standard bitswap protocol (v1.2.0).
21//! CIDs must reference a supported 256-bit transaction hash.
22
23use crate::{
24	request_responses::{IncomingRequest, OutgoingResponse, ProtocolConfig},
25	types::ProtocolName,
26	MAX_RESPONSE_SIZE,
27};
28
29use cid::{Error as CidError, Version as CidVersion};
30use futures::StreamExt;
31use log::{debug, error, trace};
32use prost::Message;
33use sc_client_api::BlockBackend;
34use sc_network_types::PeerId;
35use schema::bitswap::{
36	message::{wantlist::WantType, Block as MessageBlock, BlockPresence, BlockPresenceType},
37	Message as BitswapMessage,
38};
39use sp_core::H256;
40use sp_runtime::traits::Block as BlockT;
41use std::{io, sync::Arc, time::Duration};
42use unsigned_varint::encode as varint_encode;
43
44/// Bitswap client.
45mod client;
46/// Bitswap protobuf schema, generated from the protocol definitions.
47pub mod schema;
48
49pub use cid::Cid;
50
51pub use client::{
52	request_bitswap_blocks, request_bitswap_blocks_unverified, BitswapError, FetchOutcome,
53	BLAKE2B_256_MULTIHASH_CODE, KECCAK_256_MULTIHASH_CODE, SHA2_256_MULTIHASH_CODE,
54};
55
56pub(crate) use schema::bitswap::Message as BitswapProtoMessage;
57
58pub(crate) const LOG_TARGET: &str = "sub-libp2p::bitswap";
59
60// Use the network-wide response cap for Bitswap messages.
61const MAX_PACKET_SIZE: u64 = MAX_RESPONSE_SIZE;
62
63/// Max number of queued responses before denying requests.
64const MAX_REQUEST_QUEUE: usize = 20;
65
66/// Max number of blocks per wantlist.
67pub const MAX_WANTED_BLOCKS: usize = 16;
68
69/// Bitswap protocol name.
70pub(crate) const PROTOCOL_NAME: &str = "/ipfs/bitswap/1.2.0";
71
72/// IPFS raw multicodec used for indexed transaction payload bytes.
73pub const RAW_CODEC: u64 = 0x55;
74
75/// Check if a CID is supported by the bitswap protocol — CIDv1, 32-byte digest, with a
76/// supported multihash code (Blake2b-256, SHA2-256, or Keccak-256).
77pub fn is_cid_supported(cid: &Cid) -> bool {
78	cid.version() != CidVersion::V0 &&
79		cid.hash().size() == 32 &&
80		is_supported_multihash_code(cid.hash().code())
81}
82
83/// Return `true` if `code` is a supported multihash code.
84pub(crate) fn is_supported_multihash_code(code: u64) -> bool {
85	matches!(code, BLAKE2B_256_MULTIHASH_CODE | SHA2_256_MULTIHASH_CODE | KECCAK_256_MULTIHASH_CODE)
86}
87
88/// CID metadata without the actual content bytes.
89#[derive(PartialEq, Eq, Clone, Debug)]
90pub struct Prefix {
91	/// The version of CID.
92	pub version: CidVersion,
93	/// The codec of CID.
94	pub codec: u64,
95	/// The multihash type of CID.
96	pub mh_type: u64,
97	/// The multihash length of CID.
98	pub mh_len: u8,
99}
100
101impl From<&Cid> for Prefix {
102	fn from(cid: &Cid) -> Self {
103		Self {
104			version: cid.version(),
105			codec: cid.codec(),
106			mh_type: cid.hash().code(),
107			mh_len: cid.hash().size(),
108		}
109	}
110}
111
112impl Prefix {
113	/// Convert the prefix to encoded bytes.
114	pub fn to_bytes(&self) -> Vec<u8> {
115		let mut res = Vec::with_capacity(4);
116		let mut buf = varint_encode::u64_buffer();
117		let version = varint_encode::u64(self.version.into(), &mut buf);
118		res.extend_from_slice(version);
119		let mut buf = varint_encode::u64_buffer();
120		let codec = varint_encode::u64(self.codec, &mut buf);
121		res.extend_from_slice(codec);
122		let mut buf = varint_encode::u64_buffer();
123		let mh_type = varint_encode::u64(self.mh_type, &mut buf);
124		res.extend_from_slice(mh_type);
125		let mut buf = varint_encode::u64_buffer();
126		let mh_len = varint_encode::u64(self.mh_len as u64, &mut buf);
127		res.extend_from_slice(mh_len);
128		res
129	}
130}
131
132/// Bitswap request handler.
133pub(crate) struct BitswapRequestHandler<B> {
134	client: Arc<dyn BlockBackend<B> + Send + Sync>,
135	request_receiver: async_channel::Receiver<IncomingRequest>,
136}
137
138impl<B: BlockT> BitswapRequestHandler<B> {
139	/// Create a new [`BitswapRequestHandler`].
140	pub(crate) fn new(client: Arc<dyn BlockBackend<B> + Send + Sync>) -> (Self, ProtocolConfig) {
141		let (tx, request_receiver) = async_channel::bounded(MAX_REQUEST_QUEUE);
142
143		let config = ProtocolConfig {
144			name: ProtocolName::from(PROTOCOL_NAME),
145			fallback_names: vec![],
146			max_request_size: MAX_PACKET_SIZE,
147			max_response_size: MAX_PACKET_SIZE,
148			request_timeout: Duration::from_secs(15),
149			inbound_queue: Some(tx),
150		};
151
152		(Self { client, request_receiver }, config)
153	}
154
155	/// Run [`BitswapRequestHandler`].
156	pub(crate) async fn run(mut self) {
157		while let Some(request) = self.request_receiver.next().await {
158			let IncomingRequest { peer, payload, pending_response } = request;
159
160			match self.handle_message(&peer, &payload) {
161				Ok(response) => {
162					let response = OutgoingResponse {
163						result: Ok(response),
164						reputation_changes: Vec::new(),
165						sent_feedback: None,
166					};
167
168					match pending_response.send(response) {
169						Ok(()) => {
170							trace!(target: LOG_TARGET, "Handled bitswap request from {peer}.",)
171						},
172						Err(_) => debug!(
173							target: LOG_TARGET,
174							"Failed to handle bitswap request from {peer}: {}",
175							RequestHandlerError::SendResponse,
176						),
177					}
178				},
179				Err(err) => {
180					error!(target: LOG_TARGET, "Failed to process request from {peer}: {err}");
181
182					// TODO: adjust reputation?
183
184					let response = OutgoingResponse {
185						result: Err(()),
186						reputation_changes: vec![],
187						sent_feedback: None,
188					};
189
190					if pending_response.send(response).is_err() {
191						debug!(
192							target: LOG_TARGET,
193							"Failed to handle bitswap request from {peer}: {}",
194							RequestHandlerError::SendResponse,
195						);
196					}
197				},
198			}
199		}
200	}
201
202	/// Handle received Bitswap request
203	fn handle_message(
204		&mut self,
205		peer: &PeerId,
206		payload: &[u8],
207	) -> Result<Vec<u8>, RequestHandlerError> {
208		let request = schema::bitswap::Message::decode(payload)?;
209
210		trace!(target: LOG_TARGET, "Received request: {:?} from {}", request, peer);
211
212		let mut response = BitswapMessage::default();
213
214		let wantlist = match request.wantlist {
215			Some(wantlist) => wantlist,
216			None => {
217				debug!(target: LOG_TARGET, "Unexpected bitswap message from {}", peer);
218				return Err(RequestHandlerError::InvalidWantList);
219			},
220		};
221
222		if wantlist.entries.len() > MAX_WANTED_BLOCKS {
223			trace!(target: LOG_TARGET, "Ignored request: too many entries");
224			return Err(RequestHandlerError::TooManyEntries);
225		}
226
227		for entry in wantlist.entries {
228			let cid = match Cid::read_bytes(entry.block.as_slice()) {
229				Ok(cid) => cid,
230				Err(e) => {
231					trace!(target: LOG_TARGET, "Bad CID {:?}: {:?}", entry.block, e);
232					continue;
233				},
234			};
235
236			if !is_cid_supported(&cid) {
237				trace!(target: LOG_TARGET, "Ignoring unsupported CID {}: {}", peer, cid);
238				continue;
239			}
240
241			let mut hash = H256::default();
242			hash.as_mut().copy_from_slice(&cid.hash().digest()[0..32]);
243			let transaction = match self.client.indexed_transaction(hash) {
244				Ok(ex) => ex,
245				Err(e) => {
246					error!(target: LOG_TARGET, "Error retrieving transaction {}: {}", hash, e);
247					None
248				},
249			};
250
251			match transaction {
252				Some(transaction) => {
253					trace!(target: LOG_TARGET, "Found CID {:?}, hash {:?}", cid, hash);
254
255					if entry.want_type == WantType::Block as i32 {
256						let prefix: Prefix = (&cid).into();
257						response
258							.payload
259							.push(MessageBlock { prefix: prefix.to_bytes(), data: transaction });
260					} else {
261						response.block_presences.push(BlockPresence {
262							r#type: BlockPresenceType::Have as i32,
263							cid: cid.to_bytes(),
264						});
265					}
266				},
267				None => {
268					trace!(target: LOG_TARGET, "Missing CID {:?}, hash {:?}", cid, hash);
269
270					if entry.send_dont_have {
271						response.block_presences.push(BlockPresence {
272							r#type: BlockPresenceType::DontHave as i32,
273							cid: cid.to_bytes(),
274						});
275					}
276				},
277			}
278		}
279
280		Ok(response.encode_to_vec())
281	}
282}
283
284/// Bitswap protocol error.
285#[derive(Debug, thiserror::Error)]
286enum RequestHandlerError {
287	/// Protobuf decoding error.
288	#[error("Failed to decode request: {0}.")]
289	DecodeProto(#[from] prost::DecodeError),
290
291	/// Protobuf encoding error.
292	#[error("Failed to encode response: {0}.")]
293	EncodeProto(#[from] prost::EncodeError),
294
295	/// Client backend error.
296	#[error(transparent)]
297	Client(#[from] sp_blockchain::Error),
298
299	/// Error parsing CID
300	#[error(transparent)]
301	BadCid(#[from] CidError),
302
303	/// Packet read error.
304	#[error(transparent)]
305	Read(#[from] io::Error),
306
307	/// Error sending response.
308	#[error("Failed to send response.")]
309	SendResponse,
310
311	/// Message doesn't have a WANT list.
312	#[error("Invalid WANT list.")]
313	InvalidWantList,
314
315	/// Too many blocks requested.
316	#[error("Too many block entries in the request.")]
317	TooManyEntries,
318}
319
320#[cfg(test)]
321mod tests {
322	use super::*;
323	use futures::channel::oneshot;
324	use litep2p::types::multihash::Code as LiteP2pCode;
325	use sc_block_builder::BlockBuilderBuilder;
326	use schema::bitswap::{
327		message::{wantlist::Entry, Wantlist},
328		Message as BitswapMessage,
329	};
330	use sp_consensus::BlockOrigin;
331	use sp_runtime::codec::Encode;
332	use substrate_test_runtime::ExtrinsicBuilder;
333	use substrate_test_runtime_client::{self, prelude::*, TestClientBuilder};
334
335	#[tokio::test]
336	async fn undecodable_message() {
337		let client = substrate_test_runtime_client::new();
338		let (bitswap, config) = BitswapRequestHandler::new(Arc::new(client));
339
340		tokio::spawn(async move { bitswap.run().await });
341
342		let (tx, rx) = oneshot::channel();
343		config
344			.inbound_queue
345			.unwrap()
346			.send(IncomingRequest {
347				peer: PeerId::random(),
348				payload: vec![0x13, 0x37, 0x13, 0x38],
349				pending_response: tx,
350			})
351			.await
352			.unwrap();
353
354		if let Ok(OutgoingResponse { result, reputation_changes, sent_feedback }) = rx.await {
355			assert_eq!(result, Err(()));
356			assert_eq!(reputation_changes, Vec::new());
357			assert!(sent_feedback.is_none());
358		} else {
359			panic!("invalid event received");
360		}
361	}
362
363	#[tokio::test]
364	async fn empty_want_list() {
365		let client = substrate_test_runtime_client::new();
366		let (bitswap, mut config) = BitswapRequestHandler::new(Arc::new(client));
367
368		tokio::spawn(async move { bitswap.run().await });
369
370		let (tx, rx) = oneshot::channel();
371		config
372			.inbound_queue
373			.as_mut()
374			.unwrap()
375			.send(IncomingRequest {
376				peer: PeerId::random(),
377				payload: BitswapMessage { wantlist: None, ..Default::default() }.encode_to_vec(),
378				pending_response: tx,
379			})
380			.await
381			.unwrap();
382
383		if let Ok(OutgoingResponse { result, reputation_changes, sent_feedback }) = rx.await {
384			assert_eq!(result, Err(()));
385			assert_eq!(reputation_changes, Vec::new());
386			assert!(sent_feedback.is_none());
387		} else {
388			panic!("invalid event received");
389		}
390
391		// Empty WANT list should not cause an error
392		let (tx, rx) = oneshot::channel();
393		config
394			.inbound_queue
395			.unwrap()
396			.send(IncomingRequest {
397				peer: PeerId::random(),
398				payload: BitswapMessage {
399					wantlist: Some(Default::default()),
400					..Default::default()
401				}
402				.encode_to_vec(),
403				pending_response: tx,
404			})
405			.await
406			.unwrap();
407
408		if let Ok(OutgoingResponse { result, reputation_changes, sent_feedback }) = rx.await {
409			assert_eq!(result, Ok(BitswapMessage::default().encode_to_vec()));
410			assert_eq!(reputation_changes, Vec::new());
411			assert!(sent_feedback.is_none());
412		} else {
413			panic!("invalid event received");
414		}
415	}
416
417	#[tokio::test]
418	async fn too_long_want_list() {
419		let client = substrate_test_runtime_client::new();
420		let (bitswap, config) = BitswapRequestHandler::new(Arc::new(client));
421
422		tokio::spawn(async move { bitswap.run().await });
423
424		let (tx, rx) = oneshot::channel();
425		config
426			.inbound_queue
427			.unwrap()
428			.send(IncomingRequest {
429				peer: PeerId::random(),
430				payload: BitswapMessage {
431					wantlist: Some(Wantlist {
432						entries: (0..MAX_WANTED_BLOCKS + 1)
433							.map(|_| Entry::default())
434							.collect::<Vec<_>>(),
435						full: false,
436					}),
437					..Default::default()
438				}
439				.encode_to_vec(),
440				pending_response: tx,
441			})
442			.await
443			.unwrap();
444
445		if let Ok(OutgoingResponse { result, reputation_changes, sent_feedback }) = rx.await {
446			assert_eq!(result, Err(()));
447			assert_eq!(reputation_changes, Vec::new());
448			assert!(sent_feedback.is_none());
449		} else {
450			panic!("invalid event received");
451		}
452	}
453
454	#[tokio::test]
455	async fn transaction_not_found() {
456		let client = TestClientBuilder::with_tx_storage(u32::MAX).build();
457
458		let (bitswap, config) = BitswapRequestHandler::new(Arc::new(client));
459		tokio::spawn(async move { bitswap.run().await });
460
461		let (tx, rx) = oneshot::channel();
462		config
463			.inbound_queue
464			.unwrap()
465			.send(IncomingRequest {
466				peer: PeerId::random(),
467				payload: BitswapMessage {
468					wantlist: Some(Wantlist {
469						entries: vec![Entry {
470							block: cid::Cid::new_v1(
471								0x70,
472								cid::multihash::Multihash::wrap(
473									u64::from(LiteP2pCode::Blake2b256),
474									&[0u8; 32],
475								)
476								.unwrap(),
477							)
478							.to_bytes(),
479							..Default::default()
480						}],
481						full: false,
482					}),
483					..Default::default()
484				}
485				.encode_to_vec(),
486				pending_response: tx,
487			})
488			.await
489			.unwrap();
490
491		if let Ok(OutgoingResponse { result, reputation_changes, sent_feedback }) = rx.await {
492			assert_eq!(result, Ok(vec![]));
493			assert_eq!(reputation_changes, Vec::new());
494			assert!(sent_feedback.is_none());
495		} else {
496			panic!("invalid event received");
497		}
498	}
499
500	#[tokio::test]
501	async fn transaction_found() {
502		let client = TestClientBuilder::with_tx_storage(u32::MAX).build();
503		let mut block_builder = BlockBuilderBuilder::new(&client)
504			.on_parent_block(client.chain_info().genesis_hash)
505			.with_parent_block_number(0)
506			.build()
507			.unwrap();
508
509		// encoded extrinsic: [161, .. , 2, 6, 16, 19, 55, 19, 56]
510		let ext = ExtrinsicBuilder::new_indexed_call(vec![0x13, 0x37, 0x13, 0x38]).build();
511		let pattern_index = ext.encoded_size() - 4;
512
513		block_builder.push(ext.clone()).unwrap();
514		let block = block_builder.build().unwrap().block;
515
516		client.import(BlockOrigin::File, block).await.unwrap();
517
518		let (bitswap, config) = BitswapRequestHandler::new(Arc::new(client));
519
520		tokio::spawn(async move { bitswap.run().await });
521
522		let (tx, rx) = oneshot::channel();
523		config
524			.inbound_queue
525			.unwrap()
526			.send(IncomingRequest {
527				peer: PeerId::random(),
528				payload: BitswapMessage {
529					wantlist: Some(Wantlist {
530						entries: vec![Entry {
531							block: cid::Cid::new_v1(
532								0x70,
533								cid::multihash::Multihash::wrap(
534									u64::from(LiteP2pCode::Blake2b256),
535									&sp_crypto_hashing::blake2_256(&ext.encode()[pattern_index..]),
536								)
537								.unwrap(),
538							)
539							.to_bytes(),
540							..Default::default()
541						}],
542						full: false,
543					}),
544					..Default::default()
545				}
546				.encode_to_vec(),
547				pending_response: tx,
548			})
549			.await
550			.unwrap();
551
552		if let Ok(OutgoingResponse { result, reputation_changes, sent_feedback }) = rx.await {
553			assert_eq!(reputation_changes, Vec::new());
554			assert!(sent_feedback.is_none());
555
556			let response =
557				schema::bitswap::Message::decode(&result.expect("fetch to succeed")[..]).unwrap();
558			assert_eq!(response.payload[0].data, vec![0x13, 0x37, 0x13, 0x38]);
559		} else {
560			panic!("invalid event received");
561		}
562	}
563
564	#[tokio::test]
565	async fn transaction_not_found_sends_dont_have_when_requested() {
566		let client = TestClientBuilder::with_tx_storage(u32::MAX).build();
567		let (mut bitswap, _config) = BitswapRequestHandler::new(Arc::new(client));
568		let cid = cid::Cid::new_v1(
569			0x70,
570			cid::multihash::Multihash::wrap(u64::from(LiteP2pCode::Blake2b256), &[0u8; 32])
571				.unwrap(),
572		);
573		let request = BitswapMessage {
574			wantlist: Some(Wantlist {
575				entries: vec![Entry {
576					block: cid.to_bytes(),
577					send_dont_have: true,
578					..Default::default()
579				}],
580				full: false,
581			}),
582			..Default::default()
583		}
584		.encode_to_vec();
585
586		let response = BitswapMessage::decode(
587			bitswap.handle_message(&PeerId::random(), &request).unwrap().as_slice(),
588		)
589		.unwrap();
590
591		assert!(response.payload.is_empty());
592		assert_eq!(response.block_presences.len(), 1);
593		assert_eq!(response.block_presences[0].cid, cid.to_bytes());
594		assert_eq!(response.block_presences[0].r#type, BlockPresenceType::DontHave as i32);
595	}
596
597	#[tokio::test]
598	async fn transaction_found_sends_have_for_want_have() {
599		let client = TestClientBuilder::with_tx_storage(u32::MAX).build();
600		let mut block_builder = BlockBuilderBuilder::new(&client)
601			.on_parent_block(client.chain_info().genesis_hash)
602			.with_parent_block_number(0)
603			.build()
604			.unwrap();
605
606		let ext = ExtrinsicBuilder::new_indexed_call(vec![0x13, 0x37, 0x13, 0x38]).build();
607		let pattern_index = ext.encoded_size() - 4;
608		let cid = cid::Cid::new_v1(
609			0x70,
610			cid::multihash::Multihash::wrap(
611				u64::from(LiteP2pCode::Blake2b256),
612				&sp_crypto_hashing::blake2_256(&ext.encode()[pattern_index..]),
613			)
614			.unwrap(),
615		);
616
617		block_builder.push(ext).unwrap();
618		let block = block_builder.build().unwrap().block;
619		client.import(BlockOrigin::File, block).await.unwrap();
620
621		let (mut bitswap, _config) = BitswapRequestHandler::new(Arc::new(client));
622		let request = BitswapMessage {
623			wantlist: Some(Wantlist {
624				entries: vec![Entry {
625					block: cid.to_bytes(),
626					want_type: WantType::Have as i32,
627					..Default::default()
628				}],
629				full: false,
630			}),
631			..Default::default()
632		}
633		.encode_to_vec();
634
635		let response = BitswapMessage::decode(
636			bitswap.handle_message(&PeerId::random(), &request).unwrap().as_slice(),
637		)
638		.unwrap();
639
640		assert!(response.payload.is_empty());
641		assert_eq!(response.block_presences.len(), 1);
642		assert_eq!(response.block_presences[0].cid, cid.to_bytes());
643		assert_eq!(response.block_presences[0].r#type, BlockPresenceType::Have as i32);
644	}
645
646	#[test]
647	fn is_cid_supported_accepts_all_three_supported_hashings() {
648		use cid::multihash::Multihash;
649		for multihash_code in
650			[BLAKE2B_256_MULTIHASH_CODE, SHA2_256_MULTIHASH_CODE, KECCAK_256_MULTIHASH_CODE]
651		{
652			let digest = [9u8; 32];
653			let mh = Multihash::<64>::wrap(multihash_code, &digest).unwrap();
654			let cid = Cid::new_v1(RAW_CODEC, mh);
655			assert!(is_cid_supported(&cid), "{multihash_code} CID should be supported");
656		}
657	}
658
659	#[test]
660	fn is_cid_supported_rejects_unknown_multihash_code() {
661		use cid::multihash::Multihash;
662		let digest = [9u8; 32];
663		let mh = Multihash::<64>::wrap(0x99, &digest).unwrap();
664		let cid = Cid::new_v1(RAW_CODEC, mh);
665		assert!(!is_cid_supported(&cid));
666	}
667}