diff --git a/easytier/src/gateway/socks5.rs b/easytier/src/gateway/socks5.rs index 82ead977..2614f046 100644 --- a/easytier/src/gateway/socks5.rs +++ b/easytier/src/gateway/socks5.rs @@ -363,7 +363,10 @@ impl Socks5ServerNet { let mut smoltcp_stack_receiver = packet_recv.lock().await; while let Some(packet) = smoltcp_stack_receiver.recv().await { tracing::trace!(?packet, "receive from peer send to smoltcp packet"); - if let Err(e) = stack_sink.send(Ok(packet.payload().to_vec())).await { + if let Err(e) = stack_sink + .send(Ok(bytes::BytesMut::from(packet.payload()))) + .await + { tracing::error!("send to smoltcp stack failed: {:?}", e); } } @@ -377,10 +380,7 @@ impl Socks5ServerNet { "receive from smoltcp stack and send to peer mgr packet, len = {}", data.len() ); - let packet = ZCPacket::new_from_buf( - bytes::BytesMut::from(bytes::Bytes::from(data)), - ZCPacketType::NIC, - ); + let packet = ZCPacket::new_from_buf(data, ZCPacketType::NIC); let Some(ipv4) = Ipv4Packet::new(packet.payload()) else { tracing::error!( payload_len = packet.payload_len(), diff --git a/easytier/src/gateway/tcp_proxy.rs b/easytier/src/gateway/tcp_proxy.rs index 69de208a..dfff31ff 100644 --- a/easytier/src/gateway/tcp_proxy.rs +++ b/easytier/src/gateway/tcp_proxy.rs @@ -563,7 +563,10 @@ impl TcpProxy { self.tasks.lock().unwrap().spawn(async move { while let Some(packet) = smoltcp_stack_receiver.recv().await { tracing::trace!(?packet, "receive from peer send to smoltcp packet"); - if let Err(e) = stack_sink.send(Ok(packet.payload().to_vec())).await { + if let Err(e) = stack_sink + .send(Ok(bytes::BytesMut::from(packet.payload()))) + .await + { tracing::error!("send to smoltcp stack failed: {:?}", e); } } @@ -577,10 +580,7 @@ impl TcpProxy { ?data, "receive from smoltcp stack and send to peer mgr packet" ); - let packet = ZCPacket::new_from_buf( - bytes::BytesMut::from(bytes::Bytes::from(data)), - ZCPacketType::NIC, - ); + let packet = ZCPacket::new_from_buf(data, ZCPacketType::NIC); let Some(ipv4) = Ipv4Packet::new(packet.payload()) else { tracing::error!( payload_len = packet.payload_len(), diff --git a/easytier/src/gateway/tokio_smoltcp/channel_device.rs b/easytier/src/gateway/tokio_smoltcp/channel_device.rs index b643bc49..58f116a4 100644 --- a/easytier/src/gateway/tokio_smoltcp/channel_device.rs +++ b/easytier/src/gateway/tokio_smoltcp/channel_device.rs @@ -1,3 +1,4 @@ +use bytes::BytesMut; use futures::{Sink, Stream}; use smoltcp::phy::DeviceCapabilities; use std::{ @@ -12,15 +13,15 @@ use super::device::AsyncDevice; /// A device that send and receive packets using a channel. pub struct ChannelDevice { - recv: Receiver>>, - send: PollSender>, + recv: Receiver>, + send: PollSender, caps: DeviceCapabilities, } pub type ChannelDeviceNewRet = ( ChannelDevice, - Sender>>, - Receiver>, + Sender>, + Receiver, ); impl ChannelDevice { @@ -43,25 +44,25 @@ impl ChannelDevice { } impl Stream for ChannelDevice { - type Item = io::Result>; + type Item = io::Result; fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { self.recv.poll_recv(cx) } } -fn map_err(e: PollSendError>) -> io::Error { +fn map_err(e: PollSendError) -> io::Error { io::Error::other(e) } -impl Sink> for ChannelDevice { +impl Sink for ChannelDevice { type Error = io::Error; fn poll_ready(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { self.send.poll_reserve(cx).map_err(map_err) } - fn start_send(mut self: Pin<&mut Self>, item: Vec) -> Result<(), Self::Error> { + fn start_send(mut self: Pin<&mut Self>, item: BytesMut) -> Result<(), Self::Error> { self.send.send_item(item).map_err(map_err) } diff --git a/easytier/src/gateway/tokio_smoltcp/device.rs b/easytier/src/gateway/tokio_smoltcp/device.rs index f701ec65..51fcaa9f 100644 --- a/easytier/src/gateway/tokio_smoltcp/device.rs +++ b/easytier/src/gateway/tokio_smoltcp/device.rs @@ -1,3 +1,4 @@ +use bytes::BytesMut; use futures::{Sink, Stream}; pub use smoltcp::phy::DeviceCapabilities; use smoltcp::{ @@ -10,7 +11,7 @@ use std::{collections::VecDeque, io}; pub const DEFAULT_MAX_BURST_SIZE: usize = 100; /// A packet used in `AsyncDevice`. -pub type Packet = Vec; +pub type Packet = BytesMut; /// A device that send and receive packets asynchronously. pub trait AsyncDevice: @@ -42,13 +43,11 @@ pub struct BufferDevice { pub struct BufferRxToken(Packet); impl RxToken for BufferRxToken { - fn consume(mut self, f: F) -> R + fn consume(self, f: F) -> R where F: FnOnce(&[u8]) -> R, { - let p = &mut self.0; - - f(p) + f(&self.0[..]) } } @@ -61,7 +60,8 @@ impl<'d> TxToken for BufferTxToken<'d> { F: FnOnce(&mut [u8]) -> R, { let tx_headroom = self.0.tx_headroom; - let mut buffer = vec![0u8; tx_headroom + len]; + let mut buffer = BytesMut::with_capacity(tx_headroom + len); + buffer.resize(tx_headroom + len, 0); let result = f(&mut buffer[tx_headroom..]); self.0.send_queue.push_back(buffer);