Skip to main content

sc_network_bitswap/
handle.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 request API.
19
20use super::{is_cid_supported, Cid};
21
22use bytes::Bytes;
23use std::collections::HashSet;
24use tokio::sync::mpsc;
25
26/// Bitswap request errors.
27#[derive(Debug, thiserror::Error)]
28pub enum BitswapError {
29	/// The service is unavailable.
30	#[error("Bitswap service is closed")]
31	ServiceClosed,
32	/// A CID is unsupported.
33	#[error("invalid CID for Bitswap: {cid}")]
34	InvalidCid {
35		/// Unsupported CID.
36		cid: Cid,
37	},
38	/// The service is at capacity.
39	#[error("Bitswap service is overloaded")]
40	Overloaded,
41}
42
43/// A fetched block or request error. Block bytes are [`Bytes`] because a block may be
44/// delivered to multiple concurrent requests without copying.
45pub type FetchItem = Result<(Cid, Bytes), BitswapError>;
46
47/// Handle for submitting Bitswap requests.
48#[derive(Debug, Clone)]
49pub struct BitswapHandle {
50	cmd_tx: mpsc::Sender<BitswapCommand>,
51}
52
53impl BitswapHandle {
54	pub(crate) fn new(cmd_tx: mpsc::Sender<BitswapCommand>) -> Self {
55		Self { cmd_tx }
56	}
57
58	/// Submit a wantlist. Returns a receiver that yields `Ok((cid, bytes))` with
59	/// hash-verified bytes for each requested CID, in the order they resolve.
60	///
61	/// The stream closes once every CID has been delivered. Unresolved CIDs are retried
62	/// until the receiver is dropped.
63	///
64	/// `Err(BitswapError::ServiceClosed)` is yielded once, as the final item, if the
65	/// service shuts down mid-request.
66	pub fn request_stream(
67		&self,
68		cids: HashSet<Cid>,
69	) -> Result<mpsc::Receiver<FetchItem>, BitswapError> {
70		if cids.is_empty() {
71			let (_tx, rx) = mpsc::channel(1);
72			return Ok(rx);
73		}
74
75		for cid in &cids {
76			if !is_cid_supported(cid) {
77				return Err(BitswapError::InvalidCid { cid: *cid });
78			}
79		}
80
81		let (sink, rx) = mpsc::channel(cids.len() + 1);
82
83		self.cmd_tx.try_send(BitswapCommand::RequestStream { cids, sink }).map_err(
84			|e| match e {
85				mpsc::error::TrySendError::Full(_) => BitswapError::Overloaded,
86				mpsc::error::TrySendError::Closed(_) => BitswapError::ServiceClosed,
87			},
88		)?;
89
90		Ok(rx)
91	}
92}
93
94#[derive(Debug)]
95pub(crate) enum BitswapCommand {
96	RequestStream { cids: HashSet<Cid>, sink: mpsc::Sender<FetchItem> },
97}