#![deny(clippy::indexing_slicing, clippy::arithmetic_side_effects)]
use super::error::IpcError;
use super::memory::SharedMemory;
use core::mem;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
pub struct SharedQueue<T> {
#[allow(dead_code)]
memory: SharedMemory,
meta: *mut QueueMetadata,
buffer: *mut T,
capacity: usize,
holds_sender: bool,
holds_receiver: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SendError<T> {
Full(T),
Closed(T),
EndpointInUse(T),
}
unsafe impl<T: Send> Send for SharedQueue<T> {}
const HEADER_ALIGN: usize = 64;
#[repr(C, align(64))]
struct QueueMetadata {
head: AtomicUsize,
capacity: AtomicUsize,
_pad1: [u8; 48],
tail: AtomicUsize,
_pad2: [u8; 56],
closed: AtomicBool,
sender_claimed: AtomicBool,
receiver_claimed: AtomicBool,
_pad3: [u8; 61],
}
pub(crate) const QUEUE_META_SIZE: usize = mem::size_of::<QueueMetadata>();
const _: () = assert!(QUEUE_META_SIZE == 3 * HEADER_ALIGN);
pub(crate) fn layout_total(
meta_size: usize,
elem_size: usize,
elem_align: usize,
elem_count: usize,
) -> Result<usize, IpcError> {
if elem_count == 0 || elem_align > HEADER_ALIGN || meta_size == 0 || elem_size == 0 {
return Err(IpcError::InvalidArgument);
}
elem_count
.checked_mul(elem_size)
.and_then(|data| data.checked_add(meta_size))
.ok_or(IpcError::InvalidArgument)
}
fn layout_for<T>(capacity: usize) -> Result<usize, IpcError> {
if capacity == 0 || mem::align_of::<T>() > HEADER_ALIGN {
return Err(IpcError::InvalidArgument);
}
layout_total(
QUEUE_META_SIZE,
mem::size_of::<T>(),
mem::align_of::<T>(),
capacity,
)
}
impl<T: bytemuck::Pod> SharedQueue<T> {
pub fn create(name: &str, capacity: usize) -> Result<Self, IpcError> {
let total_size = layout_for::<T>(capacity)?;
let memory = SharedMemory::create(name, total_size)?;
let meta = unsafe { &*header_of(&memory) };
meta.head.store(0, Ordering::Relaxed);
meta.tail.store(0, Ordering::Relaxed);
meta.closed.store(false, Ordering::Relaxed);
meta.sender_claimed.store(false, Ordering::Relaxed);
meta.receiver_claimed.store(false, Ordering::Relaxed);
meta.capacity.store(capacity, Ordering::Release);
Ok(Self::attach(memory, capacity))
}
pub fn open(name: &str, capacity: usize) -> Result<Self, IpcError> {
let total_size = layout_for::<T>(capacity)?;
let memory = SharedMemory::open(name, total_size)?;
let stored = unsafe { &*header_of(&memory) }
.capacity
.load(Ordering::Acquire);
if stored != capacity {
return Err(IpcError::InvalidArgument);
}
Ok(Self::attach(memory, capacity))
}
fn attach(memory: SharedMemory, capacity: usize) -> Self {
let meta = header_of(&memory);
let buffer = unsafe { memory.ptr.add(QUEUE_META_SIZE) }.cast::<T>();
Self {
memory,
meta,
buffer,
capacity,
holds_sender: false,
holds_receiver: false,
}
}
pub fn send(&mut self, value: T) -> Result<(), SendError<T>> {
if !self.holds_sender {
if !claim(unsafe { &(*self.meta).sender_claimed }) {
return Err(SendError::EndpointInUse(value));
}
self.holds_sender = true;
}
unsafe {
if (*self.meta).closed.load(Ordering::Relaxed) {
return Err(SendError::Closed(value));
}
let head = (*self.meta).head.load(Ordering::Relaxed);
let tail = (*self.meta).tail.load(Ordering::Acquire);
if head.wrapping_sub(tail) >= self.capacity {
return Err(SendError::Full(value));
}
#[expect(
clippy::arithmetic_side_effects,
reason = "capacity >= 1 is validated at create/open via layout_for"
)]
core::ptr::write(self.buffer.add(head % self.capacity), value);
(*self.meta)
.head
.store(head.wrapping_add(1), Ordering::Release);
Ok(())
}
}
pub fn recv(&mut self) -> Result<Option<T>, IpcError> {
if !self.holds_receiver {
if !claim(unsafe { &(*self.meta).receiver_claimed }) {
return Err(IpcError::EndpointInUse);
}
self.holds_receiver = true;
}
unsafe {
let tail = (*self.meta).tail.load(Ordering::Relaxed);
let head = (*self.meta).head.load(Ordering::Acquire);
if tail == head {
return Ok(None);
}
#[expect(
clippy::arithmetic_side_effects,
reason = "capacity >= 1 is validated at create/open via layout_for"
)]
let value = core::ptr::read(self.buffer.add(tail % self.capacity));
(*self.meta)
.tail
.store(tail.wrapping_add(1), Ordering::Release);
Ok(Some(value))
}
}
}
impl<T> Drop for SharedQueue<T> {
fn drop(&mut self) {
unsafe {
if self.holds_sender {
(*self.meta).sender_claimed.store(false, Ordering::Release);
}
if self.holds_receiver {
(*self.meta)
.receiver_claimed
.store(false, Ordering::Release);
}
}
}
}
#[expect(
clippy::cast_ptr_alignment,
reason = "mmap and MapViewOfFile return page-aligned bases"
)]
fn header_of(memory: &SharedMemory) -> *mut QueueMetadata {
memory.ptr.cast::<QueueMetadata>()
}
fn claim(flag: &AtomicBool) -> bool {
flag.compare_exchange(false, true, Ordering::AcqRel, Ordering::Relaxed)
.is_ok()
}