Skip to main content

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}