Skip to main content

moirai_core/ipc/
queue.rs

1#![deny(clippy::indexing_slicing, clippy::arithmetic_side_effects)]
2
3use super::error::IpcError;
4use super::memory::SharedMemory;
5use core::mem;
6use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
7
8/// Lock-free single-producer, single-consumer queue in shared memory.
9///
10/// Each end is exclusive: the first `send` on a handle claims the queue's sender
11/// endpoint and the first `recv` its receiver endpoint, through flags in the
12/// shared header, so two handles -- in one process or several -- can never both
13/// send or both receive. A handle may hold both ends. Claims release when the
14/// handle drops; a process killed while holding one leaves it held until the
15/// creator drops the segment.
16pub struct SharedQueue<T> {
17    #[allow(dead_code)]
18    memory: SharedMemory,
19    /// Queue metadata (stored at beginning of shared memory)
20    meta: *mut QueueMetadata,
21    /// Data buffer
22    buffer: *mut T,
23    /// Capacity
24    capacity: usize,
25    /// Whether this handle holds the queue's sender endpoint
26    holds_sender: bool,
27    /// Whether this handle holds the queue's receiver endpoint
28    holds_receiver: bool,
29}
30
31/// Why [`SharedQueue::send`] did not enqueue; each variant returns the value.
32#[derive(Debug, Clone, Copy, PartialEq, Eq)]
33pub enum SendError<T> {
34    /// The ring holds `capacity` unreceived values.
35    Full(T),
36    /// The queue was closed.
37    Closed(T),
38    /// Another handle, in this process or another, already sends on this queue.
39    EndpointInUse(T),
40}
41
42// SAFETY: queue contents move between threads and processes as plain `Pod`
43// bits, so `T: Send` is required; no references into shared memory escape.
44unsafe impl<T: Send> Send for SharedQueue<T> {}
45
46/// Alignment of the metadata header. The data buffer begins at
47/// `ptr + size_of::<QueueMetadata>()`, which — because the OS maps page-aligned
48/// memory and the header is a multiple of 64 — is 64-byte aligned. Element types
49/// whose alignment exceeds this would be placed at a misaligned address, so they
50/// are rejected at construction.
51const HEADER_ALIGN: usize = 64;
52
53#[repr(C, align(64))]
54struct QueueMetadata {
55    /// Producer position, cache-line aligned
56    head: AtomicUsize,
57    /// Element capacity recorded by the creator, validated by every `open` so a
58    /// peer cannot map a differently-sized view of the same segment.
59    capacity: AtomicUsize,
60    /// Padding to isolate the producer line (8 + 8 + 48 = 64 bytes)
61    _pad1: [u8; 48],
62    /// Consumer position, cache-line aligned
63    tail: AtomicUsize,
64    /// Padding to isolate tail and closed flag (8 + 56 = 64 bytes)
65    _pad2: [u8; 56],
66    /// Queue closed flag
67    closed: AtomicBool,
68    /// Set while one handle holds the sender endpoint
69    sender_claimed: AtomicBool,
70    /// Set while one handle holds the receiver endpoint
71    receiver_claimed: AtomicBool,
72    /// Padding to align the entire structure to 64 bytes (3 + 61 = 64 bytes)
73    _pad3: [u8; 61],
74}
75
76/// Header size in bytes; the capacity field sits right after the producer
77/// position (`head`) at this offset.
78pub(crate) const QUEUE_META_SIZE: usize = mem::size_of::<QueueMetadata>();
79
80const _: () = assert!(QUEUE_META_SIZE == 3 * HEADER_ALIGN);
81
82/// Pure layout arithmetic behind [`layout_for`]: total mapping size for
83/// `meta_size` header bytes plus `elem_count * elem_size`, rejecting zero
84/// count and overflow. Split out so the fuzz targets can exercise the exact
85/// arithmetic `create`/`open` rely on without OS resources.
86pub(crate) fn layout_total(
87    meta_size: usize,
88    elem_size: usize,
89    elem_align: usize,
90    elem_count: usize,
91) -> Result<usize, IpcError> {
92    if elem_count == 0 || elem_align > HEADER_ALIGN || meta_size == 0 || elem_size == 0 {
93        return Err(IpcError::InvalidArgument);
94    }
95    elem_count
96        .checked_mul(elem_size)
97        .and_then(|data| data.checked_add(meta_size))
98        .ok_or(IpcError::InvalidArgument)
99}
100
101/// Compute the total mapping size for `capacity` elements of `T`, rejecting a
102/// zero capacity (`% capacity` would divide by zero) and any size-overflow
103/// (which would otherwise produce an undersized mapping and out-of-bounds
104/// element access). Also rejects over-aligned element types.
105fn layout_for<T>(capacity: usize) -> Result<usize, IpcError> {
106    if capacity == 0 || mem::align_of::<T>() > HEADER_ALIGN {
107        return Err(IpcError::InvalidArgument);
108    }
109    layout_total(
110        QUEUE_META_SIZE,
111        mem::size_of::<T>(),
112        mem::align_of::<T>(),
113        capacity,
114    )
115}
116
117impl<T: bytemuck::Pod> SharedQueue<T> {
118    /// Create a new shared queue under a name no live segment holds.
119    ///
120    /// Fails with [`IpcError::AlreadyExists`] when the name is taken, so a
121    /// second creator can never reinitialise the header under live handles.
122    ///
123    /// `T` is bounded by [`bytemuck::Pod`]: shared-memory contents are written by
124    /// one process and read as `T` by another, so the element type must be valid
125    /// for every bit pattern (no `bool`/`char`/enum/reference discriminants a
126    /// peer could corrupt into an invalid value).
127    ///
128    /// # Errors
129    /// Returns [`IpcError::AlreadyExists`] for a taken name,
130    /// [`IpcError::InvalidArgument`] for a zero capacity, an over-aligned `T`,
131    /// or a size overflow, and the OS error otherwise.
132    pub fn create(name: &str, capacity: usize) -> Result<Self, IpcError> {
133        let total_size = layout_for::<T>(capacity)?;
134        let memory = SharedMemory::create(name, total_size)?;
135
136        // SAFETY: `memory.ptr` is the base of a fresh OS mapping (mmap /
137        // MapViewOfFile), always page-aligned and so satisfying
138        // `QueueMetadata`'s 64-byte alignment, and `total_size` covers the
139        // header. The header is only touched through atomics, which is sound
140        // against a peer that opens the name while this runs; `capacity` is
141        // stored last with `Release` so an opener that observes it sees the
142        // rest.
143        let meta = unsafe { &*header_of(&memory) };
144        meta.head.store(0, Ordering::Relaxed);
145        meta.tail.store(0, Ordering::Relaxed);
146        meta.closed.store(false, Ordering::Relaxed);
147        meta.sender_claimed.store(false, Ordering::Relaxed);
148        meta.receiver_claimed.store(false, Ordering::Relaxed);
149        meta.capacity.store(capacity, Ordering::Release);
150
151        Ok(Self::attach(memory, capacity))
152    }
153
154    /// Open an existing shared queue. Fails with [`IpcError::InvalidArgument`] if
155    /// the segment was created with a different capacity, which would otherwise
156    /// map a view inconsistent with the creator's and fault on access.
157    ///
158    /// # Errors
159    /// Returns [`IpcError::InvalidArgument`] for a capacity mismatch (including
160    /// a creator that has not finished initialising), and the OS error when the
161    /// segment is missing or too small.
162    pub fn open(name: &str, capacity: usize) -> Result<Self, IpcError> {
163        let total_size = layout_for::<T>(capacity)?;
164        let memory = SharedMemory::open(name, total_size)?;
165
166        // SAFETY: as in `create`; `SharedMemory::open` proved the object covers
167        // `total_size >= QUEUE_META_SIZE` bytes, and the recorded capacity is
168        // read atomically because the creator writes it while peers attach.
169        let stored = unsafe { &*header_of(&memory) }
170            .capacity
171            .load(Ordering::Acquire);
172        if stored != capacity {
173            return Err(IpcError::InvalidArgument);
174        }
175
176        Ok(Self::attach(memory, capacity))
177    }
178
179    fn attach(memory: SharedMemory, capacity: usize) -> Self {
180        let meta = header_of(&memory);
181        // SAFETY: `layout_for` sized the mapping as the header followed by
182        // `capacity` elements, so the offset stays inside it; the header is a
183        // multiple of `HEADER_ALIGN`, and `layout_for` rejected any `T`
184        // aligned more strictly, so the buffer is aligned for `T`.
185        let buffer = unsafe { memory.ptr.add(QUEUE_META_SIZE) }.cast::<T>();
186        Self {
187            memory,
188            meta,
189            buffer,
190            capacity,
191            holds_sender: false,
192            holds_receiver: false,
193        }
194    }
195
196    /// Send a value.
197    ///
198    /// The first call claims the queue's sender endpoint for this handle; the
199    /// claim is released when the handle drops.
200    ///
201    /// # Errors
202    /// Returns the value inside [`SendError::Full`] when the ring is full,
203    /// [`SendError::Closed`] when the queue is closed, and
204    /// [`SendError::EndpointInUse`] when another handle already sends on this
205    /// queue.
206    pub fn send(&mut self, value: T) -> Result<(), SendError<T>> {
207        if !self.holds_sender {
208            // SAFETY: `meta` points at the live mapping's header for the
209            // lifetime of `self`; only atomics are accessed through it.
210            if !claim(unsafe { &(*self.meta).sender_claimed }) {
211                return Err(SendError::EndpointInUse(value));
212            }
213            self.holds_sender = true;
214        }
215
216        // SAFETY: holding the sender claim makes this handle the queue's only
217        // sender, in this process and every other (the flag lives in the shared
218        // header and is won by one compare-exchange). The fullness check keeps
219        // the head slot outside the consumer window, and Pod writes need no
220        // drop coordination.
221        unsafe {
222            if (*self.meta).closed.load(Ordering::Relaxed) {
223                return Err(SendError::Closed(value));
224            }
225
226            let head = (*self.meta).head.load(Ordering::Relaxed);
227            let tail = (*self.meta).tail.load(Ordering::Acquire);
228
229            if head.wrapping_sub(tail) >= self.capacity {
230                return Err(SendError::Full(value));
231            }
232
233            // SAFETY-adjacent lint note: `capacity` is >= 1 by construction
234            // (`layout_for` rejects zero at create/open), so the modulo
235            // cannot panic.
236            #[expect(
237                clippy::arithmetic_side_effects,
238                reason = "capacity >= 1 is validated at create/open via layout_for"
239            )]
240            core::ptr::write(self.buffer.add(head % self.capacity), value);
241            (*self.meta)
242                .head
243                .store(head.wrapping_add(1), Ordering::Release);
244
245            Ok(())
246        }
247    }
248
249    /// Receive a value, or `None` when the queue is empty.
250    ///
251    /// The first call claims the queue's receiver endpoint for this handle; the
252    /// claim is released when the handle drops.
253    ///
254    /// # Errors
255    /// Returns [`IpcError::EndpointInUse`] when another handle already receives
256    /// on this queue.
257    pub fn recv(&mut self) -> Result<Option<T>, IpcError> {
258        if !self.holds_receiver {
259            // SAFETY: as in `send`.
260            if !claim(unsafe { &(*self.meta).receiver_claimed }) {
261                return Err(IpcError::EndpointInUse);
262            }
263            self.holds_receiver = true;
264        }
265
266        // SAFETY: holding the receiver claim makes this handle the queue's only
267        // receiver; the emptiness check guarantees the tail slot was published
268        // by the sender, and reading it as Pod bits moves it out exactly once.
269        unsafe {
270            let tail = (*self.meta).tail.load(Ordering::Relaxed);
271            let head = (*self.meta).head.load(Ordering::Acquire);
272
273            if tail == head {
274                return Ok(None);
275            }
276
277            #[expect(
278                clippy::arithmetic_side_effects,
279                reason = "capacity >= 1 is validated at create/open via layout_for"
280            )]
281            let value = core::ptr::read(self.buffer.add(tail % self.capacity));
282            (*self.meta)
283                .tail
284                .store(tail.wrapping_add(1), Ordering::Release);
285
286            Ok(Some(value))
287        }
288    }
289}
290
291impl<T> Drop for SharedQueue<T> {
292    fn drop(&mut self) {
293        // SAFETY: `meta` points into `self.memory`, which is unmapped only after
294        // this body returns. Releasing publishes this handle's final `head` or
295        // `tail` store to the next holder's acquiring claim.
296        unsafe {
297            if self.holds_sender {
298                (*self.meta).sender_claimed.store(false, Ordering::Release);
299            }
300            if self.holds_receiver {
301                (*self.meta)
302                    .receiver_claimed
303                    .store(false, Ordering::Release);
304            }
305        }
306    }
307}
308
309/// The header at the base of a mapping.
310///
311/// The OS maps page-aligned memory, so the base satisfies `QueueMetadata`'s
312/// 64-byte alignment even though `ptr` is typed as bytes.
313#[expect(
314    clippy::cast_ptr_alignment,
315    reason = "mmap and MapViewOfFile return page-aligned bases"
316)]
317fn header_of(memory: &SharedMemory) -> *mut QueueMetadata {
318    memory.ptr.cast::<QueueMetadata>()
319}
320
321/// Win an endpoint flag: exactly one caller across all handles and processes
322/// sees `true` until the holder releases it.
323fn claim(flag: &AtomicBool) -> bool {
324    flag.compare_exchange(false, true, Ordering::AcqRel, Ordering::Relaxed)
325        .is_ok()
326}