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. Creation is
15//! exclusive, so concurrent first touches resolve to one creator and the rest
16//! attach.
17//! - Each segment is single-producer, single-consumer: at most one handle sends
18//! and one receives (a handle may do both). A second sender or receiver on the
19//! same name gets [`TransportError::Closed`]; use distinct names per direction
20//! for two-way traffic.
21//! - This is deliberately not registered in [`crate::TransportManager`]: it would
22//! collide with `InMemoryTransport` on `Address::Local`. Construct and use it
23//! directly when shared-memory IPC is wanted.
24
25use std::collections::HashMap;
26use std::sync::Mutex;
27
28use moirai_core::ipc::{IpcError, SendError, SharedQueue};
29
30use crate::{Address, Transport, TransportError, TransportResult};
31
32/// Per-message payload capacity. The backing frame is `4 + IPC_FRAME_DATA`
33/// bytes; keeping the total a 4 KiB multiple (and a multiple of the frame
34/// alignment) makes [`IpcFrame`] padding-free POD.
35pub const IPC_FRAME_DATA: usize = 4096 - core::mem::size_of::<u32>();
36
37/// Number of in-flight frames a segment's ring can hold before `send` reports
38/// [`TransportError::Full`].
39const IPC_QUEUE_CAPACITY: usize = 64;
40
41/// Yields an attach spends waiting for a concurrent creator to finish
42/// initialising a segment before reporting it unreachable.
43const ATTACH_ATTEMPTS: u32 = 1024;
44
45/// A fixed-size shared-memory frame: a length-prefixed byte payload.
46#[repr(C)]
47#[derive(Clone, Copy)]
48struct IpcFrame {
49 /// Number of valid bytes in `data` (always `<= IPC_FRAME_DATA`).
50 len: u32,
51 data: [u8; IPC_FRAME_DATA],
52}
53
54// SAFETY: `IpcFrame` is `#[repr(C)]` over a `u32` followed by `[u8; IPC_FRAME_DATA]`
55// where `IPC_FRAME_DATA` is a multiple of 4, so the struct has size `4 +
56// IPC_FRAME_DATA` (a multiple of its 4-byte alignment) with no padding, and every
57// bit pattern is a valid value — exactly the `Zeroable`/`Pod` contract. This is
58// required because `SharedQueue<T>` writes/reads `T` across the process boundary.
59unsafe impl bytemuck::Zeroable for IpcFrame {}
60unsafe impl bytemuck::Pod for IpcFrame {}
61
62/// Shared-memory IPC transport. Holds one `SharedQueue` handle per segment name,
63/// created lazily on first use.
64pub struct IpcTransport {
65 // `SharedQueue::{send,recv}` need `&mut self`, and the `Transport` trait is
66 // `&self`, so the per-segment handles live behind a `Mutex`. `Mutex<HashMap<
67 // _, SharedQueue<IpcFrame>>>` is `Send + Sync` (SharedQueue is `Send`), so
68 // `IpcTransport` satisfies the `Transport: Send + Sync` bound.
69 segments: Mutex<HashMap<String, SharedQueue<IpcFrame>>>,
70}
71
72impl IpcTransport {
73 /// Create an IPC transport with no open segments.
74 #[must_use]
75 pub fn new() -> Self {
76 Self {
77 segments: Mutex::new(HashMap::new()),
78 }
79 }
80
81 /// Attach to the segment `name`, creating it when no live segment holds the
82 /// name.
83 ///
84 /// `create` is exclusive, so two parties first-touching one name cannot both
85 /// create it: the loser reports `AlreadyExists` and retries `open`, which
86 /// succeeds once the winner finishes initialising the header. That window
87 /// spans the winner's `ftruncate` and header stores, so the retry is bounded
88 /// by `ATTACH_ATTEMPTS` yields rather than a wait.
89 fn attach(name: &str) -> TransportResult<SharedQueue<IpcFrame>> {
90 for _ in 0..ATTACH_ATTEMPTS {
91 if let Ok(queue) = SharedQueue::open(name, IPC_QUEUE_CAPACITY) {
92 return Ok(queue);
93 }
94 match SharedQueue::create(name, IPC_QUEUE_CAPACITY) {
95 Ok(queue) => return Ok(queue),
96 Err(IpcError::AlreadyExists) => std::thread::yield_now(),
97 Err(_) => return Err(TransportError::Closed),
98 }
99 }
100 Err(TransportError::Closed)
101 }
102
103 /// Borrow (attaching or creating on first use) the segment for `name`.
104 fn segment<'a>(
105 segments: &'a mut HashMap<String, SharedQueue<IpcFrame>>,
106 name: &str,
107 ) -> TransportResult<&'a mut SharedQueue<IpcFrame>> {
108 if !segments.contains_key(name) {
109 let queue = Self::attach(name)?;
110 segments.insert(name.to_string(), queue);
111 }
112 // Just inserted or already present.
113 Ok(segments
114 .get_mut(name)
115 .expect("invariant: segment present after insert"))
116 }
117}
118
119impl Default for IpcTransport {
120 fn default() -> Self {
121 Self::new()
122 }
123}
124
125impl Transport for IpcTransport {
126 fn send(&self, target: &Address, data: Vec<u8>) -> TransportResult<()> {
127 let Address::Local(name) = target else {
128 return Err(TransportError::Closed);
129 };
130 // Bound the payload to one frame; oversized messages are backpressure-class
131 // failures the caller must chunk around.
132 let len = u32::try_from(data.len())
133 .ok()
134 .filter(|&n| (n as usize) <= IPC_FRAME_DATA)
135 .ok_or(TransportError::Full)?;
136
137 let mut frame = IpcFrame {
138 len,
139 data: [0u8; IPC_FRAME_DATA],
140 };
141 frame.data[..data.len()].copy_from_slice(&data);
142
143 let mut segments = crate::lock_mutex(&self.segments);
144 let queue = Self::segment(&mut segments, name)?;
145 // SharedQueue::send returns the value back on a full ring.
146 queue.send(frame).map_err(|error| match error {
147 SendError::Full(_) => TransportError::Full,
148 SendError::Closed(_) | SendError::EndpointInUse(_) => TransportError::Closed,
149 })
150 }
151
152 fn recv(&self, source: &Address) -> TransportResult<Vec<u8>> {
153 let Address::Local(name) = source else {
154 return Err(TransportError::Closed);
155 };
156 let mut segments = crate::lock_mutex(&self.segments);
157 let queue = Self::segment(&mut segments, name)?;
158 match queue.recv().map_err(|_| TransportError::Closed)? {
159 Some(frame) => {
160 let len = frame.len as usize;
161 // Guard against a corrupt/hostile length from shared memory.
162 if len > IPC_FRAME_DATA {
163 return Err(TransportError::Closed);
164 }
165 Ok(frame.data[..len].to_vec())
166 }
167 None => Err(TransportError::Empty),
168 }
169 }
170
171 fn supports(&self, address: &Address) -> bool {
172 matches!(address, Address::Local(_))
173 }
174}
175
176#[cfg(test)]
177mod tests {
178 use super::*;
179
180 #[test]
181 fn ipc_transport_round_trips_through_shared_memory() {
182 let ipc = IpcTransport::new();
183 let addr = Address::Local("/moirai_ipc_transport_roundtrip".to_string());
184
185 ipc.send(&addr, b"hello ipc".to_vec()).expect("send");
186 ipc.send(&addr, b"second".to_vec()).expect("send");
187
188 // FIFO delivery across the shared-memory ring.
189 assert_eq!(ipc.recv(&addr).expect("recv"), b"hello ipc");
190 assert_eq!(ipc.recv(&addr).expect("recv"), b"second");
191 assert_eq!(ipc.recv(&addr), Err(TransportError::Empty));
192 }
193
194 #[test]
195 fn ipc_transport_attaches_across_separate_handles() {
196 // Two independent transports (as two processes would) sharing one segment:
197 // the first creates it, the second attaches.
198 let sender = IpcTransport::new();
199 let receiver = IpcTransport::new();
200 let addr = Address::Local("/moirai_ipc_transport_attach".to_string());
201
202 sender.send(&addr, b"cross-handle".to_vec()).expect("send");
203 assert_eq!(receiver.recv(&addr).expect("recv"), b"cross-handle");
204 }
205
206 #[test]
207 fn ipc_transport_rejects_oversized_and_non_local() {
208 let ipc = IpcTransport::new();
209 let addr = Address::Local("/moirai_ipc_transport_oversized".to_string());
210
211 let too_big = vec![0u8; IPC_FRAME_DATA + 1];
212 assert_eq!(ipc.send(&addr, too_big), Err(TransportError::Full));
213
214 let remote = Address::Remote(crate::RemoteAddress {
215 host: "127.0.0.1".to_string(),
216 port: 1,
217 service: "x".to_string(),
218 });
219 assert_eq!(ipc.send(&remote, vec![1]), Err(TransportError::Closed));
220 assert!(!ipc.supports(&remote));
221 assert!(ipc.supports(&addr));
222 }
223}