diff --git a/easytier/src/gateway/socks5.rs b/easytier/src/gateway/socks5.rs index fe98fd0b..82ead977 100644 --- a/easytier/src/gateway/socks5.rs +++ b/easytier/src/gateway/socks5.rs @@ -28,7 +28,7 @@ use crate::{ ip_reassembler::IpReassembler, tokio_smoltcp::{BufferSize, Net, NetConfig, channel_device}, }, - tunnel::packet_def::{PacketType, ZCPacket}, + tunnel::packet_def::{PacketType, ZCPacket, ZCPacketType}, }; use anyhow::Context; use dashmap::DashMap; @@ -377,13 +377,19 @@ impl Socks5ServerNet { "receive from smoltcp stack and send to peer mgr packet, len = {}", data.len() ); - let Some(ipv4) = Ipv4Packet::new(&data) else { - tracing::error!(?data, "smoltcp stack stream get non ipv4 packet"); + let packet = ZCPacket::new_from_buf( + bytes::BytesMut::from(bytes::Bytes::from(data)), + ZCPacketType::NIC, + ); + let Some(ipv4) = Ipv4Packet::new(packet.payload()) else { + tracing::error!( + payload_len = packet.payload_len(), + "smoltcp stack stream get non ipv4 packet" + ); continue; }; let dst = ipv4.get_destination(); - let packet = ZCPacket::new_with_payload(&data); let Some(peer_manager) = peer_manager.upgrade() else { tracing::warn!("peer manager is gone, smoltcp sender exited"); return; @@ -412,7 +418,8 @@ impl Socks5ServerNet { tcp_tx_size: 1024 * 128, ..Default::default() }), - ), + ) + .with_packet_tx_headroom(ZCPacketType::NIC.get_packet_offsets().payload_offset), ); let forward_tasks = Arc::new(std::sync::Mutex::new(forward_tasks)); diff --git a/easytier/src/gateway/tcp_proxy.rs b/easytier/src/gateway/tcp_proxy.rs index 6e252268..69de208a 100644 --- a/easytier/src/gateway/tcp_proxy.rs +++ b/easytier/src/gateway/tcp_proxy.rs @@ -39,6 +39,8 @@ use super::CidrSet; #[cfg(feature = "smoltcp")] use super::tokio_smoltcp::{self, Net, NetConfig, channel_device}; +#[cfg(feature = "smoltcp")] +use crate::tunnel::packet_def::ZCPacketType; #[async_trait::async_trait] pub(crate) trait NatDstConnector: Send + Sync + Clone + 'static { @@ -575,13 +577,19 @@ impl TcpProxy { ?data, "receive from smoltcp stack and send to peer mgr packet" ); - let Some(ipv4) = Ipv4Packet::new(&data) else { - tracing::error!(?data, "smoltcp stack stream get non ipv4 packet"); + let packet = ZCPacket::new_from_buf( + bytes::BytesMut::from(bytes::Bytes::from(data)), + ZCPacketType::NIC, + ); + let Some(ipv4) = Ipv4Packet::new(packet.payload()) else { + tracing::error!( + payload_len = packet.payload_len(), + "smoltcp stack stream get non ipv4 packet" + ); continue; }; let dst = ipv4.get_destination(); - let packet = ZCPacket::new_with_payload(&data); let Some(peer_mgr) = peer_mgr.upgrade() else { tracing::warn!("peer manager is gone, smoltcp sender exited"); return; @@ -610,7 +618,8 @@ impl TcpProxy { tcp_tx_size: 1024 * 16, ..Default::default() }), - ), + ) + .with_packet_tx_headroom(ZCPacketType::NIC.get_packet_offsets().payload_offset), ); net.set_any_ip(true); self.smoltcp_net.lock().await.replace(net); diff --git a/easytier/src/gateway/tokio_smoltcp/device.rs b/easytier/src/gateway/tokio_smoltcp/device.rs index 4df75d3b..f701ec65 100644 --- a/easytier/src/gateway/tokio_smoltcp/device.rs +++ b/easytier/src/gateway/tokio_smoltcp/device.rs @@ -33,6 +33,7 @@ where pub struct BufferDevice { caps: DeviceCapabilities, max_burst_size: usize, + tx_headroom: usize, recv_queue: VecDeque, send_queue: VecDeque, } @@ -59,8 +60,9 @@ impl<'d> TxToken for BufferTxToken<'d> { where F: FnOnce(&mut [u8]) -> R, { - let mut buffer = vec![0u8; len]; - let result = f(&mut buffer); + let tx_headroom = self.0.tx_headroom; + let mut buffer = vec![0u8; tx_headroom + len]; + let result = f(&mut buffer[tx_headroom..]); self.0.send_queue.push_back(buffer); @@ -98,11 +100,12 @@ impl Device for BufferDevice { } impl BufferDevice { - pub(crate) fn new(caps: DeviceCapabilities) -> BufferDevice { + pub(crate) fn new(caps: DeviceCapabilities, tx_headroom: usize) -> BufferDevice { let max_burst_size = caps.max_burst_size.unwrap_or(DEFAULT_MAX_BURST_SIZE); BufferDevice { caps, max_burst_size, + tx_headroom, recv_queue: VecDeque::with_capacity(max_burst_size), send_queue: VecDeque::with_capacity(max_burst_size), } @@ -123,3 +126,27 @@ impl BufferDevice { self.recv_queue.is_empty() } } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn buffer_device_reserves_tx_headroom() { + let mut caps = DeviceCapabilities::default(); + caps.max_burst_size = Some(1); + let mut device = BufferDevice::new(caps, 16); + + let token = device.transmit(Instant::now()).unwrap(); + token.consume(4, |buf| { + assert_eq!(buf.len(), 4); + buf.copy_from_slice(&[1, 2, 3, 4]); + }); + + let mut queue = device.take_send_queue(); + let packet = queue.pop_front().unwrap(); + assert_eq!(packet.len(), 20); + assert_eq!(&packet[..16], &[0; 16]); + assert_eq!(&packet[16..], &[1, 2, 3, 4]); + } +} diff --git a/easytier/src/gateway/tokio_smoltcp/mod.rs b/easytier/src/gateway/tokio_smoltcp/mod.rs index adb2d2f3..8e16f7dd 100644 --- a/easytier/src/gateway/tokio_smoltcp/mod.rs +++ b/easytier/src/gateway/tokio_smoltcp/mod.rs @@ -51,6 +51,7 @@ pub struct NetConfig { pub ip_addr: IpCidr, pub gateway: Vec, pub buffer_size: BufferSize, + pub packet_tx_headroom: usize, } impl NetConfig { @@ -65,8 +66,14 @@ impl NetConfig { ip_addr, gateway, buffer_size: buffer_size.unwrap_or_default(), + packet_tx_headroom: 0, } } + + pub fn with_packet_tx_headroom(mut self, packet_tx_headroom: usize) -> Self { + self.packet_tx_headroom = packet_tx_headroom; + self + } } /// `Net` is the main interface to the network stack. @@ -97,7 +104,8 @@ impl Net { } fn new2(device: D, config: NetConfig) -> Net { - let mut buffer_device = BufferDevice::new(device.capabilities().clone()); + let mut buffer_device = + BufferDevice::new(device.capabilities().clone(), config.packet_tx_headroom); let mut iface = Interface::new(config.interface_config, &mut buffer_device, Instant::now()); let ip_addr = config.ip_addr; iface.update_ip_addrs(|ip_addrs| {