use super::{StreamReceiver, StreamSender};
use crate::error::RequestError;
use crate::handler::handler::InternalRequest;
use tokio::sync::{mpsc, oneshot};
#[derive(Debug)]
pub struct Sender<P> {
pub(super) inner: mpsc::Sender<InternalRequest<P>>,
}
impl<P> Sender<P> {
pub(crate) async fn enable_server_requests(
&self,
) -> Result<(), RequestError> {
self.inner
.send(InternalRequest::EnableServerRequests)
.await
.map_err(|_| RequestError::ConnectionAlreadyClosed)
}
pub async fn request(&self, packet: P) -> Result<P, RequestError> {
let (tx, rx) = oneshot::channel();
self.inner
.send(InternalRequest::Request(packet, tx))
.await
.map_err(|_| RequestError::ConnectionAlreadyClosed)?;
rx.await.map_err(|_| RequestError::TaskFailed)?
}
pub async fn request_sender(
&self,
packet: P,
) -> Result<StreamSender<P>, RequestError> {
let (tx, rx) = mpsc::channel(10);
self.inner
.send(InternalRequest::RequestSender(packet, rx))
.await
.map_err(|_| RequestError::ConnectionAlreadyClosed)?;
Ok(StreamSender::new(tx))
}
pub async fn request_receiver(
&self,
packet: P,
) -> Result<StreamReceiver<P>, RequestError> {
let (tx, rx) = mpsc::channel(10);
self.inner
.send(InternalRequest::RequestReceiver(packet, tx))
.await
.map_err(|_| RequestError::ConnectionAlreadyClosed)?;
Ok(StreamReceiver::new(rx))
}
pub fn is_closed(&self) -> bool {
self.inner.is_closed()
}
pub async fn closed(&self) {
self.inner.closed().await
}
}
impl<P> Clone for Sender<P> {
fn clone(&self) -> Self {
Self {
inner: self.inner.clone(),
}
}
}