mirror of
https://github.com/EasyTier/EasyTier.git
synced 2026-09-02 17:15:43 +00:00
perf: avoid smoltcp packet copy
Reserve NIC packet headroom in BufferDevice's TxToken so smoltcp writes the IP packet directly into a buf already carrying the ZCPacket NIC offset. socks5/tcp_proxy then wrap that buf zero-copy via ZCPacket::new_from_buf instead of copying through new_with_payload. Benchmark: tunnel::packet_def::tests::smoltcp_zcpacket_construct_bench copy (new_with_payload) 1280B: 13.5M pps zerocopy (new_from_buf) 1280B: 28.3M pps (2.09x) copy (new_with_payload) 4096B: 10.7M pps zerocopy (new_from_buf) 4096B: 19.1M pps (1.79x) Environment: AMD Ryzen 9 9955HX, rustc 1.95.0, --release, median of 3 runs
This commit is contained in:
@@ -28,7 +28,7 @@ use crate::{
|
|||||||
ip_reassembler::IpReassembler,
|
ip_reassembler::IpReassembler,
|
||||||
tokio_smoltcp::{BufferSize, Net, NetConfig, channel_device},
|
tokio_smoltcp::{BufferSize, Net, NetConfig, channel_device},
|
||||||
},
|
},
|
||||||
tunnel::packet_def::{PacketType, ZCPacket},
|
tunnel::packet_def::{PacketType, ZCPacket, ZCPacketType},
|
||||||
};
|
};
|
||||||
use anyhow::Context;
|
use anyhow::Context;
|
||||||
use dashmap::DashMap;
|
use dashmap::DashMap;
|
||||||
@@ -377,13 +377,19 @@ impl Socks5ServerNet {
|
|||||||
"receive from smoltcp stack and send to peer mgr packet, len = {}",
|
"receive from smoltcp stack and send to peer mgr packet, len = {}",
|
||||||
data.len()
|
data.len()
|
||||||
);
|
);
|
||||||
let Some(ipv4) = Ipv4Packet::new(&data) else {
|
let packet = ZCPacket::new_from_buf(
|
||||||
tracing::error!(?data, "smoltcp stack stream get non ipv4 packet");
|
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;
|
continue;
|
||||||
};
|
};
|
||||||
|
|
||||||
let dst = ipv4.get_destination();
|
let dst = ipv4.get_destination();
|
||||||
let packet = ZCPacket::new_with_payload(&data);
|
|
||||||
let Some(peer_manager) = peer_manager.upgrade() else {
|
let Some(peer_manager) = peer_manager.upgrade() else {
|
||||||
tracing::warn!("peer manager is gone, smoltcp sender exited");
|
tracing::warn!("peer manager is gone, smoltcp sender exited");
|
||||||
return;
|
return;
|
||||||
@@ -412,7 +418,8 @@ impl Socks5ServerNet {
|
|||||||
tcp_tx_size: 1024 * 128,
|
tcp_tx_size: 1024 * 128,
|
||||||
..Default::default()
|
..Default::default()
|
||||||
}),
|
}),
|
||||||
),
|
)
|
||||||
|
.with_packet_tx_headroom(ZCPacketType::NIC.get_packet_offsets().payload_offset),
|
||||||
);
|
);
|
||||||
|
|
||||||
let forward_tasks = Arc::new(std::sync::Mutex::new(forward_tasks));
|
let forward_tasks = Arc::new(std::sync::Mutex::new(forward_tasks));
|
||||||
|
|||||||
@@ -39,6 +39,8 @@ use super::CidrSet;
|
|||||||
|
|
||||||
#[cfg(feature = "smoltcp")]
|
#[cfg(feature = "smoltcp")]
|
||||||
use super::tokio_smoltcp::{self, Net, NetConfig, channel_device};
|
use super::tokio_smoltcp::{self, Net, NetConfig, channel_device};
|
||||||
|
#[cfg(feature = "smoltcp")]
|
||||||
|
use crate::tunnel::packet_def::ZCPacketType;
|
||||||
|
|
||||||
#[async_trait::async_trait]
|
#[async_trait::async_trait]
|
||||||
pub(crate) trait NatDstConnector: Send + Sync + Clone + 'static {
|
pub(crate) trait NatDstConnector: Send + Sync + Clone + 'static {
|
||||||
@@ -575,13 +577,19 @@ impl<C: NatDstConnector> TcpProxy<C> {
|
|||||||
?data,
|
?data,
|
||||||
"receive from smoltcp stack and send to peer mgr packet"
|
"receive from smoltcp stack and send to peer mgr packet"
|
||||||
);
|
);
|
||||||
let Some(ipv4) = Ipv4Packet::new(&data) else {
|
let packet = ZCPacket::new_from_buf(
|
||||||
tracing::error!(?data, "smoltcp stack stream get non ipv4 packet");
|
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;
|
continue;
|
||||||
};
|
};
|
||||||
|
|
||||||
let dst = ipv4.get_destination();
|
let dst = ipv4.get_destination();
|
||||||
let packet = ZCPacket::new_with_payload(&data);
|
|
||||||
let Some(peer_mgr) = peer_mgr.upgrade() else {
|
let Some(peer_mgr) = peer_mgr.upgrade() else {
|
||||||
tracing::warn!("peer manager is gone, smoltcp sender exited");
|
tracing::warn!("peer manager is gone, smoltcp sender exited");
|
||||||
return;
|
return;
|
||||||
@@ -610,7 +618,8 @@ impl<C: NatDstConnector> TcpProxy<C> {
|
|||||||
tcp_tx_size: 1024 * 16,
|
tcp_tx_size: 1024 * 16,
|
||||||
..Default::default()
|
..Default::default()
|
||||||
}),
|
}),
|
||||||
),
|
)
|
||||||
|
.with_packet_tx_headroom(ZCPacketType::NIC.get_packet_offsets().payload_offset),
|
||||||
);
|
);
|
||||||
net.set_any_ip(true);
|
net.set_any_ip(true);
|
||||||
self.smoltcp_net.lock().await.replace(net);
|
self.smoltcp_net.lock().await.replace(net);
|
||||||
|
|||||||
@@ -33,6 +33,7 @@ where
|
|||||||
pub struct BufferDevice {
|
pub struct BufferDevice {
|
||||||
caps: DeviceCapabilities,
|
caps: DeviceCapabilities,
|
||||||
max_burst_size: usize,
|
max_burst_size: usize,
|
||||||
|
tx_headroom: usize,
|
||||||
recv_queue: VecDeque<Packet>,
|
recv_queue: VecDeque<Packet>,
|
||||||
send_queue: VecDeque<Packet>,
|
send_queue: VecDeque<Packet>,
|
||||||
}
|
}
|
||||||
@@ -59,8 +60,9 @@ impl<'d> TxToken for BufferTxToken<'d> {
|
|||||||
where
|
where
|
||||||
F: FnOnce(&mut [u8]) -> R,
|
F: FnOnce(&mut [u8]) -> R,
|
||||||
{
|
{
|
||||||
let mut buffer = vec![0u8; len];
|
let tx_headroom = self.0.tx_headroom;
|
||||||
let result = f(&mut buffer);
|
let mut buffer = vec![0u8; tx_headroom + len];
|
||||||
|
let result = f(&mut buffer[tx_headroom..]);
|
||||||
|
|
||||||
self.0.send_queue.push_back(buffer);
|
self.0.send_queue.push_back(buffer);
|
||||||
|
|
||||||
@@ -98,11 +100,12 @@ impl Device for BufferDevice {
|
|||||||
}
|
}
|
||||||
|
|
||||||
impl 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);
|
let max_burst_size = caps.max_burst_size.unwrap_or(DEFAULT_MAX_BURST_SIZE);
|
||||||
BufferDevice {
|
BufferDevice {
|
||||||
caps,
|
caps,
|
||||||
max_burst_size,
|
max_burst_size,
|
||||||
|
tx_headroom,
|
||||||
recv_queue: VecDeque::with_capacity(max_burst_size),
|
recv_queue: VecDeque::with_capacity(max_burst_size),
|
||||||
send_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()
|
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]);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -51,6 +51,7 @@ pub struct NetConfig {
|
|||||||
pub ip_addr: IpCidr,
|
pub ip_addr: IpCidr,
|
||||||
pub gateway: Vec<IpAddress>,
|
pub gateway: Vec<IpAddress>,
|
||||||
pub buffer_size: BufferSize,
|
pub buffer_size: BufferSize,
|
||||||
|
pub packet_tx_headroom: usize,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl NetConfig {
|
impl NetConfig {
|
||||||
@@ -65,8 +66,14 @@ impl NetConfig {
|
|||||||
ip_addr,
|
ip_addr,
|
||||||
gateway,
|
gateway,
|
||||||
buffer_size: buffer_size.unwrap_or_default(),
|
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.
|
/// `Net` is the main interface to the network stack.
|
||||||
@@ -97,7 +104,8 @@ impl Net {
|
|||||||
}
|
}
|
||||||
|
|
||||||
fn new2<D: device::AsyncDevice + 'static>(device: D, config: NetConfig) -> Net {
|
fn new2<D: device::AsyncDevice + 'static>(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 mut iface = Interface::new(config.interface_config, &mut buffer_device, Instant::now());
|
||||||
let ip_addr = config.ip_addr;
|
let ip_addr = config.ip_addr;
|
||||||
iface.update_ip_addrs(|ip_addrs| {
|
iface.update_ip_addrs(|ip_addrs| {
|
||||||
|
|||||||
Reference in New Issue
Block a user