use crate::{
RuntimeError,
futures::{
net::{
exchange::Sends,
socket::{get_option, set_keepalive, set_option},
stream::{FinishTask, Pipe, RecvTask, SendTask, Source},
},
tcp::tcp_task::AcceptTask,
},
modules::fd::Fd,
};
use std::{fmt, net::SocketAddr, sync::Arc, time::Duration};
struct Stream {
pipe: Pipe,
local: SocketAddr,
peer: SocketAddr,
}
#[derive(Clone)]
pub struct Connection {
stream: Arc<Stream>,
}
impl Connection {
pub(crate) fn new(fd: Fd, local: SocketAddr, peer: SocketAddr) -> Self {
Self {
stream: Arc::new(Stream {
pipe: Pipe::new(fd),
local,
peer,
}),
}
}
#[inline(always)]
pub(crate) fn pipe(&self) -> &Pipe {
&self.stream.pipe
}
#[inline(always)]
fn source(&self) -> Source {
Source::Tcp(self.clone())
}
pub fn send(&self, data: impl Into<Arc<[u8]>>) -> SendTask {
SendTask::new(self.source(), data.into())
}
pub fn recv(&self, max: usize) -> RecvTask {
RecvTask::some(self.source(), max)
}
pub fn recv_exact(&self, len: usize) -> RecvTask {
RecvTask::exact(self.source(), len)
}
pub fn recv_until(&self, delimiter: &[u8], max: usize) -> RecvTask {
RecvTask::until(self.source(), Arc::from(delimiter), max)
}
pub fn recv_to_end(&self) -> RecvTask {
RecvTask::to_end(self.source())
}
#[inline(always)]
pub fn local_addr(&self) -> SocketAddr {
self.stream.local
}
#[inline(always)]
pub fn peer_addr(&self) -> SocketAddr {
self.stream.peer
}
pub fn set_nodelay(&self, nodelay: bool) -> Result<(), RuntimeError> {
set_option(
self.pipe().fd(),
libc::IPPROTO_TCP,
libc::TCP_NODELAY,
nodelay as libc::c_int,
)
}
pub fn nodelay(&self) -> Result<bool, RuntimeError> {
Ok(get_option(self.pipe().fd(), libc::IPPROTO_TCP, libc::TCP_NODELAY)? != 0)
}
pub fn set_keepalive(&self, idle: Option<Duration>) -> Result<(), RuntimeError> {
set_keepalive(self.pipe().fd(), idle)
}
pub fn finish(&self) -> FinishTask {
FinishTask::new(self.source())
}
pub fn close(self) {
drop(self);
}
}
impl fmt::Debug for Connection {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("Connection")
.field("local", &self.stream.local)
.field("peer", &self.stream.peer)
.finish()
}
}
struct Listening {
fd: Fd,
local: SocketAddr,
}
#[derive(Clone)]
pub struct Listener {
socket: Arc<Listening>,
}
impl Listener {
pub(crate) fn new(fd: Fd, local: SocketAddr) -> Self {
Self {
socket: Arc::new(Listening { fd, local }),
}
}
#[inline(always)]
pub(crate) fn fd(&self) -> libc::c_int {
self.socket.fd.raw()
}
pub fn accept(&self) -> AcceptTask {
AcceptTask::new(self.clone())
}
#[inline(always)]
pub fn local_addr(&self) -> SocketAddr {
self.socket.local
}
pub fn close(self) {
drop(self);
}
}
impl fmt::Debug for Listener {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("Listener")
.field("local", &self.socket.local)
.finish()
}
}
impl Sends for Connection {
fn send_all(&self, data: Arc<[u8]>) -> SendTask {
self.send(data)
}
}