summaryrefslogtreecommitdiff
path: root/src/transport
diff options
context:
space:
mode:
Diffstat (limited to 'src/transport')
-rw-r--r--src/transport/receiver.rs13
-rw-r--r--src/transport/sender.rs12
-rw-r--r--src/transport/tcp.rs32
-rw-r--r--src/transport/udp.rs26
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);