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. 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}