use crate::iobuf::pool::{BufferPoolThreadCache, SizeClassLease};
use bytes::Bytes;
use std::{
alloc::{Layout, alloc, alloc_zeroed, dealloc, handle_alloc_error},
mem::{ManuallyDrop, MaybeUninit, align_of, offset_of, size_of},
ptr::{self, NonNull, addr_of_mut},
sync::atomic::Ordering,
};
cfg_if::cfg_if! {
if #[cfg(feature = "loom")] {
use loom::sync::atomic::{AtomicUsize, fence};
} else {
use std::sync::atomic::{AtomicUsize, fence};
}
}
const OWNER_EMPTY: usize = 0b00;
const OWNER_HEAP: usize = 0b01;
const OWNER_POOLED: usize = 0b10;
const OWNER_EXTERNAL: usize = 0b11;
const OWNER_TAG_MASK: usize = 0b11;
const MAX_REFCOUNT: usize = isize::MAX as usize;
const _: () = assert!(align_of::<HeapOwner>() >= 4);
const _: () = assert!(align_of::<PooledOwner>() >= 4);
const _: () = assert!(align_of::<ExternalOwner>() >= 4);
const _: () = assert!(size_of::<HeapOwner>().is_multiple_of(align_of::<HeapOwner>()));
const _: () = assert!(offset_of!(HeapOwner, refs) == 0);
const _: () = assert!(offset_of!(PooledOwner, refs) == 0);
const _: () = assert!(offset_of!(ExternalOwner, refs) == 0);
#[derive(Clone, Copy, Debug)]
#[repr(transparent)]
pub(crate) struct OwnerRef(*mut ());
unsafe impl Send for OwnerRef {}
unsafe impl Sync for OwnerRef {}
impl OwnerRef {
#[inline(always)]
pub(crate) const fn empty() -> Self {
Self(ptr::null_mut())
}
#[inline(always)]
pub(crate) const fn is_empty(self) -> bool {
self.0.is_null()
}
#[inline(always)]
fn tag(self) -> usize {
self.0.addr() & OWNER_TAG_MASK
}
#[inline(always)]
pub(crate) fn is_pooled(self) -> bool {
self.tag() == OWNER_POOLED
}
#[inline(always)]
pub(crate) fn is_external(self) -> bool {
self.tag() == OWNER_EXTERNAL
}
#[inline(always)]
fn from_tagged<T>(ptr: NonNull<T>, tag: usize) -> Self {
let ptr = ptr.cast::<()>().as_ptr();
Self(ptr.with_addr(ptr.addr() | tag))
}
#[inline(always)]
unsafe fn untag<T>(self) -> NonNull<T> {
let ptr = self
.0
.with_addr(self.0.addr() & !OWNER_TAG_MASK)
.cast::<T>();
unsafe { NonNull::new_unchecked(ptr) }
}
#[inline(always)]
fn split(self) -> (usize, *mut ()) {
let addr = self.0.addr();
(
addr & OWNER_TAG_MASK,
self.0.with_addr(addr & !OWNER_TAG_MASK),
)
}
#[inline(always)]
unsafe fn from_heap(header: NonNull<HeapOwner>) -> Self {
Self::from_tagged(header, OWNER_HEAP)
}
#[inline(always)]
pub(crate) unsafe fn from_pooled(slot: NonNull<PooledOwner>) -> Self {
Self::from_tagged(slot, OWNER_POOLED)
}
#[inline(always)]
unsafe fn from_external(owner: NonNull<ExternalOwner>) -> Self {
Self::from_tagged(owner, OWNER_EXTERNAL)
}
pub(crate) fn from_vec(vec: Vec<u8>) -> (NonNull<u8>, usize, Self) {
if vec.is_empty() {
return (NonNull::dangling(), 0, Self::empty());
}
match HeapOwner::try_adopt_vec(vec) {
Ok((ptr, len, _, owner)) => (ptr, len, owner),
Err(vec) => Self::from_bytes(Bytes::from(vec)),
}
}
pub(crate) fn from_bytes(bytes: Bytes) -> (NonNull<u8>, usize, Self) {
if bytes.is_empty() {
return (NonNull::dangling(), 0, Self::empty());
}
ExternalOwner::from_bytes(bytes)
}
#[inline(always)]
unsafe fn heap(self) -> NonNull<HeapOwner> {
unsafe { self.untag() }
}
#[inline(always)]
unsafe fn pooled(self) -> NonNull<PooledOwner> {
unsafe { self.untag() }
}
#[inline(always)]
unsafe fn external(self) -> NonNull<ExternalOwner> {
unsafe { self.untag() }
}
#[inline(always)]
unsafe fn refs(self) -> &'static AtomicUsize {
unsafe { self.untag::<AtomicUsize>().as_ref() }
}
#[inline(always)]
pub(crate) unsafe fn external_bytes<'a>(self) -> &'a Bytes {
unsafe { &(*self.external().as_ptr()).bytes }
}
#[inline(always)]
pub(crate) unsafe fn clone_shared(self) {
if self.is_empty() {
return;
}
let old = unsafe { self.refs() }.fetch_add(1, Ordering::Relaxed);
if old > MAX_REFCOUNT {
std::process::abort();
}
}
#[inline(always)]
pub(crate) unsafe fn drop_shared(self) {
if self.is_empty() {
return;
}
let refs = unsafe { self.refs() };
if refs.load(Ordering::Acquire) == 1 {
unsafe { self.release_unique_outlined() };
return;
}
let old = refs.fetch_sub(1, Ordering::Release);
if old == 1 {
unsafe { self.drop_shared_race_final(refs) };
}
}
#[inline(never)]
unsafe fn release_unique_outlined(self) {
unsafe { self.release_unique() };
}
#[cold]
#[inline(never)]
unsafe fn drop_shared_race_final(self, refs: &AtomicUsize) {
fence(Ordering::Acquire);
refs.store(1, Ordering::Relaxed);
unsafe { self.release_unique() };
}
#[inline(always)]
pub(crate) unsafe fn release_unique(self) {
let (tag, owner) = self.split();
if tag == OWNER_POOLED {
unsafe { PooledOwner::release_to_thread_cache(NonNull::new_unchecked(owner.cast())) };
return;
}
unsafe { self.release_unique_cold() };
}
#[inline(always)]
pub(crate) unsafe fn release_unique_mut_at(self, ptr: NonNull<u8>, cap: usize) {
let (tag, owner) = self.split();
if tag == OWNER_POOLED {
unsafe { PooledOwner::release_to_thread_cache(NonNull::new_unchecked(owner.cast())) };
return;
}
if tag == OWNER_EMPTY {
return;
}
if self.is_front_heap_for_mut(ptr) {
unsafe { HeapOwner::release_front(NonNull::new_unchecked(owner.cast()), ptr, cap) };
return;
}
unsafe { HeapOwner::release(NonNull::new_unchecked(owner.cast())) };
}
#[cold]
#[inline(never)]
unsafe fn release_unique_cold(self) {
let tag = self.tag();
if tag == OWNER_HEAP {
unsafe { HeapOwner::release(self.heap()) };
} else if tag == OWNER_EXTERNAL {
unsafe { ExternalOwner::release(self.external()) };
}
}
#[inline(always)]
pub(crate) unsafe fn ensure_heap_header_for_mut(&mut self, ptr: NonNull<u8>, cap: usize) {
if !self.is_front_heap_for_mut(ptr) {
return;
}
let base = unsafe { self.heap() };
let data_base = HeapOwner::front_data_base(base);
let alloc_size = HeapOwner::front_alloc_size(base, ptr, cap);
unsafe {
base.as_ptr().write(HeapOwner {
refs: AtomicUsize::new(1),
data_base,
alloc_size,
alloc_align: align_of::<HeapOwner>(),
});
}
}
#[inline(always)]
fn is_front_heap_for_mut(self, ptr: NonNull<u8>) -> bool {
self.tag() == OWNER_HEAP && (self.0.addr() & !OWNER_TAG_MASK) < ptr.as_ptr().addr()
}
#[inline(always)]
pub(crate) unsafe fn is_unique(self) -> bool {
unsafe { self.refs() }.load(Ordering::Acquire) == 1
}
#[inline]
pub(crate) unsafe fn data_base(self) -> NonNull<u8> {
if self.is_pooled() {
unsafe { self.pooled().as_ref().data_base }
} else {
unsafe { self.heap().as_ref().data_base }
}
}
#[inline]
pub(crate) unsafe fn usable_capacity(self) -> usize {
if self.is_pooled() {
unsafe { self.pooled().as_ref().capacity }
} else {
let header = unsafe { self.heap() };
let header_ref = unsafe { header.as_ref() };
HeapOwner::usable_capacity(header, header_ref)
}
}
#[cfg(all(test, not(feature = "loom")))]
pub(crate) unsafe fn refcount(self) -> Option<usize> {
if self.is_empty() {
return None;
}
Some(unsafe { self.refs() }.load(Ordering::Acquire))
}
#[cfg(all(test, feature = "loom"))]
pub(crate) unsafe fn refcount_relaxed(self) -> usize {
unsafe { self.refs() }.load(Ordering::Relaxed)
}
}
#[repr(C)]
pub(crate) struct HeapOwner {
refs: AtomicUsize,
data_base: NonNull<u8>,
alloc_size: usize,
alloc_align: usize,
}
impl HeapOwner {
#[inline]
pub(crate) fn allocate_aligned(
capacity: usize,
alignment: usize,
zeroed: bool,
) -> (NonNull<u8>, usize, OwnerRef) {
assert!(capacity > 0, "capacity must be greater than zero");
assert!(
alignment.is_power_of_two(),
"alignment must be a power of two"
);
let (layout, header_offset) = Self::layout(capacity, alignment);
let ptr = if zeroed {
unsafe { alloc_zeroed(layout) }
} else {
unsafe { alloc(layout) }
};
let data = NonNull::new(ptr).unwrap_or_else(|| handle_alloc_error(layout));
let owner = unsafe {
let header = data.as_ptr().add(header_offset).cast::<Self>();
header.write(Self {
refs: AtomicUsize::new(1),
data_base: data,
alloc_size: layout.size(),
alloc_align: layout.align(),
});
OwnerRef::from_heap(NonNull::new_unchecked(header))
};
(data, header_offset, owner)
}
#[inline(always)]
pub(crate) fn allocate_aligned_mut(
capacity: usize,
alignment: usize,
zeroed: bool,
) -> (NonNull<u8>, usize, OwnerRef) {
assert!(capacity > 0, "capacity must be greater than zero");
assert!(
alignment.is_power_of_two(),
"alignment must be a power of two"
);
if alignment > align_of::<Self>() {
return Self::allocate_aligned(capacity, alignment, zeroed);
}
let layout = Self::front_layout(capacity);
let ptr = if zeroed {
unsafe { alloc_zeroed(layout) }
} else {
unsafe { alloc(layout) }
};
let base = NonNull::new(ptr).unwrap_or_else(|| handle_alloc_error(layout));
let header = base.cast::<Self>();
let data = Self::front_data_base(header);
let owner = unsafe { OwnerRef::from_heap(header) };
(data, capacity, owner)
}
pub(crate) fn try_adopt_vec(
vec: Vec<u8>,
) -> Result<(NonNull<u8>, usize, usize, OwnerRef), Vec<u8>> {
let len = vec.len();
let cap = vec.capacity();
let base_addr = vec.as_ptr() as usize;
let Some(header_offset) = Self::vec_adoption_header_offset(base_addr, len, cap) else {
return Err(vec);
};
let mut vec = ManuallyDrop::new(vec);
let base = vec.as_mut_ptr();
unsafe {
let header = base.add(header_offset).cast::<Self>();
header.write(Self {
refs: AtomicUsize::new(1),
data_base: NonNull::new_unchecked(base),
alloc_size: cap,
alloc_align: 1,
});
Ok((
NonNull::new_unchecked(base),
len,
header_offset,
OwnerRef::from_heap(NonNull::new_unchecked(header)),
))
}
}
#[inline]
fn layout(capacity: usize, alignment: usize) -> (Layout, usize) {
let header_offset = capacity
.checked_next_multiple_of(align_of::<Self>())
.expect("layout size overflow");
let total = header_offset
.checked_add(size_of::<Self>())
.expect("heap layout size overflow");
let layout_alignment = alignment.max(align_of::<Self>());
let layout = Layout::from_size_align(total, layout_alignment)
.expect("heap layout size overflow or alignment not a power of two");
(layout, header_offset)
}
#[inline(always)]
const fn front_layout(capacity: usize) -> Layout {
let total = size_of::<Self>()
.checked_add(capacity)
.expect("front heap layout size overflow");
match Layout::from_size_align(total, align_of::<Self>()) {
Ok(layout) => layout,
Err(_) => panic!("front heap layout size overflow"),
}
}
#[inline(always)]
const fn front_data_base(base: NonNull<Self>) -> NonNull<u8> {
unsafe { NonNull::new_unchecked(base.as_ptr().cast::<u8>().add(size_of::<Self>())) }
}
#[inline(always)]
fn front_alloc_size(base: NonNull<Self>, ptr: NonNull<u8>, cap: usize) -> usize {
let base_addr = base.as_ptr() as usize;
let end_addr = ptr.as_ptr() as usize + cap;
assert!(end_addr >= base_addr);
end_addr - base_addr
}
#[inline(always)]
fn usable_capacity(header: NonNull<Self>, header_ref: &Self) -> usize {
let header_addr = header.as_ptr() as usize;
let data_addr = header_ref.data_base.as_ptr() as usize;
if header_addr < data_addr {
header_ref
.alloc_size
.checked_sub(data_addr - header_addr)
.expect("front heap data base must lie within allocation")
} else {
header_addr - data_addr
}
}
#[inline(always)]
fn round_down(value: usize, align: usize) -> usize {
assert!(align.is_power_of_two());
value & !(align - 1)
}
#[inline(always)]
fn vec_adoption_header_offset(base_addr: usize, len: usize, cap: usize) -> Option<usize> {
if cap < size_of::<Self>() {
return None;
}
let header_addr = Self::round_down(
base_addr.checked_add(cap - size_of::<Self>())?,
align_of::<Self>(),
);
if header_addr < base_addr || header_addr < base_addr.checked_add(len)? {
return None;
}
Some(header_addr - base_addr)
}
#[inline]
unsafe fn release(header: NonNull<Self>) {
let header_ref = unsafe { header.as_ref() };
assert_eq!(header_ref.refs.load(Ordering::Relaxed), 1);
let header_addr = header.as_ptr() as usize;
let data_addr = header_ref.data_base.as_ptr() as usize;
let base = if header_addr < data_addr {
header.cast::<u8>()
} else {
header_ref.data_base
};
let layout = unsafe {
Layout::from_size_align_unchecked(header_ref.alloc_size, header_ref.alloc_align)
};
unsafe { dealloc(base.as_ptr(), layout) };
}
#[inline(always)]
unsafe fn release_front(base: NonNull<Self>, ptr: NonNull<u8>, cap: usize) {
let alloc_size = Self::front_alloc_size(base, ptr, cap);
let layout = unsafe { Layout::from_size_align_unchecked(alloc_size, align_of::<Self>()) };
unsafe { dealloc(base.as_ptr().cast::<u8>(), layout) };
}
}
#[repr(C)]
pub struct PooledOwner {
refs: AtomicUsize,
lease: MaybeUninit<SizeClassLease>,
data_base: NonNull<u8>,
capacity: usize,
slot: u32,
}
impl PooledOwner {
#[inline]
#[allow(clippy::missing_const_for_fn)]
pub fn new(slot: u32, capacity: usize) -> Self {
Self {
refs: AtomicUsize::new(1),
lease: MaybeUninit::uninit(),
data_base: NonNull::dangling(),
capacity,
slot,
}
}
#[inline]
pub(crate) fn layout(size: usize, alignment: usize) -> Layout {
Layout::from_size_align(size, alignment)
.expect("pool layout size overflow or alignment not a power of two")
}
#[inline(always)]
unsafe fn release_to_thread_cache(owner: NonNull<Self>) {
let buffer = unsafe { PooledBuffer::from_owner(owner) };
BufferPoolThreadCache::push(buffer);
}
}
pub struct PooledBuffer {
owner: NonNull<PooledOwner>,
}
unsafe impl Send for PooledBuffer {}
unsafe impl Sync for PooledBuffer {}
impl std::fmt::Debug for PooledBuffer {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("PooledBuffer")
.field("owner", &self.owner)
.field("slot", &self.slot())
.field("ptr", &self.data_ptr())
.finish()
}
}
impl PooledBuffer {
#[inline]
pub unsafe fn new(owner: NonNull<PooledOwner>, layout: Layout, zeroed: bool) -> Self {
assert!(layout.size() > 0, "pooled data layout must be non-zero");
let ptr = if zeroed {
unsafe { alloc_zeroed(layout) }
} else {
unsafe { alloc(layout) }
};
let ptr = NonNull::new(ptr).unwrap_or_else(|| handle_alloc_error(layout));
unsafe {
assert_eq!((*owner.as_ptr()).refs.load(Ordering::Relaxed), 1);
addr_of_mut!((*owner.as_ptr()).data_base).write(ptr);
}
Self { owner }
}
#[inline(always)]
pub(crate) const unsafe fn from_owner(owner: NonNull<PooledOwner>) -> Self {
Self { owner }
}
#[cfg(any(all(test, not(feature = "loom")), feature = "bench"))]
#[inline(always)]
pub const fn as_ptr(&self) -> *mut u8 {
self.data_ptr().as_ptr()
}
#[inline(always)]
pub(crate) const fn data_ptr(&self) -> NonNull<u8> {
unsafe { self.owner.as_ref().data_base }
}
#[inline(always)]
pub(crate) const fn capacity(&self) -> usize {
unsafe { self.owner.as_ref().capacity }
}
#[inline(always)]
pub(crate) const fn slot(&self) -> u32 {
unsafe { self.owner.as_ref().slot }
}
#[inline(always)]
pub(crate) const fn into_slot(self) -> u32 {
self.slot()
}
#[inline(always)]
pub(crate) unsafe fn init_lease(&mut self, lease: SizeClassLease) {
unsafe {
addr_of_mut!((*self.owner.as_ptr()).lease).write(MaybeUninit::new(lease));
}
}
#[inline(always)]
pub(crate) const unsafe fn lease(&self) -> &SizeClassLease {
unsafe { &*self.owner.as_ref().lease.as_ptr() }
}
#[inline(always)]
pub(crate) const unsafe fn take_lease(&mut self) -> SizeClassLease {
unsafe { self.owner.as_mut().lease.assume_init_read() }
}
#[inline(always)]
pub(crate) unsafe fn owner_ref(&self) -> OwnerRef {
unsafe { OwnerRef::from_pooled(self.owner) }
}
#[inline(always)]
pub unsafe fn deallocate(self, layout: Layout) {
unsafe { dealloc(self.data_ptr().as_ptr(), layout) };
}
#[cfg(feature = "loom")]
pub(crate) fn assert_parked_sentinel(&self) {
let refs = unsafe { &self.owner.as_ref().refs };
assert_eq!(refs.load(Ordering::Relaxed), 1);
}
}
#[repr(C)]
struct ExternalOwner {
refs: AtomicUsize,
bytes: Bytes,
}
impl ExternalOwner {
fn from_bytes(bytes: Bytes) -> (NonNull<u8>, usize, OwnerRef) {
let owner = Box::new(Self {
refs: AtomicUsize::new(1),
bytes,
});
let ptr = NonNull::new(owner.bytes.as_ptr().cast_mut())
.expect("non-empty Bytes has non-null data");
let len = owner.bytes.len();
let owner = NonNull::from(Box::leak(owner));
let owner = unsafe { OwnerRef::from_external(owner) };
(ptr, len, owner)
}
#[inline]
unsafe fn release(owner: NonNull<Self>) {
let owner_ref = unsafe { owner.as_ref() };
assert_eq!(owner_ref.refs.load(Ordering::Relaxed), 1);
drop(unsafe { Box::from_raw(owner.as_ptr()) });
}
}
#[cfg(all(test, not(feature = "loom")))]
mod tests {
use super::*;
use crate::iobuf::page_size;
use commonware_utils::NZUsize;
#[test]
fn test_heap_layout_places_tail_header_after_data() {
let page = page_size();
let (data, usable, owner) = HeapOwner::allocate_aligned(4096, page, false);
assert!((data.as_ptr() as usize).is_multiple_of(page));
assert_eq!(usable, 4096);
let header = unsafe { owner.heap() };
assert!((header.as_ptr() as usize).is_multiple_of(align_of::<HeapOwner>()));
assert!(header.as_ptr() as usize >= data.as_ptr() as usize + 4096);
assert_eq!(unsafe { owner.data_base() }, data);
assert_eq!(unsafe { owner.usable_capacity() }, 4096);
unsafe { owner.release_unique() };
}
#[test]
fn test_heap_zeroed_only_exposes_usable_region() {
let (data, _, owner) = HeapOwner::allocate_aligned(64, page_size(), true);
let bytes = unsafe { std::slice::from_raw_parts(data.as_ptr(), 64) };
assert_eq!(bytes, &[0u8; 64]);
unsafe { owner.release_unique() };
}
#[test]
fn test_heap_unaligned_capacity_rounds_usable_region_up() {
let (_, usable, owner) = HeapOwner::allocate_aligned(10, 1, false);
let capacity = unsafe { owner.usable_capacity() };
assert_eq!(capacity, 10usize.next_multiple_of(align_of::<HeapOwner>()));
assert_eq!(usable, capacity);
unsafe { owner.release_unique() };
}
#[test]
fn test_front_layout_accepts_maximum_valid_capacity() {
let capacity = isize::MAX as usize - size_of::<HeapOwner>() - (align_of::<HeapOwner>() - 1);
let layout = HeapOwner::front_layout(capacity);
assert_eq!(layout.size(), size_of::<HeapOwner>() + capacity);
assert_eq!(layout.align(), align_of::<HeapOwner>());
}
#[test]
#[should_panic(expected = "front heap layout size overflow")]
fn test_front_layout_rejects_isize_max_overflow() {
let _ = HeapOwner::front_layout(isize::MAX as usize);
}
#[test]
#[should_panic(expected = "front heap layout size overflow")]
fn test_front_layout_rejects_usize_overflow() {
let _ = HeapOwner::front_layout(usize::MAX);
}
#[test]
#[should_panic(expected = "layout size overflow")]
fn test_layout_rejects_rounding_overflow() {
let _ = HeapOwner::layout(usize::MAX, 64);
}
#[test]
#[should_panic(expected = "heap layout size overflow")]
fn test_layout_rejects_header_add_overflow() {
let _ = HeapOwner::layout(usize::MAX - 31, 64);
}
#[test]
#[should_panic(expected = "heap layout size overflow or alignment not a power of two")]
fn test_layout_rejects_isize_max_overflow() {
let _ = HeapOwner::layout(isize::MAX as usize, 64);
}
#[test]
fn test_front_heap_mut_drop_does_not_read_reserved_header() {
let (data, cap, owner) = HeapOwner::allocate_aligned_mut(64, 1, false);
assert_eq!(cap, 64);
assert!((data.as_ptr() as usize).is_multiple_of(align_of::<HeapOwner>()));
assert!(owner.is_front_heap_for_mut(data));
unsafe { owner.release_unique_mut_at(data, 64) };
let (data, cap, owner) = HeapOwner::allocate_aligned_mut(64, 1, false);
assert_eq!(cap, 64);
let advanced = unsafe { data.add(17) };
assert!(owner.is_front_heap_for_mut(advanced));
unsafe { owner.release_unique_mut_at(advanced, 47) };
}
#[test]
fn test_front_heap_materializes_before_shared_owner() {
let (data, _, mut owner) = HeapOwner::allocate_aligned_mut(64, 1, false);
let advanced = unsafe { data.add(17) };
unsafe { owner.ensure_heap_header_for_mut(advanced, 47) };
assert!(owner.is_front_heap_for_mut(advanced));
assert_eq!(unsafe { owner.data_base() }, data);
assert_eq!(unsafe { owner.usable_capacity() }, 64);
assert_eq!(unsafe { owner.refcount() }, Some(1));
unsafe { owner.drop_shared() };
}
#[test]
fn test_mut_allocator_uses_tail_header_for_high_alignment() {
let page = page_size();
let (data, cap, owner) = HeapOwner::allocate_aligned_mut(64, page, false);
assert_eq!(cap, 64);
assert!((data.as_ptr() as usize).is_multiple_of(page));
assert!(!owner.is_front_heap_for_mut(data));
unsafe { owner.release_unique_mut_at(data, 64) };
}
#[test]
fn test_front_heap_zeroed_exposes_zeroed_data_region() {
let (data, _, owner) = HeapOwner::allocate_aligned_mut(64, 1, true);
let bytes = unsafe { std::slice::from_raw_parts(data.as_ptr(), 64) };
assert_eq!(bytes, &[0u8; 64]);
unsafe { owner.release_unique_mut_at(data, 64) };
}
#[test]
fn test_external_owner_refcount() {
let (_, len, owner) = OwnerRef::from_bytes(Bytes::from_static(b"abc"));
assert_eq!(len, 3);
assert!(owner.is_external());
assert_eq!(unsafe { owner.refcount() }, Some(1));
unsafe { owner.clone_shared() };
assert_eq!(unsafe { owner.refcount() }, Some(2));
unsafe { owner.drop_shared() };
assert_eq!(unsafe { owner.refcount() }, Some(1));
unsafe { owner.drop_shared() };
}
#[test]
fn test_external_owner_keeps_inner_bytes_alive() {
let payload = Bytes::from(vec![7u8; 32]);
let inner_ptr = payload.as_ptr();
let (ptr, len, owner) = OwnerRef::from_bytes(payload);
assert_eq!(ptr.as_ptr().cast_const(), inner_ptr);
assert_eq!(len, 32);
let inner = unsafe { owner.external_bytes() };
assert_eq!(inner.as_ref(), &[7u8; 32]);
unsafe { owner.drop_shared() };
}
#[test]
fn test_empty_vec_has_no_owner() {
let (_, len, owner) = OwnerRef::from_vec(Vec::new());
assert_eq!(len, 0);
assert!(owner.is_empty());
assert_eq!(unsafe { owner.refcount() }, None);
}
#[test]
fn test_empty_bytes_has_no_owner() {
let (_, len, owner) = OwnerRef::from_bytes(Bytes::new());
assert_eq!(len, 0);
assert!(owner.is_empty());
}
#[test]
fn test_vec_adoption_with_spare_capacity() {
let mut vec = Vec::with_capacity(256);
vec.extend_from_slice(&[1u8, 2, 3, 4]);
let base_addr = vec.as_ptr() as usize;
let cap = vec.capacity();
let (ptr, len, owner) = OwnerRef::from_vec(vec);
assert_eq!(ptr.as_ptr() as usize, base_addr);
assert_eq!(len, 4);
assert!(!owner.is_external());
assert!(!owner.is_pooled());
assert!(!owner.is_empty());
let expected_header =
base_addr + HeapOwner::vec_adoption_header_offset(base_addr, len, cap).unwrap();
assert_eq!(unsafe { owner.data_base() }.as_ptr() as usize, base_addr);
let usable = unsafe { owner.usable_capacity() };
assert_eq!(usable, expected_header - base_addr);
let payload = unsafe { std::slice::from_raw_parts(ptr.as_ptr(), len) };
assert_eq!(payload, &[1, 2, 3, 4]);
unsafe { owner.release_unique() };
}
#[test]
fn test_vec_adoption_exact_size_falls_back_to_external() {
let vec = vec![5u8, 6, 7];
assert_eq!(vec.len(), vec.capacity());
let (ptr, len, owner) = OwnerRef::from_vec(vec);
assert_eq!(len, 3);
assert!(owner.is_external());
let payload = unsafe { std::slice::from_raw_parts(ptr.as_ptr(), len) };
assert_eq!(payload, &[5, 6, 7]);
unsafe { owner.drop_shared() };
}
#[test]
fn test_vec_adoption_boundary_matches_placement_rule() {
for spare in 0..(size_of::<HeapOwner>() + 2 * align_of::<HeapOwner>()) {
let len = 16;
let mut vec = Vec::with_capacity(len + spare);
vec.extend_from_slice(&[9u8; 16]);
let base_addr = vec.as_ptr() as usize;
let cap = vec.capacity();
let (ptr, out_len, owner) = OwnerRef::from_vec(vec);
assert_eq!(out_len, len, "spare={spare} cap={cap}");
if owner.is_external() {
assert!(
cap - len < size_of::<HeapOwner>() + align_of::<HeapOwner>() - 1,
"spare={spare} cap={cap} base={base_addr:#x}"
);
assert!(
HeapOwner::vec_adoption_header_offset(base_addr, len, cap).is_none(),
"spare={spare} cap={cap} base={base_addr:#x}"
);
} else {
let header_offset = unsafe { owner.usable_capacity() };
assert!(header_offset >= len, "spare={spare} cap={cap}");
assert!((base_addr + header_offset).is_multiple_of(align_of::<HeapOwner>()));
assert!(header_offset + size_of::<HeapOwner>() <= cap);
assert_eq!(ptr.as_ptr() as usize, base_addr);
}
let payload = unsafe { std::slice::from_raw_parts(ptr.as_ptr(), out_len) };
assert_eq!(payload, &[9u8; 16], "spare={spare} cap={cap}");
unsafe { owner.drop_shared() };
}
}
#[test]
fn test_vec_adoption_rejects_header_before_base() {
let align = align_of::<HeapOwner>();
let base_addr = align - 1;
let cap = size_of::<HeapOwner>();
let header_addr = HeapOwner::round_down(base_addr + cap - size_of::<HeapOwner>(), align);
assert!(header_addr < base_addr);
assert_eq!(
HeapOwner::vec_adoption_header_offset(base_addr, 0, cap),
None
);
}
#[test]
fn test_heap_owner_shared_clone_and_drop() {
let (_, _, owner) = HeapOwner::allocate_aligned(64, 1, false);
unsafe { owner.clone_shared() };
assert_eq!(unsafe { owner.refcount() }, Some(2));
assert!(!unsafe { owner.is_unique() });
unsafe { owner.drop_shared() };
assert!(unsafe { owner.is_unique() });
unsafe { owner.drop_shared() };
}
#[test]
fn test_pooled_layout_is_data_only() {
let size = 1024;
let layout = PooledOwner::layout(size, NZUsize!(64).get());
assert_eq!(layout.size(), size);
assert!(layout.align() >= 64);
assert!(align_of::<PooledOwner>() >= 4);
}
}
#[cfg(all(test, feature = "loom"))]
mod loom_tests {
use super::*;
use loom::{
cell::UnsafeCell,
sync::{Arc, atomic::AtomicUsize},
thread,
};
struct Tracker(Arc<AtomicUsize>);
impl Drop for Tracker {
fn drop(&mut self) {
self.0.fetch_add(1, Ordering::SeqCst);
}
}
impl AsRef<[u8]> for Tracker {
fn as_ref(&self) -> &[u8] {
&[1, 2, 3]
}
}
#[test]
fn shared_clone_drop_releases_exactly_once() {
loom::model(|| {
let released = Arc::new(AtomicUsize::new(0));
let bytes = Bytes::from_owner(Tracker(released.clone()));
let (_, len, owner) = OwnerRef::from_bytes(bytes);
assert_eq!(len, 3);
unsafe { owner.clone_shared() };
unsafe { owner.clone_shared() };
let t1 = thread::spawn(move || {
unsafe { owner.drop_shared() };
});
let t2 = thread::spawn(move || {
unsafe { owner.drop_shared() };
});
unsafe { owner.drop_shared() };
t1.join().unwrap();
t2.join().unwrap();
assert_eq!(released.load(Ordering::SeqCst), 1);
});
}
#[test]
fn clone_races_concurrent_drop() {
loom::model(|| {
let released = Arc::new(AtomicUsize::new(0));
let bytes = Bytes::from_owner(Tracker(released.clone()));
let (_, _, owner) = OwnerRef::from_bytes(bytes);
unsafe { owner.clone_shared() };
let t1 = thread::spawn(move || {
unsafe { owner.clone_shared() };
unsafe { owner.drop_shared() };
unsafe { owner.drop_shared() };
});
unsafe { owner.drop_shared() };
t1.join().unwrap();
assert_eq!(released.load(Ordering::SeqCst), 1);
});
}
#[test]
fn is_unique_races_final_drop() {
loom::model(|| {
let released = Arc::new(AtomicUsize::new(0));
let bytes = Bytes::from_owner(Tracker(released.clone()));
let (_, _, owner) = OwnerRef::from_bytes(bytes);
let payload = Arc::new(UnsafeCell::new(0usize));
unsafe { owner.clone_shared() };
let t1 = thread::spawn({
let payload = payload.clone();
move || {
payload.with_mut(|cell| unsafe { *cell = 1 });
unsafe { owner.drop_shared() };
}
});
loop {
if unsafe { owner.is_unique() } {
payload.with(|cell| assert_eq!(unsafe { *cell }, 1));
assert_eq!(released.load(Ordering::SeqCst), 0);
break;
}
thread::yield_now();
}
unsafe { owner.drop_shared() };
t1.join().unwrap();
assert_eq!(released.load(Ordering::SeqCst), 1);
});
}
struct PayloadCell(UnsafeCell<usize>);
unsafe impl Send for PayloadCell {}
unsafe impl Sync for PayloadCell {}
struct TrackedPayload {
payload: Arc<PayloadCell>,
released: Arc<AtomicUsize>,
}
impl Drop for TrackedPayload {
fn drop(&mut self) {
self.payload.0.with(|cell| {
assert_eq!(unsafe { *cell }, 1);
});
self.released.fetch_add(1, Ordering::SeqCst);
}
}
impl AsRef<[u8]> for TrackedPayload {
fn as_ref(&self) -> &[u8] {
&[1, 2, 3]
}
}
#[test]
fn payload_writes_happen_before_race_final_release() {
loom::model(|| {
let released = Arc::new(AtomicUsize::new(0));
let payload = Arc::new(PayloadCell(UnsafeCell::new(0)));
let bytes = Bytes::from_owner(TrackedPayload {
payload: payload.clone(),
released: released.clone(),
});
let (_, _, owner) = OwnerRef::from_bytes(bytes);
unsafe { owner.clone_shared() };
let t1 = thread::spawn({
move || {
payload.0.with_mut(|cell| unsafe { *cell = 1 });
unsafe { owner.drop_shared() };
}
});
unsafe { owner.drop_shared() };
t1.join().unwrap();
assert_eq!(released.load(Ordering::SeqCst), 1);
});
}
#[test]
fn payload_writes_happen_before_fast_path_release() {
loom::model(|| {
let released = Arc::new(AtomicUsize::new(0));
let payload = Arc::new(PayloadCell(UnsafeCell::new(0)));
let bytes = Bytes::from_owner(TrackedPayload {
payload: payload.clone(),
released: released.clone(),
});
let (_, _, owner) = OwnerRef::from_bytes(bytes);
unsafe { owner.clone_shared() };
let t1 = thread::spawn({
move || {
payload.0.with_mut(|cell| unsafe { *cell = 1 });
unsafe { owner.drop_shared() };
}
});
while unsafe { owner.refcount_relaxed() } != 1 {
thread::yield_now();
}
unsafe { owner.drop_shared() };
t1.join().unwrap();
assert_eq!(released.load(Ordering::SeqCst), 1);
});
}
}