strat9_kernel/ipc/
mailbox.rs1use alloc::boxed::Box;
16use core::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
17
18const FREELIST_CAPACITY: usize = 32;
21
22use super::transport::{
23 IpcConsumer, IpcError, IpcProducer, IpcTransport, TransportCapabilities, TransportLevel,
24};
25use crate::ipc::message::IpcMessage;
26
27#[cfg(target_arch = "x86_64")]
32const TAG_SHIFT: usize = 48;
33#[cfg(target_arch = "x86_64")]
34const TAG_MASK: usize = 0xFFFF_0000_0000_0000;
35#[cfg(target_arch = "x86_64")]
36const PTR_MASK: usize = !TAG_MASK;
37
38#[cfg(target_arch = "riscv64")]
40const TAG_SHIFT: usize = 56;
41#[cfg(target_arch = "riscv64")]
42const TAG_MASK: usize = 0xFF00_0000_0000_0000;
43#[cfg(target_arch = "riscv64")]
44const PTR_MASK: usize = !TAG_MASK;
45
46static TAG_COUNTER: AtomicUsize = AtomicUsize::new(0);
47
48fn tag_ptr(ptr: usize) -> usize {
50 let tag = TAG_COUNTER.fetch_add(1, Ordering::Relaxed) & 0xFFFF;
51 (ptr & PTR_MASK) | (tag << TAG_SHIFT)
52}
53
54fn untag_ptr(tagged: usize) -> *mut MailboxMessage {
56 (tagged & PTR_MASK) as *mut MailboxMessage
57}
58
59#[derive(Debug, Clone, Copy, PartialEq, Eq)]
65pub enum MailboxError {
66 AllocFailed,
68}
69
70#[repr(C)]
76pub struct MailboxMessage {
77 next: AtomicUsize,
79 pub data: IpcMessage,
81}
82
83#[derive(Debug)]
108struct NodePool {
109 head: AtomicU64,
110}
111
112const POOL_PTR_MASK: u64 = 0x0000_FFFF_FFFF_FFFF;
114const POOL_GEN_SHIFT: u64 = 48;
116
117impl NodePool {
118 const fn new() -> Self {
119 NodePool {
120 head: AtomicU64::new(0),
121 }
122 }
123
124 fn preallocate(&self, count: usize) {
126 for _ in 0..count {
127 let msg = MailboxMessage {
128 next: AtomicUsize::new(0),
129 data: IpcMessage::new(0),
130 };
131 let ptr = Box::into_raw(Box::new(msg)) as u64;
132 self.push_raw(ptr as *mut MailboxMessage);
133 }
134 }
135
136 fn try_pop_raw(&self) -> Option<*mut MailboxMessage> {
138 loop {
139 let tagged = self.head.load(Ordering::Acquire);
140 let ptr = (tagged & POOL_PTR_MASK) as *mut MailboxMessage;
141 if ptr.is_null() {
142 return None;
143 }
144 let gen = tagged >> POOL_GEN_SHIFT;
145 let next = unsafe { (*ptr).next.load(Ordering::Relaxed) } as u64;
146 let new_tagged = (next & POOL_PTR_MASK) | ((gen + 1) << POOL_GEN_SHIFT);
147 if self
148 .head
149 .compare_exchange_weak(tagged, new_tagged, Ordering::Acquire, Ordering::Relaxed)
150 .is_ok()
151 {
152 return Some(ptr);
153 }
154 }
155 }
156
157 fn push_raw(&self, ptr: *mut MailboxMessage) {
159 loop {
160 let tagged = self.head.load(Ordering::Relaxed);
161 let gen = tagged >> POOL_GEN_SHIFT;
162 unsafe {
163 (*ptr)
164 .next
165 .store((tagged & POOL_PTR_MASK) as usize, Ordering::Relaxed);
166 }
167 let new_tagged = (ptr as u64 & POOL_PTR_MASK) | ((gen + 1) << POOL_GEN_SHIFT);
168 if self
169 .head
170 .compare_exchange_weak(tagged, new_tagged, Ordering::Release, Ordering::Relaxed)
171 .is_ok()
172 {
173 return;
174 }
175 }
176 }
177}
178
179#[derive(Debug)]
180pub struct IntrusiveMailbox {
181 head: AtomicUsize,
182 pool: NodePool,
184}
185
186impl IntrusiveMailbox {
187 pub fn new() -> Self {
189 let mb = IntrusiveMailbox {
190 head: AtomicUsize::new(0),
191 pool: NodePool::new(),
192 };
193 mb.pool.preallocate(FREELIST_CAPACITY);
195 mb
196 }
197
198 pub const fn new_empty() -> Self {
202 IntrusiveMailbox {
203 head: AtomicUsize::new(0),
204 pool: NodePool::new(),
205 }
206 }
207
208 pub fn preallocate_nodes(&self, count: usize) {
210 self.pool.preallocate(count);
211 }
212
213 pub fn push(&self, msg: &[u8]) -> Result<(), MailboxError> {
219 let node_ptr = if let Some(ptr) = self.pool.try_pop_raw() {
221 let node = unsafe { &mut *ptr };
223 let len = msg.len().min(256);
224 node.data = IpcMessage::new(0);
225 node.data.payload[..len].copy_from_slice(&msg[..len]);
226 ptr as usize
227 } else {
228 let node = MailboxMessage::try_from_slice(msg).ok_or(MailboxError::AllocFailed)?;
230 Box::into_raw(Box::new(node)) as usize
231 };
232
233 loop {
234 let current = self.head.load(Ordering::Acquire);
235 unsafe {
236 (*(node_ptr as *mut MailboxMessage))
237 .next
238 .store(current & PTR_MASK, Ordering::Relaxed);
239 }
240 let new_tagged = tag_ptr(node_ptr);
241 if self
242 .head
243 .compare_exchange_weak(current, new_tagged, Ordering::Release, Ordering::Relaxed)
244 .is_ok()
245 {
246 return Ok(());
247 }
248 }
249 }
250
251 pub fn pop(&self) -> Option<IpcMessage> {
257 loop {
258 let current = self.head.load(Ordering::Acquire);
259 if current & PTR_MASK == 0 {
260 return None;
261 }
262 let current_ptr = untag_ptr(current);
263 let next = unsafe { (*current_ptr).next.load(Ordering::Relaxed) };
264 let new_tagged = tag_ptr(next);
265 if self
266 .head
267 .compare_exchange_weak(current, new_tagged, Ordering::Acquire, Ordering::Relaxed)
268 .is_ok()
269 {
270 let msg = unsafe { (*current_ptr).data };
274 self.pool.push_raw(current_ptr);
276 return Some(msg);
277 }
278 }
279 }
280
281 pub fn is_empty(&self) -> bool {
283 self.head.load(Ordering::Relaxed) & PTR_MASK == 0
284 }
285}
286
287impl MailboxMessage {
288 fn try_from_slice(data: &[u8]) -> Option<MailboxMessage> {
290 let len = data.len().min(256);
291 let mut msg = MailboxMessage {
292 next: AtomicUsize::new(0),
293 data: IpcMessage::new(0),
294 };
295 msg.data.payload[..len].copy_from_slice(&data[..len]);
296 Some(msg)
297 }
298}
299
300impl IpcTransport for IntrusiveMailbox {
305 fn level(&self) -> TransportLevel {
306 TransportLevel::TypeSafe
307 }
308
309 fn capabilities(&self) -> TransportCapabilities {
310 TransportCapabilities {
311 max_message_size: 256,
312 blocking: false,
313 zero_copy: false,
314 vectored: false,
315 directions: 1,
316 estimated_cost_cycles: 10,
317 }
318 }
319
320 fn name(&self) -> &'static str {
321 "mailbox"
322 }
323}
324
325impl IpcProducer for IntrusiveMailbox {
326 fn send(&self, msg: &[u8]) -> Result<(), IpcError> {
327 self.push(msg).map_err(|_| IpcError::TransportFailed)
328 }
329
330 fn try_send(&self, msg: &[u8]) -> Result<(), IpcError> {
331 self.send(msg)
332 }
333}
334
335impl IpcConsumer for IntrusiveMailbox {
336 fn recv(&self, buf: &mut [u8]) -> Result<usize, IpcError> {
337 match self.pop() {
338 Some(msg) => {
339 let len = msg.payload.len().min(buf.len());
340 buf[..len].copy_from_slice(&msg.payload[..len]);
341 Ok(len)
342 }
343 None => Err(IpcError::WouldBlock),
344 }
345 }
346
347 fn try_recv(&self, buf: &mut [u8]) -> Result<Option<usize>, IpcError> {
348 match self.recv(buf) {
349 Ok(n) => Ok(Some(n)),
350 Err(IpcError::WouldBlock) => Ok(None),
351 Err(e) => Err(e),
352 }
353 }
354}
355
356#[cfg(test)]
361mod tests {
362 use super::*;
363
364 #[test]
365 fn push_pop_single() {
366 let mb = IntrusiveMailbox::new();
367 mb.push(b"hello").unwrap();
368 let msg = mb.pop().unwrap();
369 assert_eq!(&msg.payload[..5], b"hello");
370 }
371
372 #[test]
373 fn push_pop_lifo_order() {
374 let mb = IntrusiveMailbox::new();
375 mb.push(b"first").unwrap();
376 mb.push(b"second").unwrap();
377 let msg2 = mb.pop().unwrap();
379 assert_eq!(&msg2.payload[..6], b"second");
380 let msg1 = mb.pop().unwrap();
381 assert_eq!(&msg1.payload[..5], b"first");
382 }
383
384 #[test]
385 fn pop_empty() {
386 let mb = IntrusiveMailbox::new();
387 assert!(mb.pop().is_none());
388 }
389}