From a017dc133c4595d19cdbd4e322153d348acfb139 Mon Sep 17 00:00:00 2001 From: Luna Yao <40349250+ZnqbuZ@users.noreply.github.com> Date: Thu, 1 Oct 2026 03:32:45 +0200 Subject: [PATCH] feat(buf): add BufPool, BufMargins, and BufList for zero-copy packet buffering (#2625) --- easytier-core/src/tunnel/buf.rs | 128 +++++++++++++++++++++++++ easytier-core/src/tunnel/mod.rs | 1 + easytier/src/utils/buf.rs | 159 ++++++++++++++++++++++++++++++++ easytier/src/utils/mod.rs | 1 + 4 files changed, 289 insertions(+) create mode 100644 easytier-core/src/tunnel/buf.rs create mode 100644 easytier/src/utils/buf.rs diff --git a/easytier-core/src/tunnel/buf.rs b/easytier-core/src/tunnel/buf.rs new file mode 100644 index 00000000..689d4ff6 --- /dev/null +++ b/easytier-core/src/tunnel/buf.rs @@ -0,0 +1,128 @@ +use std::collections::VecDeque; +use std::io::IoSlice; + +use bytes::{Buf, BufMut, Bytes, BytesMut}; + +#[derive(Debug, Default)] +pub struct BufList { + bufs: VecDeque, +} + +impl BufList { + pub fn new() -> BufList { + BufList { + bufs: VecDeque::new(), + } + } + + #[inline] + pub fn push(&mut self, buf: T) { + debug_assert!(buf.has_remaining()); + self.bufs.push_back(buf); + } + + #[inline] + pub fn pop(&mut self) -> Option { + self.bufs.pop_front() + } + + #[inline] + pub fn len(&self) -> usize { + self.bufs.len() + } + + #[inline] + pub fn is_empty(&self) -> bool { + self.bufs.is_empty() + } +} + +impl Extend for BufList { + fn extend>(&mut self, iter: I) { + self.bufs.extend( + iter.into_iter() + .inspect(|buf| debug_assert!(buf.has_remaining())), + ); + } +} + +impl Buf for BufList { + #[inline] + fn remaining(&self) -> usize { + self.bufs.iter().map(|buf| buf.remaining()).sum() + } + + #[inline] + fn chunk(&self) -> &[u8] { + self.bufs.front().map(Buf::chunk).unwrap_or_default() + } + + #[inline] + fn chunks_vectored<'t>(&'t self, dst: &mut [IoSlice<'t>]) -> usize { + if dst.is_empty() { + return 0; + } + let mut vecs = 0; + for buf in &self.bufs { + vecs += buf.chunks_vectored(&mut dst[vecs..]); + if vecs == dst.len() { + break; + } + } + vecs + } + + #[inline] + fn advance(&mut self, mut cnt: usize) { + while cnt > 0 { + { + let front = &mut self.bufs[0]; + let rem = front.remaining(); + if rem > cnt { + front.advance(cnt); + return; + } else { + front.advance(rem); + cnt -= rem; + } + } + self.bufs.pop_front(); + } + } + + #[inline] + fn copy_to_bytes(&mut self, len: usize) -> Bytes { + // Our inner buffer may have an optimized version of copy_to_bytes, and if the whole + // request can be fulfilled by the front buffer, we can take advantage. + match self.bufs.front_mut() { + Some(front) if front.remaining() == len => { + let b = front.copy_to_bytes(len); + self.bufs.pop_front(); + b + } + Some(front) if front.remaining() > len => front.copy_to_bytes(len), + _ => { + assert!(len <= self.remaining(), "`len` greater than remaining"); + let mut bm = BytesMut::with_capacity(len); + bm.put(self.take(len)); + bm.freeze() + } + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_buf_list() { + let mut list = BufList::new(); + list.push(Bytes::from_static(b"hello ")); + list.push(Bytes::from_static(b"world")); + assert_eq!(list.remaining(), 11); + let bytes = list.copy_to_bytes(11); + assert_eq!(&bytes[..], b"hello world"); + assert_eq!(list.remaining(), 0); + } +} diff --git a/easytier-core/src/tunnel/mod.rs b/easytier-core/src/tunnel/mod.rs index 1a6b046c..7ccc8168 100644 --- a/easytier-core/src/tunnel/mod.rs +++ b/easytier-core/src/tunnel/mod.rs @@ -15,6 +15,7 @@ use crate::{foundation::time::error::Elapsed, packet::ZCPacket, proto::common::T pub use crate::socket::IpVersion; +pub mod buf; pub(crate) mod encrypt; pub mod filter; pub mod framed; diff --git a/easytier/src/utils/buf.rs b/easytier/src/utils/buf.rs new file mode 100644 index 00000000..fe4ed932 --- /dev/null +++ b/easytier/src/utils/buf.rs @@ -0,0 +1,159 @@ +use bytes::{BufMut, BytesMut}; +use derive_more::{From, Into}; +use std::mem::MaybeUninit; +use std::ptr::copy_nonoverlapping; + +pub use easytier_core::tunnel::buf::BufList; + +#[derive(Debug, Clone, Copy, Default, From, Into, PartialEq, Eq)] +pub struct BufMargins { + pub header: usize, + pub trailer: usize, +} + +impl BufMargins { + #[inline(always)] + pub fn size(&self) -> usize { + self.header + self.trailer + } +} + +#[derive(Debug)] +pub struct BufPool { + pool: BytesMut, + pub min_capacity: usize, +} + +impl BufPool { + #[inline(always)] + pub fn new(min_capacity: usize) -> Self { + Self { + pool: BytesMut::with_capacity(min_capacity), + min_capacity, + } + } + + #[inline(always)] + pub fn reserve(&mut self, additional: usize) { + if self.pool.capacity() - self.pool.len() < additional { + self.pool.reserve(additional.max(self.min_capacity)); + } + } + + #[inline(always)] + pub fn split(&mut self) -> BytesMut { + self.pool.split() + } + + #[inline] + pub fn write(&mut self, chunk: &[u8], margins: BufMargins) { + let len = margins.size() + chunk.len(); + self.reserve(len); + unsafe { + copy_nonoverlapping( + chunk.as_ptr(), + self.pool.chunk_mut().as_mut_ptr().add(margins.header), + chunk.len(), + ); + self.pool.advance_mut(len); + } + } + + #[inline(always)] + pub fn buf(&mut self, chunk: &[u8], margins: BufMargins) -> BytesMut { + self.write(chunk, margins); + self.pool.split() + } + + #[inline(always)] + pub fn writer(&mut self, capacity: usize, margins: BufMargins) -> BufPoolWriter<'_> { + assert!(capacity >= margins.size()); + self.reserve(capacity); + BufPoolWriter { + pool: self, + capacity, + margins, + } + } +} + +#[derive(Debug)] +pub struct BufPoolWriter<'t> { + pool: &'t mut BufPool, + capacity: usize, + margins: BufMargins, +} + +impl<'t> BufPoolWriter<'t> { + #[inline(always)] + pub fn reserve(&mut self, additional: usize) { + if self.capacity < additional { + self.pool.reserve(additional); + self.capacity += additional; + } + } + + #[inline(always)] + pub fn split(&mut self) -> BytesMut { + self.pool.split() + } + + #[inline(always)] + pub fn as_slice(&mut self) -> &mut [MaybeUninit] { + unsafe { + self.pool + .pool + .spare_capacity_mut() + .get_unchecked_mut(self.margins.header..self.capacity - self.margins.trailer) + } + } + + #[inline(always)] + pub fn commit(&mut self, written: usize) { + let len = self.margins.size() + written; + assert!(self.capacity >= len); + self.capacity -= len; + unsafe { + self.pool.pool.advance_mut(len); + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_buf_pool_write() { + let mut pool = BufPool::new(1024); + let margins = BufMargins { + header: 4, + trailer: 2, + }; + let data = b"hello world"; + pool.write(data, margins); + let buf = pool.split(); + assert_eq!(buf.len(), 4 + data.len() + 2); + assert_eq!(&buf[4..4 + data.len()], data); + } + + #[test] + fn test_buf_pool_writer() { + let mut pool = BufPool::new(1024); + let margins = BufMargins { + header: 10, + trailer: 6, + }; + let mut writer = pool.writer(64, margins); + let slice = writer.as_slice(); + assert_eq!(slice.len(), 64 - 10 - 6); + let data = b"packet content"; + unsafe { + std::ptr::copy_nonoverlapping(data.as_ptr(), slice.as_mut_ptr() as *mut u8, data.len()); + } + writer.commit(data.len()); + let buf = writer.split(); + assert_eq!(buf.len(), 10 + data.len() + 6); + assert_eq!(&buf[10..10 + data.len()], data); + } +} diff --git a/easytier/src/utils/mod.rs b/easytier/src/utils/mod.rs index f9956096..d9deb752 100644 --- a/easytier/src/utils/mod.rs +++ b/easytier/src/utils/mod.rs @@ -1,3 +1,4 @@ +pub mod buf; #[cfg(feature = "management")] pub mod panic; pub mod string;