use std::{any, cell::UnsafeCell, io, sync::Arc, task::Poll};
use ntex_io::{Filter, FilterBuf, FilterLayer, Io, Layer};
use tls_rustls::{ClientConfig, ClientConnection, pki_types::ServerName};
use super::stream::Stream;
#[derive(Debug)]
pub struct TlsClientFilter {
session: UnsafeCell<ClientConnection>,
}
impl FilterLayer for TlsClientFilter {
fn query(&self, id: any::TypeId) -> Option<Box<dyn any::Any>> {
self.stream(|s| s.query(id))
}
fn process_read_buf(&self, buf: &FilterBuf<'_>) -> io::Result<()> {
self.stream(|s| s.process_read_buf(buf))
}
fn process_write_buf(&self, buf: &FilterBuf<'_>) -> io::Result<()> {
self.stream(|s| s.process_write_buf(buf))
}
fn shutdown(&self, buf: &FilterBuf<'_>) -> io::Result<Poll<()>> {
self.stream(|s| s.shutdown(buf))
}
}
impl TlsClientFilter {
pub async fn create<F: Filter>(
io: Io<F>,
cfg: Arc<ClientConfig>,
domain: ServerName<'static>,
) -> Result<Io<Layer<TlsClientFilter, F>>, io::Error> {
let mut session = ClientConnection::new(cfg, domain).map_err(io::Error::other)?;
session.set_buffer_limit(Some(io.cfg().write_page_size().capacity()));
let io = io.add_filter(TlsClientFilter {
session: UnsafeCell::new(session),
});
loop {
let (wants_write, handshaking) = {
let s = unsafe { &*io.filter().session.get() };
(s.wants_write(), s.is_handshaking())
};
if wants_write {
io.flush(false).await?;
}
if handshaking {
io.read_notify()
.await?
.ok_or_else(|| io::Error::new(io::ErrorKind::NotConnected, "disconnected"))?;
} else {
return Ok(io);
}
}
}
fn stream<F, R>(&self, f: F) -> R
where
F: FnOnce(&mut Stream<'_, ClientConnection>) -> R,
{
let mut s = Stream::new(unsafe { &mut *self.session.get() });
f(&mut s)
}
}