1use 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
44mod client;
46pub 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
60const MAX_PACKET_SIZE: u64 = MAX_RESPONSE_SIZE;
62
63const MAX_REQUEST_QUEUE: usize = 20;
65
66pub const MAX_WANTED_BLOCKS: usize = 16;
68
69pub(crate) const PROTOCOL_NAME: &str = "/ipfs/bitswap/1.2.0";
71
72pub const RAW_CODEC: u64 = 0x55;
74
75pub 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
83pub(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#[derive(PartialEq, Eq, Clone, Debug)]
90pub struct Prefix {
91 pub version: CidVersion,
93 pub codec: u64,
95 pub mh_type: u64,
97 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 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
132pub(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 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 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 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 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#[derive(Debug, thiserror::Error)]
286enum RequestHandlerError {
287 #[error("Failed to decode request: {0}.")]
289 DecodeProto(#[from] prost::DecodeError),
290
291 #[error("Failed to encode response: {0}.")]
293 EncodeProto(#[from] prost::EncodeError),
294
295 #[error(transparent)]
297 Client(#[from] sp_blockchain::Error),
298
299 #[error(transparent)]
301 BadCid(#[from] CidError),
302
303 #[error(transparent)]
305 Read(#[from] io::Error),
306
307 #[error("Failed to send response.")]
309 SendResponse,
310
311 #[error("Invalid WANT list.")]
313 InvalidWantList,
314
315 #[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 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 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}