use std::sync::atomic::{fence, AtomicU64, Ordering};
use crate::error::{Result, RingfireError};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
#[repr(C)]
pub struct BlobRef {
pub offset: u64,
pub len: u32,
pub flags: u32,
}
impl BlobRef {
pub const EMPTY: Self = Self {
offset: 0,
len: 0,
flags: 0,
};
#[inline]
pub fn is_empty(&self) -> bool {
self.len == 0
}
}
#[repr(C, align(64))]
pub struct ArenaHeader {
pub capacity: u64,
pub mask: u64,
pub reserved: AtomicU64,
pub _pad: [u8; 40],
}
const _: () = {
assert!(std::mem::size_of::<ArenaHeader>() == 64);
assert!(std::mem::align_of::<ArenaHeader>() == 64);
};
pub struct PayloadArena {
header: *mut ArenaHeader,
data: *mut u8,
capacity: usize,
mask: usize,
}
unsafe impl Send for PayloadArena {}
unsafe impl Sync for PayloadArena {}
impl PayloadArena {
pub unsafe fn init(ptr: *mut u8, capacity: usize) -> Result<Self> {
if !capacity.is_power_of_two() || capacity < 64 {
return Err(RingfireError::InvalidCapacity(capacity as u64));
}
let header = ptr as *mut ArenaHeader;
unsafe {
(*header).capacity = capacity as u64;
(*header).mask = (capacity - 1) as u64;
(*header).reserved = AtomicU64::new(0);
(*header)._pad = [0u8; 40];
}
let data = unsafe { ptr.add(std::mem::size_of::<ArenaHeader>()) };
Ok(Self {
header,
data,
capacity,
mask: capacity - 1,
})
}
pub unsafe fn from_ptr(ptr: *mut u8) -> Result<Self> {
let header = ptr as *mut ArenaHeader;
let capacity = unsafe { (*header).capacity as usize };
if !capacity.is_power_of_two() || capacity < 64 {
return Err(RingfireError::InvalidCapacity(capacity as u64));
}
let mask = unsafe { (*header).mask as usize };
if mask != capacity - 1 {
return Err(RingfireError::CorruptLayout("arena mask does not match its capacity"));
}
let data = unsafe { ptr.add(std::mem::size_of::<ArenaHeader>()) };
Ok(Self {
header,
data,
capacity,
mask,
})
}
#[inline]
pub fn capacity(&self) -> usize {
self.capacity
}
#[inline]
pub fn reserve(&self, len: usize, flags: u32) -> Result<BlobRef> {
if len == 0 {
return Ok(BlobRef::EMPTY);
}
let max_len = self.capacity.min(u32::MAX as usize);
if len > max_len {
return Err(RingfireError::ArenaPayloadTooLarge {
len,
max_capacity: max_len,
});
}
let slot_size = (len + 63) & !63;
let header = unsafe { &*self.header };
loop {
let curr = header.reserved.load(Ordering::Relaxed);
let from_ix = (curr & (self.mask as u64)) as usize;
let (actual_offset, next) = if from_ix + len > self.capacity {
let aligned_to_end = (curr | (self.mask as u64)) + 1;
(aligned_to_end, aligned_to_end + slot_size as u64)
} else {
(curr, curr + slot_size as u64)
};
if header
.reserved
.compare_exchange_weak(curr, next, Ordering::AcqRel, Ordering::Relaxed)
.is_ok()
{
fence(Ordering::Release);
return Ok(BlobRef {
offset: actual_offset,
len: len as u32,
flags,
});
}
}
}
#[inline]
pub fn write_blob(&self, data: &[u8], flags: u32) -> Result<BlobRef> {
let blob_ref = self.reserve(data.len(), flags)?;
if !blob_ref.is_empty() {
let dest = self.slice_mut(blob_ref);
dest.copy_from_slice(data);
}
Ok(blob_ref)
}
#[inline]
pub fn write_blob_with<F, R>(&self, len: usize, flags: u32, f: F) -> Result<(BlobRef, R)>
where
F: FnOnce(&mut [u8]) -> R,
{
let blob_ref = self.reserve(len, flags)?;
let result = if !blob_ref.is_empty() {
let dest = self.slice_mut(blob_ref);
f(dest)
} else {
f(&mut [])
};
Ok((blob_ref, result))
}
#[inline]
pub fn read_blob(&self, blob_ref: BlobRef, out: &mut [u8]) -> Result<usize> {
let len = blob_ref.len as usize;
if len == 0 {
return Ok(0);
}
if !self.is_in_bounds(blob_ref) {
return Err(RingfireError::CorruptLayout("blob reference outside the arena"));
}
if out.len() < len {
return Err(RingfireError::BufferTooSmall {
required: len,
provided: out.len(),
});
}
let src = self.slice(blob_ref);
out[..len].copy_from_slice(src);
Ok(len)
}
#[inline]
pub fn view_blob<R>(&self, blob_ref: BlobRef, f: impl FnOnce(&[u8]) -> R) -> R {
assert!(self.is_in_bounds(blob_ref), "blob reference outside the arena");
if blob_ref.is_empty() {
f(&[])
} else {
f(self.slice(blob_ref))
}
}
#[inline]
pub fn is_lapped(&self, blob_ref: BlobRef) -> bool {
fence(Ordering::Acquire);
let curr = unsafe { (*self.header).reserved.load(Ordering::Relaxed) };
curr.saturating_sub(blob_ref.offset) > self.capacity as u64
}
#[inline]
pub fn is_in_bounds(&self, blob_ref: BlobRef) -> bool {
let start = (blob_ref.offset & self.mask as u64) as usize;
start + blob_ref.len as usize <= self.capacity
}
#[inline]
fn slice(&self, blob_ref: BlobRef) -> &[u8] {
let offset = (blob_ref.offset & (self.mask as u64)) as usize;
unsafe { std::slice::from_raw_parts(self.data.add(offset), blob_ref.len as usize) }
}
#[inline]
#[allow(clippy::mut_from_ref)]
fn slice_mut(&self, blob_ref: BlobRef) -> &mut [u8] {
let offset = (blob_ref.offset & (self.mask as u64)) as usize;
unsafe { std::slice::from_raw_parts_mut(self.data.add(offset), blob_ref.len as usize) }
}
}
#[cfg(test)]
#[path = "../tests/unit/arena.rs"]
mod tests;