use bytes::Bytes;
use tokio::sync::mpsc;
use tokio_util::sync::CancellationToken;
pub(crate) struct IdDataSenderBridge {
pub id: Bytes,
pub inner: DataSenderBridge,
}
pub(crate) struct DataSenderBridge {
chan: DataSender,
shutdown: CancellationToken,
}
impl DataSenderBridge {
pub(crate) fn new(chan: DataSender, shutdown: CancellationToken) -> Self {
Self { chan, shutdown }
}
pub(crate) async fn send_data(
&self,
data: Vec<u8>,
) -> Result<(), mpsc::error::SendError<BridgeData>> {
self.chan.send(BridgeData::Data(data)).await
}
pub(crate) async fn send_sender(
&self,
sender: mpsc::Sender<Vec<u8>>,
) -> Result<(), mpsc::error::SendError<BridgeData>> {
self.chan.send(BridgeData::Sender(sender)).await
}
pub(crate) fn close(&self) {
self.shutdown.cancel();
}
}
pub(crate) type DataSender = mpsc::Sender<BridgeData>;
pub(crate) enum BridgeData {
Sender(mpsc::Sender<Vec<u8>>),
Data(Vec<u8>),
}