diff options
Diffstat (limited to 'src/transport')
| -rw-r--r-- | src/transport/receiver.rs | 13 | ||||
| -rw-r--r-- | src/transport/sender.rs | 12 | ||||
| -rw-r--r-- | src/transport/tcp.rs | 32 | ||||
| -rw-r--r-- | src/transport/udp.rs | 26 |
4 files changed, 53 insertions, 30 deletions
diff --git a/src/transport/receiver.rs b/src/transport/receiver.rs index 3a7c948..a165110 100644 --- a/src/transport/receiver.rs +++ b/src/transport/receiver.rs @@ -3,3 +3,16 @@ use crate::{codec::Codec, context::Context, transport::Result}; pub trait PacketReceiver<Data, Uid: PartialEq>: Send + Sync { fn recv<P: Codec<Data>>(&self, ctx: &Context<Data>) -> Result<(Uid, P)>; } + +pub trait PacketReceiverBuf<Data, Uid: PartialEq>: Send + Sync { + fn recv_buf(&self, ctx: &Context<Data>) -> Result<(Uid, Vec<u8>)>; +} + +impl<Data, Uid: PartialEq, R: PacketReceiverBuf<Data, Uid>> PacketReceiver<Data, Uid> for R { + fn recv<P: Codec<Data>>(&self, ctx: &Context<Data>) -> Result<(Uid, P)> { + let (uid, bytes) = self.recv_buf(ctx)?; + let mut reader = bytes.as_slice(); + let packet = P::decode(&mut reader, ctx)?; + Ok((uid, packet)) + } +}
\ No newline at end of file diff --git a/src/transport/sender.rs b/src/transport/sender.rs index 804914f..3c83c6e 100644 --- a/src/transport/sender.rs +++ b/src/transport/sender.rs @@ -3,3 +3,15 @@ use crate::{codec::Codec, context::Context, transport::Result}; pub trait PacketSender<Data>: Send + Sync { fn send<P: Codec<Data>>(&self, packet: P, ctx: &Context<Data>) -> Result<()>; } + +pub trait PacketSenderBuf<Data>: Send + Sync { + fn send_buf(&self, buf: &[u8]) -> Result<()>; +} + +impl<Data, S: PacketSenderBuf<Data>> PacketSender<Data> for S { + fn send<P: Codec<Data>>(&self, packet: P, ctx: &Context<Data>) -> Result<()> { + let mut buf = Vec::new(); + packet.encode(&mut buf, ctx)?; + self.send_buf(&buf) + } +}
\ No newline at end of file diff --git a/src/transport/tcp.rs b/src/transport/tcp.rs index 19d02d6..12ecaa0 100644 --- a/src/transport/tcp.rs +++ b/src/transport/tcp.rs @@ -1,35 +1,30 @@ use std::hash::Hash; -use std::io::Write; +use std::io::{Read, Write}; use std::marker::PhantomData; use std::net::{SocketAddr, TcpListener as StdTcpListener, TcpStream}; use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::{Arc, Mutex}; -use crate::codec::Codec; use crate::context::Context; use crate::transport::connection::Connection; use crate::transport::error::TransportError; use crate::transport::listener::Listener; -use crate::transport::receiver::PacketReceiver; -use crate::transport::sender::PacketSender; +use crate::transport::receiver::PacketReceiverBuf; +use crate::transport::sender::PacketSenderBuf; use crate::transport::{Result, Transport}; -// TODO : check Tcp Transport - #[derive(Clone)] pub struct TcpSender { stream: Arc<Mutex<TcpStream>>, } -impl<Data> PacketSender<Data> for TcpSender { - fn send<P: Codec<Data>>(&self, packet: P, ctx: &Context<Data>) -> Result<()> { - let mut buf = Vec::new(); - packet.encode(&mut buf, ctx)?; +impl<Data> PacketSenderBuf<Data> for TcpSender { + fn send_buf(&self, buf: &[u8]) -> Result<()> { let mut guard = self .stream .lock() .map_err(|_| TransportError::LockPoisoned)?; - guard.write_all(&buf)?; + guard.write_all(buf)?; guard.flush()?; Ok(()) } @@ -40,17 +35,24 @@ pub struct TcpReceiver<Uid> { uid: Uid, } -impl<Data, Uid> PacketReceiver<Data, Uid> for TcpReceiver<Uid> +impl<Data, Uid> PacketReceiverBuf<Data, Uid> for TcpReceiver<Uid> where Uid: PartialEq + Clone + Send + Sync, { - fn recv<P: Codec<Data>>(&self, ctx: &Context<Data>) -> Result<(Uid, P)> { + fn recv_buf(&self, _ctx: &Context<Data>) -> Result<(Uid, Vec<u8>)> { let mut guard = self .stream .lock() .map_err(|_| TransportError::LockPoisoned)?; - let packet = P::decode(&mut *guard, ctx)?; - Ok((self.uid.clone(), packet)) + + // Read 4-byte big-endian length prefix + let mut len_buf = [0u8; 4]; + guard.read_exact(&mut len_buf)?; + let len = u32::from_be_bytes(len_buf) as usize; + + let mut buf = vec![0u8; len]; + guard.read_exact(&mut buf)?; + Ok((self.uid.clone(), buf)) } } diff --git a/src/transport/udp.rs b/src/transport/udp.rs index acb36f2..9b312cf 100644 --- a/src/transport/udp.rs +++ b/src/transport/udp.rs @@ -10,11 +10,11 @@ use crate::context::Context; use crate::transport::Result; use crate::transport::connection::Connection; use crate::transport::listener::Listener; -use crate::transport::receiver::PacketReceiver; -use crate::transport::sender::PacketSender; +use crate::transport::receiver::PacketReceiverBuf; +use crate::transport::sender::PacketSenderBuf; use crate::transport::{Transport, TransportError}; -// TODO : check Udp Transport +const UDP_BUF_SIZE: usize = 1500; #[derive(Clone)] pub struct UdpPeerSender { @@ -22,11 +22,9 @@ pub struct UdpPeerSender { peer: SocketAddr, } -impl<Data> PacketSender<Data> for UdpPeerSender { - fn send<P: crate::codec::Codec<Data>>(&self, packet: P, ctx: &Context<Data>) -> Result<()> { - let mut buf = Vec::new(); - packet.encode(&mut buf, ctx)?; - self.sock.send_to(&buf, self.peer)?; +impl<Data> PacketSenderBuf<Data> for UdpPeerSender { + fn send_buf(&self, buf: &[u8]) -> Result<()> { + self.sock.send_to(buf, self.peer)?; Ok(()) } } @@ -57,11 +55,11 @@ impl<Uid> UdpReceiver<Uid> { } } -impl<Data, Uid> PacketReceiver<Data, Uid> for UdpReceiver<Uid> +impl<Data, Uid> PacketReceiverBuf<Data, Uid> for UdpReceiver<Uid> where Uid: PartialEq + Clone + Send + Sync, { - fn recv<P: crate::codec::Codec<Data>>(&self, ctx: &Context<Data>) -> Result<(Uid, P)> { + fn recv_buf(&self, _ctx: &Context<Data>) -> Result<(Uid, Vec<u8>)> { let bytes = match &self.inner { UdpReceiverInner::Channel(rx) => rx .lock() @@ -69,15 +67,13 @@ where .recv() .map_err(|_| TransportError::ChannelClosed)?, UdpReceiverInner::Socket(sock) => { - let mut buf = vec![0u8; 65535]; + let mut buf = vec![0u8; UDP_BUF_SIZE]; let (n, _) = sock.recv_from(&mut buf)?; buf.truncate(n); buf } }; - let mut reader = bytes.as_slice(); - let packet = P::decode(&mut reader, ctx)?; - Ok((self.uid.clone(), packet)) + Ok((self.uid.clone(), bytes)) } } @@ -100,7 +96,7 @@ where fn accept(&mut self) -> Result<Connection<Uid, Data, UdpPeerSender, UdpReceiver<Uid>>> { loop { - let mut buf = vec![0u8; 65535]; + let mut buf = vec![0u8; UDP_BUF_SIZE]; let (n, src) = self.sock.recv_from(&mut buf)?; buf.truncate(n); |
