ringfire 0.5.2

Zero-copy lock-free inter-process communication (IPC) ring buffer and shared memory bus in Rust, with ring mirroring across hosts
Documentation
//! # PayloadArena & BlobRef
//!
//! Ultra-high-performance contiguous byte arena in shared memory for variable-sized
//! and large IPC messages (e.g. L2/L3 order books, network packets, serialized frames).
//!
//! Inspired by Firedancer's `dcache` architecture:
//! - Contiguous cyclic byte arena mapped alongside the ring buffer.
//! - Fixed-size ring buffer slots carry only lightweight metadata (`BlobRef`, 16 bytes).
//! - Guaranteed zero memory fragmentation: allocations that would wrap around the ring
//!   boundary are automatically aligned to index 0, ensuring every `BlobRef` maps
//!   to a strictly contiguous memory slice (`&[u8]`).

use std::sync::atomic::{fence, AtomicU64, Ordering};
use crate::error::{Result, RingfireError};

/// 16-byte descriptor referencing a payload stored in `PayloadArena`.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
#[repr(C)]
pub struct BlobRef {
    /// Global cumulative offset in the arena
    pub offset: u64,
    /// Exact payload length in bytes
    pub len: u32,
    /// Application-defined flags (e.g. codec, schema, compression tag)
    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
    }
}

/// Header for `PayloadArena` located at the start of the arena section.
/// 64-byte cache-line aligned.
#[repr(C, align(64))]
pub struct ArenaHeader {
    /// Total arena capacity in bytes (must be a power of two)
    pub capacity: u64,
    /// Bitmask for circular indexing (`capacity - 1`)
    pub mask: u64,
    /// Cumulative reserved byte counter
    pub reserved: AtomicU64,
    /// Padding to exactly 64 bytes
    pub _pad: [u8; 40],
}

const _: () = {
    assert!(std::mem::size_of::<ArenaHeader>() == 64);
    assert!(std::mem::align_of::<ArenaHeader>() == 64);
};

/// High-throughput shared memory byte arena.
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 {
    /// Initialize a new `PayloadArena` over pre-mapped shared memory.
    ///
    /// # Safety
    /// `ptr` must point to valid, writable memory of at least `size_of::<ArenaHeader>() + capacity` bytes.
    /// `capacity` must be a power of two and a multiple of 64.
    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,
        })
    }

    /// Open an existing `PayloadArena` from pre-mapped shared memory.
    ///
    /// # Safety
    /// `ptr` must point to an initialized `PayloadArena` region.
    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,
        })
    }

    /// Total capacity of the payload arena in bytes.
    #[inline]
    pub fn capacity(&self) -> usize {
        self.capacity
    }

    /// Atomically reserve `len` bytes in the arena.
    ///
    /// If the allocation would cross the circular buffer boundary, it skips
    /// to index 0, ensuring that every returned `BlobRef` maps to a contiguous slice.
    #[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,
            });
        }

        // Align allocation up to 64 bytes (cache line)
        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;

            // Check if contiguous slice fits before end of arena
            let (actual_offset, next) = if from_ix + len > self.capacity {
                // Wrap around: align to capacity boundary (index 0 on next lap)
                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()
            {
                // Readers detect overwritten blobs by re-reading `reserved` after copying:
                // the reservation must be visible before any byte of the new blob.
                fence(Ordering::Release);
                return Ok(BlobRef {
                    offset: actual_offset,
                    len: len as u32,
                    flags,
                });
            }
        }
    }

    /// Write payload bytes into the arena, returning a `BlobRef`.
    #[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)
    }

    /// Zero-copy write directly into reserved arena slice.
    #[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))
    }

    /// Copy payload bytes out of the arena into `out`.
    #[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)
    }

    /// Inspect payload in-place without copying.
    #[inline]
    ///
    /// # Panics
    /// If `blob_ref` does not lie inside the arena (see [`Self::is_in_bounds`]).
    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))
        }
    }

    /// Check if a `BlobRef` has been (or is being) overwritten by the producer wrapping around.
    ///
    /// To validate a copy, call this *after* reading the bytes: the acquire fence orders
    /// the copy before the re-read of the reservation counter.
    #[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
    }

    /// Whether `blob_ref` describes a slice that lies entirely inside the arena.
    /// A descriptor read from shared memory must pass this before it is dereferenced.
    #[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;