use crate::error::{RequestError, TaskError};
use crate::handler::{
Configurator, Receiver, Sender, StreamReceiver, StreamSender, TaskHandle,
};
use crate::packet::{Packet, PlainBytes};
use crate::plain;
use crate::server::Request;
use crate::util::{ByteStream, PinnedFuture};
#[cfg(feature = "encrypted")]
use crate::{encrypted, packet::EncryptedBytes};
#[cfg(feature = "encrypted")]
use crypto::signature as sign;
use std::io;
use std::time::Duration;
#[derive(Debug, Clone)]
pub struct Config {
pub timeout: Duration,
pub body_limit: u32,
}
pub struct ReconStrat<S> {
pub(crate) inner:
Box<dyn FnMut(usize) -> PinnedFuture<'static, io::Result<S>> + Send>,
}
impl<S> ReconStrat<S> {
pub fn new<F: 'static>(f: F) -> Self
where
F: FnMut(usize) -> PinnedFuture<'static, io::Result<S>> + Send,
{
Self { inner: Box::new(f) }
}
}
#[derive(Debug)]
pub struct Connection<P> {
sender: Sender<P>,
receiver: Option<Receiver<P>>,
receiver_enabled: bool,
config: Configurator<Config>,
task: TaskHandle,
}
impl<P> Connection<P> {
pub fn new<S>(
byte_stream: S,
cfg: Config,
recon_strat: Option<ReconStrat<S>>,
) -> Self
where
S: ByteStream,
P: Packet<PlainBytes> + Send + 'static,
P::Header: Send,
{
plain::client(byte_stream, cfg, recon_strat)
}
#[cfg(feature = "encrypted")]
#[cfg_attr(docsrs, doc(cfg(feature = "encrypted")))]
pub fn new_encrypted<S>(
byte_stream: S,
cfg: Config,
recon_strat: Option<ReconStrat<S>>,
sign: sign::PublicKey,
) -> Self
where
S: ByteStream,
P: Packet<EncryptedBytes> + Send + 'static,
P::Header: Send,
{
encrypted::client(byte_stream, cfg, recon_strat, sign)
}
pub(crate) fn new_raw(
sender: Sender<P>,
receiver: Receiver<P>,
config: Configurator<Config>,
task: TaskHandle,
) -> Self {
Self {
sender,
receiver: Some(receiver),
receiver_enabled: false,
config,
task,
}
}
pub fn update_config(&self, cfg: Config) {
self.config.update(cfg);
}
pub fn configurator(&self) -> Configurator<Config> {
self.config.clone()
}
pub async fn enable_server_requests(&mut self) -> Result<(), RequestError> {
self.sender.enable_server_requests().await?;
self.receiver_enabled = true;
Ok(())
}
pub fn is_server_requests_enabled(&self) -> bool {
self.receiver_enabled
}
pub fn take_receiver(&mut self) -> Option<Receiver<P>> {
assert!(self.receiver_enabled);
self.receiver.take()
}
pub async fn receive(&mut self) -> Option<Request<P>> {
assert!(self.receiver_enabled);
self.receiver.as_mut().unwrap().receive().await
}
pub fn clone_sender(&self) -> Sender<P> {
self.sender.clone()
}
pub async fn request(&self, packet: P) -> Result<P, RequestError> {
self.sender.request(packet).await
}
pub async fn request_sender(
&self,
packet: P,
) -> Result<StreamSender<P>, RequestError> {
self.sender.request_sender(packet).await
}
pub async fn request_receiver(
&self,
packet: P,
) -> Result<StreamReceiver<P>, RequestError> {
self.sender.request_receiver(packet).await
}
pub async fn closed(&mut self) {
self.task.closed().await
}
pub async fn wait(self) -> Result<(), TaskError> {
self.task.wait().await
}
pub async fn close(self) -> Result<(), TaskError> {
self.task.close().await
}
}