use crate::core::{Error, Result};
use core::ptr::NonNull;
use core::sync::atomic::AtomicU64;
use crate::memory::ring::SpscRing;
use crate::memory::slot;
pub const REGION_MAGIC: u64 = 0x464C_5942_5F52_4701;
pub const REGION_VERSION: u16 = 1;
const HEADER_OFFSET: usize = 0;
const HEADER_BYTES: usize = 64;
const PRODUCER_OFFSET: usize = 64; const CONSUMER_OFFSET: usize = 128; pub const SLOTS_OFFSET: usize = 192;
#[repr(C)]
struct RegionHeader {
magic: u64, version: u16, flags: u16, _pad0: u32, slot_count: u32, slot_size: u32, _pad1: [u8; 8], region_id: u128, _pad2: [u8; 16], }
const _: () = assert!(
core::mem::size_of::<RegionHeader>() == HEADER_BYTES,
"RegionHeader must be exactly 64 bytes"
);
pub struct Region {
ptr: NonNull<u8>,
len: usize,
slot_count: usize,
slot_size: usize,
ring: SpscRing,
}
unsafe impl Send for Region {}
unsafe impl Sync for Region {}
impl Region {
pub fn anonymous(slot_count: usize, slot_size: usize) -> Result<Self> {
Self::validate_params(slot_count, slot_size)?;
let len = SLOTS_OFFSET + slot_count * slot_size;
let raw = unsafe {
libc::mmap(
core::ptr::null_mut(),
len,
libc::PROT_READ | libc::PROT_WRITE,
libc::MAP_PRIVATE | libc::MAP_ANON,
-1,
0,
)
};
if raw == libc::MAP_FAILED {
return Err(Error::from(std::io::Error::last_os_error()));
}
let ptr = unsafe { NonNull::new_unchecked(raw as *mut u8) };
unsafe {
let header = RegionHeader {
magic: REGION_MAGIC,
version: REGION_VERSION,
flags: 0,
_pad0: 0,
slot_count: slot_count as u32,
slot_size: slot_size as u32,
_pad1: [0; 8],
region_id: 0,
_pad2: [0; 16],
};
core::ptr::write(ptr.as_ptr().add(HEADER_OFFSET) as *mut RegionHeader, header);
}
let ring = unsafe {
let head = NonNull::new_unchecked(ptr.as_ptr().add(PRODUCER_OFFSET) as *mut AtomicU64);
let tail = NonNull::new_unchecked(ptr.as_ptr().add(CONSUMER_OFFSET) as *mut AtomicU64);
SpscRing::new(head, tail, slot_count)
};
Ok(Self {
ptr,
len,
slot_count,
slot_size,
ring,
})
}
fn validate_params(slot_count: usize, slot_size: usize) -> Result<()> {
if slot_count == 0 || !slot_count.is_power_of_two() {
return Err(Error::config("slot_count must be a non-zero power of two"));
}
if slot_count > u32::MAX as usize {
return Err(Error::config("slot_count exceeds u32::MAX"));
}
if slot_size < slot::HEADER_SIZE {
return Err(Error::config("slot_size must be at least HEADER_SIZE (32)"));
}
if slot_size > u32::MAX as usize {
return Err(Error::config("slot_size exceeds u32::MAX"));
}
if !slot_size.is_multiple_of(slot::CACHE_LINE) {
return Err(Error::config(
"slot_size must be a multiple of CACHE_LINE (64)",
));
}
slot_count
.checked_mul(slot_size)
.and_then(|n| n.checked_add(SLOTS_OFFSET))
.ok_or_else(|| Error::config("region size overflows usize"))?;
Ok(())
}
fn slot_mut(&mut self, index: usize) -> &mut [u8] {
let offset = SLOTS_OFFSET + index * self.slot_size;
unsafe { core::slice::from_raw_parts_mut(self.ptr.as_ptr().add(offset), self.slot_size) }
}
fn slot_ref(&self, index: usize) -> &[u8] {
let offset = SLOTS_OFFSET + index * self.slot_size;
unsafe { core::slice::from_raw_parts(self.ptr.as_ptr().add(offset), self.slot_size) }
}
pub fn push(&mut self, fill: impl FnOnce(&mut [u8]) -> Result<()>) -> Result<bool> {
match self.ring.try_push() {
None => Ok(false),
Some(idx) => {
let buf = self.slot_mut(idx);
fill(buf)?;
self.ring.commit_push();
Ok(true)
}
}
}
pub fn pop(&mut self, consume: impl FnOnce(&[u8])) -> bool {
match self.ring.try_pop() {
None => false,
Some(idx) => {
let buf = self.slot_ref(idx);
consume(buf);
self.ring.commit_pop();
true
}
}
}
pub fn slot_count(&self) -> usize {
self.slot_count
}
pub fn slot_size(&self) -> usize {
self.slot_size
}
pub fn len(&self) -> usize {
self.ring.len()
}
pub fn is_empty(&self) -> bool {
self.ring.is_empty()
}
pub fn is_full(&self) -> bool {
self.ring.is_full()
}
}
impl Drop for Region {
fn drop(&mut self) {
unsafe {
libc::munmap(self.ptr.as_ptr() as *mut libc::c_void, self.len);
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::memory::slot::{self as s, FLAG_VALID, HEADER_SIZE, SlotHeader};
fn default_region() -> Region {
Region::anonymous(8, s::slot_size(64).unwrap()).unwrap()
}
fn fill_slot(seq: u64) -> impl FnOnce(&mut [u8]) -> Result<()> {
move |buf| {
let header = SlotHeader::new(1, FLAG_VALID, seq, 0, 4);
s::encode(&header, b"data", buf)
}
}
#[test]
fn create_and_drop() {
let region = default_region();
assert_eq!(region.slot_count(), 8);
assert!(region.is_empty());
}
#[test]
fn invalid_slot_count_rejected() {
assert!(Region::anonymous(0, 64).is_err());
assert!(Region::anonymous(3, 64).is_err()); }
#[test]
fn invalid_slot_size_rejected() {
assert!(Region::anonymous(4, HEADER_SIZE - 1).is_err());
assert!(Region::anonymous(4, 48).is_err()); }
#[test]
fn push_pop_roundtrip() {
let mut region = default_region();
let pushed = region.push(fill_slot(99)).unwrap();
assert!(pushed);
assert_eq!(region.len(), 1);
let mut seq = 0u64;
let popped = region.pop(|buf| {
let (hdr, _payload) = s::decode(buf).unwrap();
seq = hdr.sequence;
});
assert!(popped);
assert_eq!(seq, 99);
assert!(region.is_empty());
}
#[test]
fn full_ring_push_returns_false() {
let mut region = default_region();
for i in 0..8 {
assert!(region.push(fill_slot(i)).unwrap());
}
assert!(region.is_full());
assert!(!region.push(fill_slot(99)).unwrap());
}
#[test]
fn empty_ring_pop_returns_false() {
let mut region = default_region();
assert!(!region.pop(|_| {}));
}
#[test]
fn wrap_around_preserves_order() {
let mut region = Region::anonymous(4, s::slot_size(64).unwrap()).unwrap();
let mut results = Vec::new();
for round in 0..3u64 {
for i in 0..4u64 {
region.push(fill_slot(round * 4 + i)).unwrap();
}
for _ in 0..4 {
region.pop(|buf| {
let (hdr, _) = s::decode(buf).unwrap();
results.push(hdr.sequence);
});
}
}
let expected: Vec<u64> = (0..12).collect();
assert_eq!(results, expected);
}
#[test]
fn mmap_lifecycle() {
{
let mut region = Region::anonymous(4, 64).unwrap();
region.push(fill_slot(1)).unwrap();
} }
}