mirror of
https://github.com/EasyTier/EasyTier.git
synced 2026-10-08 10:56:13 -08:00
feat(buf): add BufPool, BufMargins, and BufList for zero-copy packet buffering (#2625)
This commit is contained in:
1 parent
6f68ac3758
commit
a017dc133c
4 files changed
+289
No files matched your search
@@ -0,0 +1,128 @@
|
||||
use std::collections::VecDeque;
|
||||
use std::io::IoSlice;
|
||||
|
||||
use bytes::{Buf, BufMut, Bytes, BytesMut};
|
||||
|
||||
#[derive(Debug, Default)]
|
||||
pub struct BufList<T> {
|
||||
bufs: VecDeque<T>,
|
||||
}
|
||||
|
||||
impl<T: Buf> BufList<T> {
|
||||
pub fn new() -> BufList<T> {
|
||||
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<T> {
|
||||
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<T: Buf> Extend<T> for BufList<T> {
|
||||
fn extend<I: IntoIterator<Item = T>>(&mut self, iter: I) {
|
||||
self.bufs.extend(
|
||||
iter.into_iter()
|
||||
.inspect(|buf| debug_assert!(buf.has_remaining())),
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: Buf> Buf for BufList<T> {
|
||||
#[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);
|
||||
}
|
||||
}
|
||||
@@ -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;
|
||||
|
||||
@@ -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<u8>] {
|
||||
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);
|
||||
}
|
||||
}
|
||||
@@ -1,3 +1,4 @@
|
||||
pub mod buf;
|
||||
#[cfg(feature = "management")]
|
||||
pub mod panic;
|
||||
pub mod string;
|
||||
|
||||
Reference in new issue
Block a user