use super::{is_cid_supported, Cid};
use bytes::Bytes;
use std::collections::HashSet;
use tokio::sync::mpsc;
#[derive(Debug, thiserror::Error)]
pub enum BitswapError {
#[error("Bitswap service is closed")]
ServiceClosed,
#[error("invalid CID for Bitswap: {cid}")]
InvalidCid {
cid: Cid,
},
#[error("Bitswap service is overloaded")]
Overloaded,
}
pub type FetchItem = Result<(Cid, Bytes), BitswapError>;
#[derive(Debug, Clone)]
pub struct BitswapHandle {
cmd_tx: mpsc::Sender<BitswapCommand>,
}
impl BitswapHandle {
pub(crate) fn new(cmd_tx: mpsc::Sender<BitswapCommand>) -> Self {
Self { cmd_tx }
}
pub fn request_stream(
&self,
cids: HashSet<Cid>,
) -> Result<mpsc::Receiver<FetchItem>, BitswapError> {
if cids.is_empty() {
let (_tx, rx) = mpsc::channel(1);
return Ok(rx);
}
for cid in &cids {
if !is_cid_supported(cid) {
return Err(BitswapError::InvalidCid { cid: *cid });
}
}
let (sink, rx) = mpsc::channel(cids.len() + 1);
self.cmd_tx.try_send(BitswapCommand::RequestStream { cids, sink }).map_err(
|e| match e {
mpsc::error::TrySendError::Full(_) => BitswapError::Overloaded,
mpsc::error::TrySendError::Closed(_) => BitswapError::ServiceClosed,
},
)?;
Ok(rx)
}
}
#[derive(Debug)]
pub(crate) enum BitswapCommand {
RequestStream { cids: HashSet<Cid>, sink: mpsc::Sender<FetchItem> },
}