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}