use serde::{Deserialize, Serialize, de::DeserializeOwned};
use std::{error::Error, fmt};
mod io;
mod receiver;
mod sender;
pub use receiver::{ErasedReceiver, PortDeserializer, Receiver, RecvError};
pub use sender::{Closed, ErasedSender, PortSerializer, SendError, SendErrorKind, Sender};
use crate::{RemoteSend, chmux, codec};
const BIG_DATA_CHUNK_QUEUE: usize = 32;
const BIG_DATA_LIMIT: i8 = 16;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub enum ConnectError {
Connect(chmux::ConnectError),
Listen(chmux::ListenerError),
NoConnectRequest,
}
impl fmt::Display for ConnectError {
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
match self {
ConnectError::Connect(err) => write!(f, "connect error: {err}"),
ConnectError::Listen(err) => write!(f, "listen error: {err}"),
ConnectError::NoConnectRequest => write!(f, "no connect request received"),
}
}
}
impl Error for ConnectError {}
impl From<chmux::ConnectError> for ConnectError {
fn from(err: chmux::ConnectError) -> Self {
Self::Connect(err)
}
}
impl From<chmux::ListenerError> for ConnectError {
fn from(err: chmux::ListenerError) -> Self {
Self::Listen(err)
}
}
pub async fn connect<Tx, Rx, Codec>(
client: &chmux::Client, listener: &mut chmux::Listener,
) -> Result<(Sender<Tx, Codec>, Receiver<Rx, Codec>), ConnectError>
where
Tx: RemoteSend,
Rx: RemoteSend,
Codec: codec::Codec,
{
async fn connect_raw(
client: &chmux::Client, listener: &mut chmux::Listener,
) -> Result<(chmux::Sender, chmux::Receiver), ConnectError> {
let (client_sr, listener_sr) = tokio::join!(client.connect(), listener.accept());
let (raw_sender, _) = client_sr?;
let (_, raw_receiver) = listener_sr?.ok_or(ConnectError::NoConnectRequest)?;
Ok((raw_sender, raw_receiver))
}
let (raw_sender, raw_receiver) = connect_raw(client, listener).await?;
Ok((Sender::new(raw_sender), Receiver::new(raw_receiver)))
}
pub trait BaseExt<T, Codec> {
fn with_max_item_size(self, max_item_size: usize) -> (Sender<T, Codec>, Receiver<T, Codec>);
}
impl<T, Codec> BaseExt<T, Codec> for (Sender<T, Codec>, Receiver<T, Codec>)
where
T: Serialize + DeserializeOwned + Send + 'static,
Codec: codec::Codec,
{
fn with_max_item_size(self, max_item_size: usize) -> (Sender<T, Codec>, Receiver<T, Codec>) {
let (mut tx, mut rx) = self;
tx.set_max_item_size(max_item_size);
rx.set_max_item_size(max_item_size);
(tx, rx)
}
}