use crate::alloc::alloc::{self, Layout, LayoutError};
use std::sync::atomic::Ordering::{Acquire, Relaxed, Release};
use std::sync::atomic::{self, AtomicU32};
use std::{cell::Cell, cmp, mem, num::NonZeroUsize, ptr, ptr::NonNull, slice};
use crate::{BytePageSize, buf::UninitSlice, storage::INLINE_CAP, storage::Storage};
#[derive(Debug)]
pub(crate) struct SharedVec {
pub(crate) offset: u32,
pub(crate) len: u32,
pub(crate) capacity: u32,
pub(crate) remaining: u32,
pub(crate) ref_count: AtomicU32,
pub(crate) size: BytePageSize,
}
#[derive(Debug)]
pub(crate) struct StorageVec(pub(crate) NonNull<SharedVec>);
const KIND_VEC: usize = 0b01;
const KIND_OFFSET_BITS: usize = 2;
pub const METADATA_SIZE: usize = mem::size_of::<SharedVec>();
const METADATA_SIZE_U32: u32 = METADATA_SIZE as u32;
pub(crate) const MAX_CAPACITY: usize = u32::MAX as usize - METADATA_SIZE;
impl StorageVec {
pub(crate) fn with_capacity(capacity: usize) -> StorageVec {
StorageVec(SharedVec::create(BytePageSize::Unset, capacity, &[]))
}
pub(crate) fn sized(size: BytePageSize) -> StorageVec {
let cached = CACHE
.try_with(|c| {
let mut cst = c.take()?;
let item = cst.cache[size as usize].pop();
c.set(Some(cst));
item
})
.ok()
.flatten();
if let Some(mut item) = cached {
unsafe {
(*item.as_inner()).size = size;
}
item
} else {
StorageVec(SharedVec::create(size, size.capacity(), &[]))
}
}
pub(crate) fn from_slice(capacity: usize, src: &[u8]) -> StorageVec {
StorageVec(SharedVec::create(BytePageSize::Unset, capacity, src))
}
pub(crate) fn unsize(&mut self) {
unsafe { (*self.0.as_ptr()).size = BytePageSize::Unset }
}
pub(crate) fn as_ref(&self) -> &[u8] {
unsafe { slice::from_raw_parts(self.as_ptr(), self.len()) }
}
pub(crate) fn as_mut(&mut self) -> &mut [u8] {
unsafe { slice::from_raw_parts_mut(self.as_ptr(), self.len()) }
}
pub(crate) unsafe fn as_raw(&mut self) -> &mut [u8] {
slice::from_raw_parts_mut(self.as_ptr(), self.capacity())
}
pub(crate) unsafe fn as_ptr(&self) -> *mut u8 {
(self.0.as_ptr().cast::<u8>()).add((*self.0.as_ptr()).offset as usize)
}
fn as_inner(&mut self) -> *mut SharedVec {
self.0.as_ptr()
}
pub(crate) fn put_u8(&mut self, n: u8) {
let len = self.len();
unsafe {
let inner = self.as_inner();
(*inner).len += 1;
(*inner).remaining -= 1;
*self.as_ptr().add(len) = n;
}
}
#[inline]
pub(crate) fn spare_mut(&mut self) -> &mut UninitSlice {
unsafe { UninitSlice::from_raw_parts_mut(self.as_ptr().add(self.len()), self.remaining()) }
}
#[inline]
pub(crate) fn put_slice_partial(&mut self, src: &[u8]) -> usize {
let cnt = cmp::min(src.len(), self.remaining());
self.spare_mut()[..cnt].copy_from_slice(&src[..cnt]);
unsafe { self.set_len(self.len() + cnt) };
cnt
}
pub(crate) fn len(&self) -> usize {
unsafe { (*self.0.as_ptr()).len as usize }
}
pub(crate) fn capacity(&self) -> usize {
unsafe {
let inner = self.0.as_ref();
(inner.capacity + METADATA_SIZE_U32 - inner.offset) as usize
}
}
pub(crate) fn remaining(&self) -> usize {
unsafe { (*self.0.as_ptr()).remaining as usize }
}
pub(crate) fn is_full(&self) -> bool {
unsafe { (*self.0.as_ptr()).remaining == 0 }
}
pub(crate) fn is_unique(&self) -> bool {
unsafe { (*self.0.as_ptr()).is_unique() }
}
pub(crate) unsafe fn from_unique_view(
ptr: *mut SharedVec,
offset: usize,
len: usize,
) -> Option<StorageVec> {
if !(*ptr).is_unique() {
return None;
}
let end = (*ptr).capacity + METADATA_SIZE_U32;
(*ptr).offset = offset as u32;
(*ptr).len = len as u32;
(*ptr).remaining = end - (offset + len) as u32;
Some(StorageVec(NonNull::new_unchecked(ptr)))
}
pub(crate) fn shallow_freeze(&self) -> Storage {
unsafe {
if self.len() <= INLINE_CAP {
Storage::from_ptr_inline(self.as_ptr(), self.len())
} else {
let inner = self.0.as_ref();
let ref_cnt = inner.ref_count.fetch_add(1, Relaxed);
if ref_cnt == u32::MAX {
abort();
}
let offset = inner.offset as usize;
Storage {
ptr: (self.0.as_ptr().cast::<u8>()).add(offset),
len: self.len(),
offset: NonZeroUsize::new_unchecked((offset << KIND_OFFSET_BITS) ^ KIND_VEC),
}
}
}
}
pub(crate) fn freeze(self) -> Storage {
unsafe {
if self.len() <= INLINE_CAP {
Storage::from_ptr_inline(self.as_ptr(), self.len())
} else {
let inner = self.0.as_ref();
let offset = inner.offset as usize;
let inner = Storage {
ptr: (self.0.as_ptr().cast::<u8>()).add(offset),
len: self.len(),
offset: NonZeroUsize::new_unchecked((offset << KIND_OFFSET_BITS) ^ KIND_VEC),
};
mem::forget(self);
inner
}
}
}
pub(crate) fn split_to(&mut self, at: usize) -> Storage {
unsafe {
let ptr = self.as_ptr();
let other = if at <= INLINE_CAP {
Storage::from_ptr_inline(ptr, at)
} else {
let inner = self.as_inner();
let ref_cnt = (*inner).ref_count.fetch_add(1, Relaxed);
if ref_cnt == u32::MAX {
abort();
}
let offset = (*inner).offset as usize;
Storage {
ptr: (self.0.as_ptr().cast::<u8>()).add(offset),
len: at,
offset: NonZeroUsize::new_unchecked((offset << KIND_OFFSET_BITS) ^ KIND_VEC),
}
};
self.set_start(at);
other
}
}
pub(crate) fn truncate(&mut self, len: usize) {
unsafe {
if len == 0 {
let inner = self.as_inner();
if (*inner).is_unique() && (*inner).offset != METADATA_SIZE_U32 {
(*inner).len = 0;
(*inner).offset = METADATA_SIZE_U32;
(*inner).remaining = (*inner).capacity;
return;
}
}
if len < self.len() {
self.set_len(len);
}
}
}
pub(crate) fn resize(&mut self, new_len: usize, value: u8) {
let len = self.len();
if new_len > len {
let additional = new_len - len;
self.reserve(additional);
unsafe {
let dst = self.as_raw()[len..].as_mut_ptr();
ptr::write_bytes(dst, value, additional);
self.set_len(new_len);
}
} else {
self.truncate(new_len);
}
}
#[inline]
pub(crate) fn reserve_capacity(&mut self, capacity: usize) {
if capacity > self.len() {
*self = StorageVec(SharedVec::create(
BytePageSize::Unset,
capacity,
self.as_ref(),
));
}
}
#[inline]
pub(crate) fn reserve(&mut self, additional: usize) {
if additional <= self.remaining() {
return;
}
self.reserve_inner(additional, false);
}
#[inline]
pub(crate) fn reserve_exact(&mut self, additional: usize) {
if additional <= self.remaining() {
return;
}
self.reserve_inner(additional, true);
}
fn reserve_inner(&mut self, additional: usize, exact: bool) {
unsafe {
let inner = self.as_inner();
let len = (*inner).len as usize;
let new_cap = len
.checked_add(additional)
.expect("buffer capacity overflow");
let grow_cap = if exact {
new_cap
} else {
grown_capacity(len, new_cap)
};
if (*inner).is_unique() {
let capacity = (*inner).capacity as usize;
if capacity >= new_cap {
let offset = (*inner).offset;
(*inner).offset = METADATA_SIZE_U32;
(*inner).remaining = (capacity - len) as u32;
if len != 0 {
let ptr = self.0.as_ptr().cast::<u8>();
ptr::copy(ptr.add(offset as usize), ptr.add(METADATA_SIZE), len);
}
return;
}
if (*inner).size == BytePageSize::Unset {
self.realloc(len, capacity, grow_cap);
return;
}
}
*self = StorageVec(SharedVec::create(
BytePageSize::Unset,
grow_cap,
self.as_ref(),
));
}
}
unsafe fn realloc(&mut self, len: usize, capacity: usize, new_cap: usize) {
assert!(
new_cap <= MAX_CAPACITY,
"buffer capacity {new_cap} exceeds maximum {MAX_CAPACITY}"
);
let old_layout = shared_vec_layout(capacity).unwrap();
let new_layout = shared_vec_layout(new_cap).unwrap();
unsafe {
let ptr = self.0.as_ptr();
let offset = (*ptr).offset as usize;
if offset != METADATA_SIZE {
if len != 0 {
let data = ptr.cast::<u8>();
ptr::copy(data.add(offset), data.add(METADATA_SIZE), len);
}
(*ptr).offset = METADATA_SIZE_U32;
(*ptr).remaining = (capacity - len) as u32;
}
let new_ptr = alloc::realloc(ptr.cast(), old_layout, new_layout.size());
if new_ptr.is_null() {
alloc::handle_alloc_error(new_layout);
}
#[allow(clippy::cast_ptr_alignment)]
let inner = new_ptr.cast::<SharedVec>();
let capacity = (new_layout.size() - METADATA_SIZE) as u32;
(*inner).capacity = capacity;
(*inner).remaining = capacity - len as u32;
self.0 = NonNull::new_unchecked(inner);
}
}
#[inline]
pub(crate) unsafe fn set_len(&mut self, len: usize) {
let capacity = self.capacity();
let inner = self.as_inner();
assert!(len <= capacity);
(*inner).len = len as u32;
(*inner).remaining = (capacity - len) as u32;
}
pub(crate) unsafe fn set_start(&mut self, start: usize) {
if start != 0 {
let inner = self.as_inner();
assert!(
start <= (*inner).len as usize,
"cannot advance past the end of the buffer, cnt:{start} len:{}",
(*inner).len,
);
let start = start as u32;
(*inner).offset += start;
(*inner).len -= start;
}
}
}
unsafe impl Send for StorageVec {}
unsafe impl Sync for StorageVec {}
impl Drop for StorageVec {
fn drop(&mut self) {
release_shared_vec(self.0.as_ptr());
}
}
thread_local! {
static CACHE: Cell<Option<Box<Cache>>> = Cell::new(Some(Box::default()));
}
pub(crate) fn set_pages_cache(size: usize) {
let _ = CACHE.try_with(|c| {
if let Some(mut cst) = c.take() {
cst.size = size;
c.set(Some(cst));
}
});
}
const DEFAULT_PAGES_CACHE: usize = 16;
struct Cache {
size: usize,
cache: [Vec<StorageVec>; 7],
}
impl Default for Cache {
fn default() -> Self {
Self {
size: DEFAULT_PAGES_CACHE,
cache: Default::default(),
}
}
}
impl SharedVec {
pub(crate) fn create(size: BytePageSize, cap: usize, src: &[u8]) -> NonNull<SharedVec> {
assert!(
cap >= src.len(),
"SharedVec capacity {cap} is smaller than data length {}",
src.len()
);
let ptr = Self::alloc_with_capacity(size, cap, src.len() as u32);
unsafe {
let dst = ptr.add(METADATA_SIZE);
let sl = slice::from_raw_parts_mut(dst, src.len());
sl.copy_from_slice(src);
#[allow(clippy::cast_ptr_alignment)]
NonNull::new_unchecked(ptr.cast::<SharedVec>())
}
}
fn alloc_with_capacity(size: BytePageSize, cap: usize, len: u32) -> *mut u8 {
assert!(
cap <= MAX_CAPACITY,
"buffer capacity {cap} exceeds maximum {MAX_CAPACITY}"
);
let layout = shared_vec_layout(cap).unwrap();
unsafe {
let ptr = alloc::alloc(layout);
if ptr.is_null() {
alloc::handle_alloc_error(layout);
}
let capacity = (layout.size() - METADATA_SIZE) as u32;
#[cfg(feature = "overuse")]
if cap > 1081344 {
log::debug!("Buffer size {capacity}\n{:?}", backtrace::Backtrace::new());
}
#[allow(clippy::cast_ptr_alignment)]
ptr::write(
ptr.cast::<SharedVec>(),
SharedVec {
len,
capacity,
size,
remaining: capacity - len,
offset: METADATA_SIZE_U32,
ref_count: AtomicU32::new(1),
},
);
ptr
}
}
fn is_unique(&self) -> bool {
self.ref_count.load(Acquire) == 1
}
pub(crate) unsafe fn capacity(ptr: *const SharedVec) -> usize {
ptr::addr_of!((*ptr).capacity).read() as usize
}
}
pub(crate) fn release_shared_vec(ptr: *mut SharedVec) {
unsafe {
if (*ptr).ref_count.fetch_sub(1, Release) != 1 {
return;
}
atomic::fence(Acquire);
let capacity = (*ptr).capacity;
let size = (*ptr).size;
if size != BytePageSize::Unset {
let cached = CACHE.try_with(|c| {
let Some(mut cst) = c.take() else {
return false;
};
let res = if cst.cache[size as usize].len() < cst.size {
(*ptr).len = 0;
(*ptr).offset = METADATA_SIZE_U32;
(*ptr).remaining = capacity;
(*ptr).ref_count = AtomicU32::new(1);
(*ptr).size = BytePageSize::Unset;
cst.cache[size as usize].push(StorageVec(NonNull::new_unchecked(ptr)));
true
} else {
false
};
c.set(Some(cst));
res
});
if matches!(cached, Ok(true)) {
return;
}
}
ptr::drop_in_place(ptr);
let layout = shared_vec_layout(capacity as usize).unwrap();
alloc::dealloc(ptr.cast(), layout);
}
}
fn grown_capacity(len: usize, required: usize) -> usize {
cmp::max(required, cmp::min(len.saturating_mul(2), MAX_CAPACITY))
}
const fn shared_vec_layout(cap: usize) -> Result<Layout, LayoutError> {
let s_layout = match Layout::from_size_align(cap, Layout::new::<u8>().align()) {
Ok(l) => l,
Err(e) => return Err(e),
};
match Layout::new::<SharedVec>().pad_to_align().extend(s_layout) {
Ok((l, _)) => Ok(l),
Err(err) => Err(err),
}
}
#[inline(never)]
#[cold]
pub(crate) fn abort() -> ! {
std::process::abort()
}
#[cfg(test)]
#[allow(clippy::assert_is_empty)]
mod tests {
use super::*;
use crate::*;
#[test]
fn cached() {
super::CACHE.with(|cache| cache.set(Some(Box::default())));
let mut st = StorageVec::sized(BytePageSize::Size8);
assert_eq!(unsafe { (*st.0.as_ptr()).size }, BytePageSize::Size8);
st.put_u8(b'h');
let addr = st.0;
drop(st);
let st = StorageVec::sized(BytePageSize::Size8);
assert_eq!(addr, st.0);
}
#[test]
fn default_cache_limit_per_page_size() {
super::CACHE.with(|cache| cache.set(Some(Box::default())));
let pages: Vec<_> = (0..20)
.map(|_| StorageVec::sized(BytePageSize::Size4))
.collect();
drop(pages);
let cached = super::CACHE.with(|c| {
let cst = c.take().unwrap();
let len = cst.cache[BytePageSize::Size4 as usize].len();
c.set(Some(cst));
len
});
assert_eq!(cached, super::DEFAULT_PAGES_CACHE);
assert_eq!(super::DEFAULT_PAGES_CACHE, 16);
}
#[test]
fn page_allocation_is_category_size() {
for (size, alloc) in [
(BytePageSize::Size4, 4 * 1024),
(BytePageSize::Size16, 16 * 1024),
(BytePageSize::Size64, 64 * 1024),
] {
let st = StorageVec::sized(size);
assert_eq!(st.capacity(), size.capacity());
let layout = shared_vec_layout(st.capacity()).unwrap();
assert_eq!(layout.size(), alloc, "{size:?}");
}
}
#[test]
fn is_unique_synchronizes_with_release() {
let mut st = StorageVec::with_capacity(64);
for _ in 0..=INLINE_CAP {
st.put_u8(1);
}
let other = st.shallow_freeze();
assert!(!other.is_inline());
let handle = std::thread::spawn(move || {
let val = other.as_ref()[0];
drop(other);
val
});
while !st.is_unique() {
std::thread::yield_now();
}
st.as_mut()[0] = 2;
assert_eq!(handle.join().unwrap(), 1);
}
fn pattern(len: usize) -> Vec<u8> {
(0..len).map(|i| i as u8).collect()
}
#[test]
fn reserve_grows_unique_buffer() {
let data = pattern(100);
let mut st = StorageVec::from_slice(100, &data);
st.reserve(1000);
assert_eq!(st.as_ref(), &data[..]);
assert!(st.capacity() >= 1100);
assert_eq!(st.remaining(), st.capacity() - st.len());
let mut st = StorageVec::from_slice(100, &data);
unsafe { st.set_start(30) };
st.reserve(1000);
assert_eq!(st.as_ref(), &data[30..]);
assert!(st.capacity() >= 1070);
assert_eq!(st.remaining(), st.capacity() - st.len());
for i in 0..1000 {
st.put_u8(i as u8);
}
assert_eq!(&st.as_ref()[..70], &data[30..]);
assert_eq!(st.len(), 1070);
let mut st = StorageVec::from_slice(100, &data);
unsafe { st.set_start(100) };
st.reserve(1000);
assert!(st.as_ref().is_empty());
assert_eq!(st.remaining(), st.capacity());
}
#[test]
fn reserve_keeps_shared_views() {
let data = pattern(100);
let mut st = StorageVec::from_slice(100, &data);
let view = st.shallow_freeze();
assert!(!view.is_inline());
st.reserve(1000);
st.as_mut()[0] = 0xff;
assert_eq!(view.as_ref(), &data[..]);
assert_eq!(&st.as_ref()[1..], &data[1..]);
}
#[test]
fn reserve_pooled_page_leaves_page_to_cache() {
super::CACHE.with(|cache| cache.set(Some(Box::default())));
let mut st = StorageVec::sized(BytePageSize::Size8);
let page = st.0;
let data = pattern(st.capacity());
for b in &data {
st.put_u8(*b);
}
st.reserve(1);
assert_eq!(unsafe { (*st.0.as_ptr()).size }, BytePageSize::Unset);
assert_ne!(st.0, page);
assert_eq!(st.as_ref(), &data[..]);
let st2 = StorageVec::sized(BytePageSize::Size8);
assert_eq!(st2.0, page);
}
fn cached_pages(size: BytePageSize) -> usize {
super::CACHE.with(|c| {
let cst = c.take().unwrap();
let len = cst.cache[size as usize].len();
c.set(Some(cst));
len
})
}
#[test]
fn pages_cache_size() {
super::CACHE.with(|cache| cache.set(Some(Box::default())));
crate::set_pages_cache(1);
drop((
StorageVec::sized(BytePageSize::Size4),
StorageVec::sized(BytePageSize::Size4),
));
assert_eq!(cached_pages(BytePageSize::Size4), 1);
crate::set_pages_cache(0);
drop(StorageVec::sized(BytePageSize::Size8));
assert_eq!(cached_pages(BytePageSize::Size8), 0);
crate::set_pages_cache(16);
let st = StorageVec::sized(BytePageSize::Size16);
let cache = super::CACHE.with(Cell::take);
drop(st);
super::CACHE.with(|c| c.set(cache));
assert_eq!(cached_pages(BytePageSize::Size16), 0);
let cache = super::CACHE.with(Cell::take);
crate::set_pages_cache(3);
super::CACHE.with(|c| c.set(cache));
assert_eq!(
super::CACHE.with(|c| {
let cst = c.take().unwrap();
let size = cst.size;
c.set(Some(cst));
size
}),
16
);
}
#[test]
fn truncate_reclaims_unique_buffer() {
let mut b = BytesMut::with_capacity(128);
let cap = b.capacity();
b.extend_from_slice(&[1; 64]);
b.advance(0);
b.advance(32);
assert_eq!(b.capacity(), cap - 32);
b.truncate(0);
assert_eq!(b.capacity(), cap);
b.extend_from_slice(&[1; 64]);
b.advance(32);
let other = b.split_to(30);
b.truncate(0);
assert_eq!(b.capacity(), cap - 62);
drop(other);
}
#[test]
fn resize_shrinks() {
let mut b = BytesMut::copy_from_slice(b"hello world");
b.resize(5, 0);
assert_eq!(&b[..], b"hello");
b.resize(7, b'!');
assert_eq!(&b[..], b"hello!!");
}
#[test]
fn reserve_reclaims_front_space() {
let mut b = BytesMut::with_capacity(128);
let cap = b.capacity();
b.extend_from_slice(&[1; 40]);
b.extend_from_slice(&[2; 10]);
b.advance(40);
let spare = b.capacity() - b.len();
b.reserve(spare + 1);
assert_eq!(&b[..], &[2; 10]);
assert_eq!(b.capacity(), cap);
assert!(b.is_unique());
}
#[test]
fn reserve_exact() {
let mut b = BytesMut::with_capacity(64);
b.extend_from_slice(&[1; 64]);
let cap = b.capacity();
b.reserve_exact(0);
assert_eq!(b.capacity(), cap);
b.reserve_exact(1000);
assert_eq!(&b[..], &[1; 64]);
assert!(b.capacity() >= 1064);
assert!(b.capacity() < 1064 + 2 * METADATA_SIZE);
let other = b.split_to(32);
b.reserve_exact(2000);
assert_eq!(&b[..], &[1; 32]);
assert!(b.capacity() >= 2032);
assert!(b.capacity() < 2032 + 2 * METADATA_SIZE);
assert!(b.is_unique());
assert_eq!(&other[..], &[1; 32]);
}
}