moirai_transport/ipc.rs
1//! Same-machine inter-process [`Transport`] over shared memory.
2//!
3//! Messages are carried as fixed-size [`IpcFrame`]s through a
4//! [`moirai_core::ipc::SharedQueue`] (a lock-free SPSC-style ring in a named
5//! shared-memory segment), one segment per [`Address::Local`] name. Unlike
6//! [`crate::InMemoryTransport`] — which routes within a single process via
7//! channels — this crosses the process boundary, so two processes mapping the
8//! same segment name exchange bytes directly.
9//!
10//! Scope and limits:
11//! - A single message is at most [`IPC_FRAME_DATA`] bytes; larger payloads are
12//! rejected with [`TransportError::Full`] (fragmentation is intentionally out of
13//! scope — callers chunk).
14//! - The first party to touch a segment creates it; others attach. Two processes
15//! first-touching the *same* segment concurrently is a creation race — in that
16//! case arrange for one side (typically the receiver) to create the segment
17//! before the other attaches, or use distinct names per direction.
18//! - This is deliberately not registered in [`crate::TransportManager`]: it would
19//! collide with `InMemoryTransport` on `Address::Local`. Construct and use it
20//! directly when shared-memory IPC is wanted.
21
22use std::collections::HashMap;
23use std::sync::Mutex;
24
25use moirai_core::ipc::SharedQueue;
26
27use crate::{Address, Transport, TransportError, TransportResult};
28
29/// Per-message payload capacity. The backing frame is `4 + IPC_FRAME_DATA`
30/// bytes; keeping the total a 4 KiB multiple (and a multiple of the frame
31/// alignment) makes [`IpcFrame`] padding-free POD.
32pub const IPC_FRAME_DATA: usize = 4096 - core::mem::size_of::<u32>();
33
34/// Number of in-flight frames a segment's ring can hold before `send` reports
35/// [`TransportError::Full`].
36const IPC_QUEUE_CAPACITY: usize = 64;
37
38/// A fixed-size shared-memory frame: a length-prefixed byte payload.
39#[repr(C)]
40#[derive(Clone, Copy)]
41struct IpcFrame {
42 /// Number of valid bytes in `data` (always `<= IPC_FRAME_DATA`).
43 len: u32,
44 data: [u8; IPC_FRAME_DATA],
45}
46
47// SAFETY: `IpcFrame` is `#[repr(C)]` over a `u32` followed by `[u8; IPC_FRAME_DATA]`
48// where `IPC_FRAME_DATA` is a multiple of 4, so the struct has size `4 +
49// IPC_FRAME_DATA` (a multiple of its 4-byte alignment) with no padding, and every
50// bit pattern is a valid value — exactly the `Zeroable`/`Pod` contract. This is
51// required because `SharedQueue<T>` writes/reads `T` across the process boundary.
52unsafe impl bytemuck::Zeroable for IpcFrame {}
53unsafe impl bytemuck::Pod for IpcFrame {}
54
55/// Shared-memory IPC transport. Holds one `SharedQueue` handle per segment name,
56/// created lazily on first use.
57pub struct IpcTransport {
58 // `SharedQueue::{send,recv}` need `&mut self`, and the `Transport` trait is
59 // `&self`, so the per-segment handles live behind a `Mutex`. `Mutex<HashMap<
60 // _, SharedQueue<IpcFrame>>>` is `Send + Sync` (SharedQueue is `Send`), so
61 // `IpcTransport` satisfies the `Transport: Send + Sync` bound.
62 segments: Mutex<HashMap<String, SharedQueue<IpcFrame>>>,
63}
64
65impl IpcTransport {
66 /// Create an IPC transport with no open segments.
67 #[must_use]
68 pub fn new() -> Self {
69 Self {
70 segments: Mutex::new(HashMap::new()),
71 }
72 }
73
74 /// Borrow (attaching or creating on first use) the segment for `name`.
75 fn segment<'a>(
76 segments: &'a mut HashMap<String, SharedQueue<IpcFrame>>,
77 name: &str,
78 ) -> TransportResult<&'a mut SharedQueue<IpcFrame>> {
79 if !segments.contains_key(name) {
80 // Attach to an existing segment, else create it. Capacity must match
81 // the creator's; every party uses IPC_QUEUE_CAPACITY so attach
82 // succeeds.
83 let queue = SharedQueue::open(name, IPC_QUEUE_CAPACITY)
84 .or_else(|_| SharedQueue::create(name, IPC_QUEUE_CAPACITY))
85 .map_err(|_| TransportError::Closed)?;
86 segments.insert(name.to_string(), queue);
87 }
88 // Just inserted or already present.
89 Ok(segments
90 .get_mut(name)
91 .expect("invariant: segment present after insert"))
92 }
93}
94
95impl Default for IpcTransport {
96 fn default() -> Self {
97 Self::new()
98 }
99}
100
101impl Transport for IpcTransport {
102 fn send(&self, target: &Address, data: Vec<u8>) -> TransportResult<()> {
103 let Address::Local(name) = target else {
104 return Err(TransportError::Closed);
105 };
106 // Bound the payload to one frame; oversized messages are backpressure-class
107 // failures the caller must chunk around.
108 let len = u32::try_from(data.len())
109 .ok()
110 .filter(|&n| (n as usize) <= IPC_FRAME_DATA)
111 .ok_or(TransportError::Full)?;
112
113 let mut frame = IpcFrame {
114 len,
115 data: [0u8; IPC_FRAME_DATA],
116 };
117 frame.data[..data.len()].copy_from_slice(&data);
118
119 let mut segments = crate::lock_mutex(&self.segments);
120 let queue = Self::segment(&mut segments, name)?;
121 // SharedQueue::send returns the value back on a full ring.
122 queue.send(frame).map_err(|_| TransportError::Full)
123 }
124
125 fn recv(&self, source: &Address) -> TransportResult<Vec<u8>> {
126 let Address::Local(name) = source else {
127 return Err(TransportError::Closed);
128 };
129 let mut segments = crate::lock_mutex(&self.segments);
130 let queue = Self::segment(&mut segments, name)?;
131 match queue.recv() {
132 Some(frame) => {
133 let len = frame.len as usize;
134 // Guard against a corrupt/hostile length from shared memory.
135 if len > IPC_FRAME_DATA {
136 return Err(TransportError::Closed);
137 }
138 Ok(frame.data[..len].to_vec())
139 }
140 None => Err(TransportError::Empty),
141 }
142 }
143
144 fn supports(&self, address: &Address) -> bool {
145 matches!(address, Address::Local(_))
146 }
147}
148
149#[cfg(test)]
150mod tests {
151 use super::*;
152
153 #[test]
154 fn ipc_transport_round_trips_through_shared_memory() {
155 let ipc = IpcTransport::new();
156 let addr = Address::Local("/moirai_ipc_transport_roundtrip".to_string());
157
158 ipc.send(&addr, b"hello ipc".to_vec()).expect("send");
159 ipc.send(&addr, b"second".to_vec()).expect("send");
160
161 // FIFO delivery across the shared-memory ring.
162 assert_eq!(ipc.recv(&addr).expect("recv"), b"hello ipc");
163 assert_eq!(ipc.recv(&addr).expect("recv"), b"second");
164 assert_eq!(ipc.recv(&addr), Err(TransportError::Empty));
165 }
166
167 #[test]
168 fn ipc_transport_attaches_across_separate_handles() {
169 // Two independent transports (as two processes would) sharing one segment:
170 // the first creates it, the second attaches.
171 let sender = IpcTransport::new();
172 let receiver = IpcTransport::new();
173 let addr = Address::Local("/moirai_ipc_transport_attach".to_string());
174
175 sender.send(&addr, b"cross-handle".to_vec()).expect("send");
176 assert_eq!(receiver.recv(&addr).expect("recv"), b"cross-handle");
177 }
178
179 #[test]
180 fn ipc_transport_rejects_oversized_and_non_local() {
181 let ipc = IpcTransport::new();
182 let addr = Address::Local("/moirai_ipc_transport_oversized".to_string());
183
184 let too_big = vec![0u8; IPC_FRAME_DATA + 1];
185 assert_eq!(ipc.send(&addr, too_big), Err(TransportError::Full));
186
187 let remote = Address::Remote(crate::RemoteAddress {
188 host: "127.0.0.1".to_string(),
189 port: 1,
190 service: "x".to_string(),
191 });
192 assert_eq!(ipc.send(&remote, vec![1]), Err(TransportError::Closed));
193 assert!(!ipc.supports(&remote));
194 assert!(ipc.supports(&addr));
195 }
196}