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}