Skip to main content

orbit_core/ring/
shm.rs

1//! `ShmRing` — POSIX SHM-backed ring buffer (V1 substrate).
2//!
3//! The runtime substrate that finally makes Orbit's "fleet sees each
4//! other" promise real. Cross-process visibility is a memory-coherent
5//! `mmap` with shared atomics; no kernel round-trips per access.
6//!
7//! ## Layout (single SHM segment per `OrbitTyped` kind)
8//!
9//! ```text
10//! ┌──────────────────────────┬───────────────────────────┬──────────────┐
11//! │ ShmRingHeader (64B)      │ LaneHeader[0..L]          │ Lane slots   │
12//! │ kind/spec/topology       │ one head per writer lane  │ L × N slots  │
13//! └──────────────────────────┴───────────────────────────┴──────────────┘
14//! ```
15//!
16//! Shared topology retains the lock-free multi-writer claim-before-commit
17//! protocol. Shared-ordered topology serializes writers with a
18//! process-recoverable kernel advisory lock, commits the slot, and only then
19//! advances the head.
20//! Per-node topology assigns disjoint slots to every node and requires one
21//! active process per node id. Concurrent tasks inside that process are
22//! serialized locally, and the lane head is published only after the slot
23//! reaches its committed sequence. A process dying during a per-node write
24//! therefore leaves no visible hole for other lanes.
25//!
26//! ## Slot publication protocol
27//!
28//! 1. choose/claim the lane counter
29//! 2. `slot = slots[counter % capacity]`
30//! 3. `slot.seq = 2*counter + 1` (odd → writing)
31//! 4. fill `id`, `kind`, `ver`, `payload_len`, `payload[..len]`
32//! 5. `slot.seq = 2*counter + 2` (even → committed)
33//! 6. for per-node lanes, release-store `head = counter + 1`
34//!
35//! ## Read protocol
36//!
37//! 1. read `seq_pre = slot.seq` (Acquire)
38//! 2. if odd → writer in progress, return None / retry
39//! 3. read all fields into locals
40//! 4. read `seq_post = slot.seq` (Acquire)
41//! 5. if `seq_pre != seq_post` → torn write, retry
42//! 6. validate counter encoded in seq matches the one we want
43//!
44//! Tearing is detected, never silently merged. A reader that loses
45//! the race against the writer simply retries (or returns None for
46//! head-read).
47//!
48//! ## Per-kind payload size
49//!
50//! Every `OrbitTyped::KIND` declares its own `RingSpec`. The payload
51//! capacity determines that SHM segment's slot stride; unrelated lanes
52//! no longer pay for one global maximum. Capacity and payload capacity
53//! are persisted in the header and verified by every attaching process.
54
55#![cfg(unix)]
56
57use std::ptr;
58use std::sync::Mutex;
59use std::sync::atomic::{AtomicU8, AtomicU32, AtomicU64, Ordering, fence};
60
61use bytes::Bytes;
62
63use crate::NodeId;
64use crate::id::NetId64;
65use crate::ring::{Frame, RingSpec, RingTopology};
66use crate::shm::{self, ShmRegion};
67
68// ─────────────────────────────────────────────────────────────────────
69// On-disk (well, on-SHM) layout
70// ─────────────────────────────────────────────────────────────────────
71
72const MAGIC: u32 = 0x4F524254; // "ORBT" big-endian when displayed
73const VERSION: u32 = 3;
74const SLOT_ALIGNMENT: usize = 64;
75
76/// Cache-line aligned to keep the header on its own line.
77#[repr(C, align(64))]
78struct ShmRingHeader {
79    version_counter: AtomicU64,
80    capacity: u64,
81    magic: u32,
82    version: u32,
83    payload_capacity: u32,
84    slot_stride: u32,
85    kind: u8,
86    topology: u8,
87    lane_count: u16,
88    /// Shared futex/umtx word used by notification-enabled ring users.
89    /// It lives in the existing reserved header space, so the V3
90    /// layout remains compatible with already-created segments.
91    notification_generation: AtomicU32,
92    /// Explicitly fills one cache line; field order avoids implicit
93    /// alignment padding that would otherwise make this header 128B.
94    _reserved: [u8; 24],
95}
96
97#[repr(C, align(64))]
98struct ShmLaneHeader {
99    write_pos: AtomicU64,
100    _reserved: [u8; 56],
101}
102
103/// Fixed prefix of one dynamically-strided SHM slot.
104#[repr(C)]
105struct ShmSlotHeader {
106    /// LMAX-style sequence: odd while a writer is mid-flight, even
107    /// once committed. Reader checks seq before/after content read
108    /// to detect torn writes.
109    seq: AtomicU64,
110    /// `NetId64::raw()` of the frame that occupies this slot.
111    id: AtomicU64,
112    /// `Frame::ver` — caller-supplied version / tick.
113    ver: AtomicU64,
114    /// Length of the meaningful prefix of `payload`.
115    payload_len: AtomicU32,
116    /// `Frame::kind` — the message-class byte (state/event/cmd/…).
117    kind: AtomicU8,
118    _reserved: [u8; 3],
119}
120
121// The content fields are atomic for the same reason `seq` is: a reader may be
122// copying this slot while a writer overwrites it. The seqlock *detects* that
123// afterwards, which is not the same as not having a race — under the memory
124// model a torn read through a plain `u64` is undefined however the hardware
125// behaves, and Miri or TSan would say so. `Relaxed` is what they need: all the
126// ordering is already carried by `seq` and the two fences around it, so these
127// compile to the same plain loads and stores they were, with the race removed
128// rather than papered over.
129
130const SLOT_HEADER_SIZE: usize = std::mem::size_of::<ShmSlotHeader>();
131const HEADER_SIZE: usize = std::mem::size_of::<ShmRingHeader>();
132const LANE_HEADER_SIZE: usize = std::mem::size_of::<ShmLaneHeader>();
133
134const _: () = assert!(HEADER_SIZE == 64);
135const _: () = assert!(LANE_HEADER_SIZE == 64);
136const _: () = assert!(SLOT_HEADER_SIZE == 32);
137
138fn invalid_input(message: impl Into<String>) -> std::io::Error {
139    std::io::Error::new(std::io::ErrorKind::InvalidInput, message.into())
140}
141
142fn lane_count_for(spec: RingSpec, fleet_capacity: u16) -> std::io::Result<usize> {
143    if fleet_capacity == 0 {
144        return Err(invalid_input("ShmRing fleet capacity must be > 0"));
145    }
146    Ok(match spec.topology {
147        RingTopology::Shared | RingTopology::SharedOrdered => 1,
148        RingTopology::PerNode => usize::from(fleet_capacity),
149    })
150}
151
152fn checked_layout(spec: RingSpec, fleet_capacity: u16) -> std::io::Result<(usize, usize, usize)> {
153    if spec.capacity == 0 {
154        return Err(invalid_input("ShmRing capacity must be > 0"));
155    }
156    if !spec.capacity.is_power_of_two() {
157        return Err(invalid_input("ShmRing capacity must be a power of two"));
158    }
159    if spec.payload_capacity > u32::MAX as usize {
160        return Err(invalid_input("ShmRing payload capacity must fit in u32"));
161    }
162
163    let unaligned = SLOT_HEADER_SIZE
164        .checked_add(spec.payload_capacity)
165        .ok_or_else(|| invalid_input("ShmRing slot size overflow"))?;
166    let slot_stride = unaligned
167        .checked_add(SLOT_ALIGNMENT - 1)
168        .map(|value| value & !(SLOT_ALIGNMENT - 1))
169        .ok_or_else(|| invalid_input("ShmRing slot stride overflow"))?;
170    if slot_stride > u32::MAX as usize {
171        return Err(invalid_input("ShmRing slot stride must fit in u32"));
172    }
173    let lane_count = lane_count_for(spec, fleet_capacity)?;
174    let lane_headers_size = lane_count
175        .checked_mul(LANE_HEADER_SIZE)
176        .ok_or_else(|| invalid_input("ShmRing lane header size overflow"))?;
177    let slots_offset = HEADER_SIZE
178        .checked_add(lane_headers_size)
179        .ok_or_else(|| invalid_input("ShmRing slots offset overflow"))?;
180    let slots_per_lane = spec
181        .capacity
182        .checked_mul(slot_stride)
183        .ok_or_else(|| invalid_input("ShmRing lane size overflow"))?;
184    let slots_size = lane_count
185        .checked_mul(slots_per_lane)
186        .ok_or_else(|| invalid_input("ShmRing slots size overflow"))?;
187    let segment_size = slots_offset
188        .checked_add(slots_size)
189        .ok_or_else(|| invalid_input("ShmRing segment size overflow"))?;
190    Ok((slot_stride, slots_offset, segment_size))
191}
192
193/// Compute the SHM segment size required by a ring spec.
194pub fn segment_size_for_spec(spec: RingSpec) -> std::io::Result<usize> {
195    segment_size_for_spec_and_fleet(spec, 1)
196}
197
198/// Compute the SHM segment size required for a fleet-aware ring spec.
199pub fn segment_size_for_spec_and_fleet(
200    spec: RingSpec,
201    fleet_capacity: u16,
202) -> std::io::Result<usize> {
203    checked_layout(spec, fleet_capacity).map(|(_, _, segment_size)| segment_size)
204}
205
206// ─────────────────────────────────────────────────────────────────────
207// ShmRing
208// ─────────────────────────────────────────────────────────────────────
209
210/// Cross-process ring buffer backed by a POSIX SHM segment.
211///
212/// Use [`ShmRing::open_or_create`] to attach to (or create) the
213/// segment by name. Multiple processes calling this with the same name and
214/// matching [`RingSpec`] share the underlying memory; the first process to
215/// call it does the one-time header initialization.
216pub struct ShmRing {
217    region: ShmRegion,
218    kind: u8,
219    capacity: usize,
220    payload_capacity: usize,
221    topology: RingTopology,
222    lane_count: usize,
223    slot_stride: usize,
224    slots_offset: usize,
225    lane_stride: usize,
226    write_locks: Vec<Mutex<()>>,
227}
228
229impl ShmRing {
230    /// Open or create a SHM-backed ring under `fleet_name` for type
231    /// kind `kind` with `spec`. The first process to call
232    /// this initializes the header; later attachers reuse it.
233    pub fn open_or_create(fleet_name: &str, kind: u8, spec: RingSpec) -> std::io::Result<Self> {
234        Self::open_or_create_for_fleet(fleet_name, kind, spec, 1)
235    }
236
237    /// Open or create a SHM-backed ring using `fleet_capacity` physical
238    /// writer lanes when `spec` is [`RingTopology::PerNode`].
239    pub fn open_or_create_for_fleet(
240        fleet_name: &str,
241        kind: u8,
242        spec: RingSpec,
243        fleet_capacity: u16,
244    ) -> std::io::Result<Self> {
245        let lane_count = lane_count_for(spec, fleet_capacity)?;
246        let (slot_stride, slots_offset, size) = checked_layout(spec, fleet_capacity)?;
247        let lane_stride = spec
248            .capacity
249            .checked_mul(slot_stride)
250            .ok_or_else(|| invalid_input("ShmRing lane stride overflow"))?;
251        let name = shm::ring_segment_name(fleet_name, kind);
252
253        // Locked whatever the topology. Creation is two steps that a peer can
254        // arrive between — `shm_open(O_CREAT|O_EXCL)` publishes the name, and
255        // the header below is written after it — and the attach path reads that
256        // header and rejects a segment whose magic is still zero, with no
257        // retry. Only `SharedOrdered` took this lock, so every other ring was
258        // one scheduling accident away from a peer failing to boot with
259        // "wrong magic 0x00000000"; N workers starting together is exactly the
260        // case that produces it. `SharedOrdered` still keeps the lock for its
261        // writes — this is the same lock, held for a different reason.
262        let (region, _initialization_lock) = ShmRegion::open_or_create_locked(&name, size)?;
263
264        // Initialize header on first creation; subsequent attachers
265        // skip and rely on whatever the creator wrote.
266        if region.created() {
267            // SAFETY: region is mapped, header lives at offset 0,
268            // size is right because we just created with this size.
269            unsafe {
270                let header_ptr = region.as_ptr() as *mut ShmRingHeader;
271                ptr::write(
272                    header_ptr,
273                    ShmRingHeader {
274                        version_counter: AtomicU64::new(0),
275                        capacity: spec.capacity as u64,
276                        magic: MAGIC,
277                        version: VERSION,
278                        payload_capacity: spec.payload_capacity as u32,
279                        slot_stride: slot_stride as u32,
280                        kind,
281                        topology: spec.topology as u8,
282                        lane_count: lane_count as u16,
283                        notification_generation: AtomicU32::new(0),
284                        _reserved: [0; 24],
285                    },
286                );
287
288                for lane in 0..lane_count {
289                    let lane_ptr = region.as_ptr().add(HEADER_SIZE + lane * LANE_HEADER_SIZE)
290                        as *mut ShmLaneHeader;
291                    ptr::write(
292                        lane_ptr,
293                        ShmLaneHeader {
294                            write_pos: AtomicU64::new(0),
295                            _reserved: [0; 56],
296                        },
297                    );
298                }
299
300                // Zero out the slot region so all `seq` values start
301                // at 0 (== "never written, even, slot empty"). 0 is
302                // valid as both a u64 atomic and as bytes for our
303                // POD slot shape.
304                let slots_ptr = region.as_ptr().add(slots_offset);
305                ptr::write_bytes(slots_ptr, 0, lane_count * lane_stride);
306            }
307        } else {
308            // Attaching to an existing segment — sanity-check the header.
309            // SAFETY: region is mapped; header lives at offset 0.
310            let header = unsafe { &*(region.as_ptr() as *const ShmRingHeader) };
311            if header.magic != MAGIC {
312                return Err(std::io::Error::new(
313                    std::io::ErrorKind::InvalidData,
314                    format!(
315                        "SHM segment {} has wrong magic 0x{:08X} (expected 0x{:08X})",
316                        name, header.magic, MAGIC
317                    ),
318                ));
319            }
320            if header.version != VERSION {
321                return Err(std::io::Error::new(
322                    std::io::ErrorKind::InvalidData,
323                    format!(
324                        "SHM segment {} version {} != local {}",
325                        name, header.version, VERSION
326                    ),
327                ));
328            }
329            if header.kind != kind {
330                return Err(std::io::Error::new(
331                    std::io::ErrorKind::InvalidData,
332                    format!(
333                        "SHM segment {} kind {} != requested {}",
334                        name, header.kind, kind
335                    ),
336                ));
337            }
338            if header.topology != spec.topology as u8 {
339                return Err(std::io::Error::new(
340                    std::io::ErrorKind::InvalidData,
341                    format!(
342                        "SHM segment {} topology {} != requested {}",
343                        name, header.topology, spec.topology as u8
344                    ),
345                ));
346            }
347            if header.lane_count as usize != lane_count {
348                return Err(std::io::Error::new(
349                    std::io::ErrorKind::InvalidData,
350                    format!(
351                        "SHM segment {} lane count {} != requested {}",
352                        name, header.lane_count, lane_count
353                    ),
354                ));
355            }
356            if header.capacity as usize != spec.capacity {
357                return Err(std::io::Error::new(
358                    std::io::ErrorKind::InvalidData,
359                    format!(
360                        "SHM segment {} capacity {} != requested {}",
361                        name, header.capacity, spec.capacity
362                    ),
363                ));
364            }
365            if header.payload_capacity as usize != spec.payload_capacity {
366                return Err(std::io::Error::new(
367                    std::io::ErrorKind::InvalidData,
368                    format!(
369                        "SHM segment {} payload capacity {} != requested {}",
370                        name, header.payload_capacity, spec.payload_capacity
371                    ),
372                ));
373            }
374            if header.slot_stride as usize != slot_stride {
375                return Err(std::io::Error::new(
376                    std::io::ErrorKind::InvalidData,
377                    format!(
378                        "SHM segment {} slot stride {} != requested {}",
379                        name, header.slot_stride, slot_stride
380                    ),
381                ));
382            }
383        }
384
385        Ok(Self {
386            region,
387            kind,
388            capacity: spec.capacity,
389            payload_capacity: spec.payload_capacity,
390            topology: spec.topology,
391            lane_count,
392            slot_stride,
393            slots_offset,
394            lane_stride,
395            write_locks: (0..lane_count).map(|_| Mutex::new(())).collect(),
396        })
397    }
398
399    /// True when this handle was the one that *created* the SHM
400    /// segment (vs attaching to an already-existing one).
401    pub fn created(&self) -> bool {
402        self.region.created()
403    }
404
405    /// Remove the SHM segment name. Existing mappings stay valid;
406    /// new opens will fail until `open_or_create` recreates it.
407    /// Call from the owner process at fleet shutdown.
408    pub fn unlink(&self) -> std::io::Result<()> {
409        self.region.unlink()
410    }
411
412    /// The KIND byte this ring carries.
413    pub fn kind(&self) -> u8 {
414        self.kind
415    }
416
417    /// Slot count per lane.
418    pub fn capacity(&self) -> usize {
419        self.capacity
420    }
421
422    /// Number of physical writer lanes in this segment.
423    pub fn lane_count(&self) -> usize {
424        self.lane_count
425    }
426
427    /// Maximum inline payload bytes for this ring lane.
428    pub fn payload_capacity(&self) -> usize {
429        self.payload_capacity
430    }
431
432    pub fn spec(&self) -> RingSpec {
433        RingSpec {
434            capacity: self.capacity,
435            payload_capacity: self.payload_capacity,
436            topology: self.topology,
437        }
438    }
439
440    fn lane_header(&self, lane: usize) -> &ShmLaneHeader {
441        debug_assert!(lane < self.lane_count);
442        unsafe {
443            &*(self
444                .region
445                .as_ptr()
446                .add(HEADER_SIZE + lane * LANE_HEADER_SIZE) as *const ShmLaneHeader)
447        }
448    }
449
450    #[cfg(any(target_os = "linux", target_os = "freebsd", target_os = "macos"))]
451    pub(crate) fn notification_generation(&self) -> &AtomicU32 {
452        // SAFETY: the mapped header was initialized before this ring handle
453        // was returned and remains mapped for the lifetime of `self`.
454        unsafe { &(*(self.region.as_ptr() as *const ShmRingHeader)).notification_generation }
455    }
456
457    fn slot_ptr(&self, lane: usize, idx: usize) -> *mut ShmSlotHeader {
458        debug_assert!(lane < self.lane_count);
459        debug_assert!(idx < self.capacity);
460        // SAFETY: slots region begins at `slots_offset`; lane and idx
461        // are bounded by the validated mapping layout.
462        unsafe {
463            let base = self
464                .region
465                .as_ptr()
466                .add(self.slots_offset + lane * self.lane_stride);
467            base.add(idx * self.slot_stride).cast::<ShmSlotHeader>()
468        }
469    }
470
471    /// The payload bytes, as the atomics they have to be for the same reason
472    /// the header fields are.
473    ///
474    /// The cost is a byte at a time instead of a `memcpy`, over payloads this
475    /// workspace sizes in tens to hundreds of bytes — 18 for a metric sample,
476    /// 1024 at the largest — published at snapshot rather than request rate.
477    /// Word-at-a-time would need the payload region padded to a word multiple,
478    /// which changes the segment layout and so the compatibility check every
479    /// attaching process makes; worth doing if a payload ever grows enough to
480    /// notice, and not before.
481    unsafe fn payload_ptr(slot_ptr: *mut ShmSlotHeader) -> *mut AtomicU8 {
482        unsafe {
483            slot_ptr
484                .cast::<u8>()
485                .add(SLOT_HEADER_SIZE)
486                .cast::<AtomicU8>()
487        }
488    }
489
490    /// Head of the sole shared lane, or lane zero for a per-node ring.
491    pub fn head(&self) -> u64 {
492        self.lane_header(0).write_pos.load(Ordering::Acquire)
493    }
494
495    /// Current committed head for `node_id`'s lane.
496    pub fn lane_head(&self, node_id: NodeId) -> u64 {
497        let lane = self.lane_index(node_id);
498        self.lane_header(lane).write_pos.load(Ordering::Acquire)
499    }
500
501    /// Allocate one non-zero semantic version shared by all writer lanes and
502    /// processes attached to this ring.
503    pub fn next_version(&self) -> u64 {
504        let header = unsafe { &*(self.region.as_ptr() as *const ShmRingHeader) };
505        header
506            .version_counter
507            .fetch_add(1, Ordering::AcqRel)
508            .checked_add(1)
509            .expect("SHM ring semantic version exhausted")
510    }
511
512    /// Last semantic version allocated by any process attached to this ring.
513    pub fn current_version(&self) -> u64 {
514        let header = unsafe { &*(self.region.as_ptr() as *const ShmRingHeader) };
515        header.version_counter.load(Ordering::Acquire)
516    }
517
518    /// Append a frame. Atomically reserves the next counter, mints
519    /// the [`NetId64`], writes the slot. Returns the minted id.
520    pub fn write(
521        &self,
522        node_id: NodeId,
523        frame_kind: u8,
524        ver: u64,
525        payload: Bytes,
526    ) -> std::io::Result<NetId64> {
527        if payload.len() > self.payload_capacity {
528            return Err(std::io::Error::new(
529                std::io::ErrorKind::InvalidInput,
530                format!(
531                    "payload {} > ring payload capacity {}",
532                    payload.len(),
533                    self.payload_capacity
534                ),
535            ));
536        }
537
538        let lane = self.lane_index(node_id);
539        match self.topology {
540            RingTopology::Shared => {
541                let counter = self
542                    .lane_header(lane)
543                    .write_pos
544                    .fetch_add(1, Ordering::AcqRel);
545                Ok(self.write_slot(lane, node_id, counter, frame_kind, ver, &payload))
546            }
547            RingTopology::PerNode => {
548                let _write = self.write_locks[lane]
549                    .lock()
550                    .unwrap_or_else(|error| error.into_inner());
551                let lane_header = self.lane_header(lane);
552                let counter = lane_header.write_pos.load(Ordering::Relaxed);
553                let id = self.write_slot(lane, node_id, counter, frame_kind, ver, &payload);
554                lane_header
555                    .write_pos
556                    .store(counter.wrapping_add(1), Ordering::Release);
557                Ok(id)
558            }
559            RingTopology::SharedOrdered => {
560                let _write = self.write_locks[lane]
561                    .lock()
562                    .unwrap_or_else(|error| error.into_inner());
563                let _cross_process = self.region.lock_exclusive()?;
564                let lane_header = self.lane_header(lane);
565                let counter = lane_header.write_pos.load(Ordering::Relaxed);
566                let id = self.write_slot(lane, node_id, counter, frame_kind, ver, &payload);
567                lane_header
568                    .write_pos
569                    .store(counter.wrapping_add(1), Ordering::Release);
570                Ok(id)
571            }
572        }
573    }
574
575    /// Append one contiguous batch to a lane and return consecutive ids.
576    ///
577    /// Per-node and shared-ordered lane heads advance only after the complete
578    /// batch has committed. The batch may not exceed the ring capacity.
579    pub fn write_batch(
580        &self,
581        node_id: NodeId,
582        frame_kind: u8,
583        ver: u64,
584        payloads: Vec<Bytes>,
585    ) -> std::io::Result<Vec<NetId64>> {
586        if payloads.len() > self.capacity {
587            return Err(std::io::Error::new(
588                std::io::ErrorKind::InvalidInput,
589                format!("batch {} > ring capacity {}", payloads.len(), self.capacity),
590            ));
591        }
592        if let Some(payload) = payloads
593            .iter()
594            .find(|payload| payload.len() > self.payload_capacity)
595        {
596            return Err(std::io::Error::new(
597                std::io::ErrorKind::InvalidInput,
598                format!(
599                    "payload {} > ring payload capacity {}",
600                    payload.len(),
601                    self.payload_capacity
602                ),
603            ));
604        }
605        if payloads.is_empty() {
606            return Ok(Vec::new());
607        }
608
609        let lane = self.lane_index(node_id);
610        let write_slots = |start: u64, payloads: Vec<Bytes>| {
611            payloads
612                .into_iter()
613                .enumerate()
614                .map(|(offset, payload)| {
615                    self.write_slot(
616                        lane,
617                        node_id,
618                        start.wrapping_add(offset as u64),
619                        frame_kind,
620                        ver,
621                        &payload,
622                    )
623                })
624                .collect::<Vec<_>>()
625        };
626
627        match self.topology {
628            RingTopology::Shared => {
629                let start = self
630                    .lane_header(lane)
631                    .write_pos
632                    .fetch_add(payloads.len() as u64, Ordering::AcqRel);
633                Ok(write_slots(start, payloads))
634            }
635            RingTopology::PerNode => {
636                let _write = self.write_locks[lane]
637                    .lock()
638                    .unwrap_or_else(|error| error.into_inner());
639                let lane_header = self.lane_header(lane);
640                let start = lane_header.write_pos.load(Ordering::Relaxed);
641                let ids = write_slots(start, payloads);
642                lane_header
643                    .write_pos
644                    .store(start.wrapping_add(ids.len() as u64), Ordering::Release);
645                Ok(ids)
646            }
647            RingTopology::SharedOrdered => {
648                let _write = self.write_locks[lane]
649                    .lock()
650                    .unwrap_or_else(|error| error.into_inner());
651                let _cross_process = self.region.lock_exclusive()?;
652                let lane_header = self.lane_header(lane);
653                let start = lane_header.write_pos.load(Ordering::Relaxed);
654                let ids = write_slots(start, payloads);
655                lane_header
656                    .write_pos
657                    .store(start.wrapping_add(ids.len() as u64), Ordering::Release);
658                Ok(ids)
659            }
660        }
661    }
662
663    /// Read the slot whose counter matches `id.counter()`. Returns
664    /// `None` if the slot has been overwritten, was never written,
665    /// or a torn read could not be reconciled across two retries.
666    pub fn read(&self, id: NetId64) -> Option<Frame> {
667        if id.kind() != self.kind {
668            return None;
669        }
670        let lane = self.lane_index_for_frame(id)?;
671        let counter = id.counter();
672        let slot_idx = (counter as usize) & (self.capacity - 1);
673        let slot_ptr = self.slot_ptr(lane, slot_idx);
674
675        // Two retries — torn writes happen but resolve quickly.
676        for _ in 0..3 {
677            let Some(frame) = (unsafe { read_committed_frame(slot_ptr, self.payload_capacity) })
678            else {
679                continue;
680            };
681            if frame.id.counter() == counter {
682                return Some(frame);
683            } else {
684                // Slot now holds a different (later) id — wraparound.
685                return None;
686            }
687        }
688        None
689    }
690
691    /// Read the most recent frame (head - 1). Returns `None` if no
692    /// write has happened yet, or a torn read could not resolve.
693    pub fn read_head(&self) -> Option<Frame> {
694        let head = self.head();
695        if head == 0 {
696            return None;
697        }
698        let counter = head - 1;
699        let slot_idx = (counter as usize) & (self.capacity - 1);
700        let slot_ptr = self.slot_ptr(0, slot_idx);
701
702        for _ in 0..3 {
703            if let Some(frame) = unsafe { read_committed_frame(slot_ptr, self.payload_capacity) } {
704                return Some(frame);
705            }
706        }
707        None
708    }
709
710    /// Read whatever frame currently occupies the slot at
711    /// `counter % capacity`, regardless of which counter is stored
712    /// there. Used by walking readers that need slot-by-slot access
713    /// without knowing the writer's `NetId64` ahead of time.
714    pub fn read_at(&self, counter: u64) -> Option<Frame> {
715        self.read_lane_index_at(0, counter)
716    }
717
718    /// Read the frame currently occupying `node_id`'s lane slot.
719    pub fn read_lane_at(&self, node_id: NodeId, counter: u64) -> Option<Frame> {
720        let lane = self.lane_index(node_id);
721        self.read_lane_index_at(lane, counter)
722    }
723
724    fn read_lane_index_at(&self, lane: usize, counter: u64) -> Option<Frame> {
725        let slot_idx = (counter as usize) & (self.capacity - 1);
726        let slot_ptr = self.slot_ptr(lane, slot_idx);
727        for _ in 0..3 {
728            if let Some(frame) = unsafe { read_committed_frame(slot_ptr, self.payload_capacity) } {
729                return Some(frame);
730            }
731        }
732        None
733    }
734
735    pub(crate) fn read_state_at(&self, counter: u64) -> crate::ring::cursor::RingRead {
736        self.read_lane_index_state_at(0, counter)
737    }
738
739    pub(crate) fn read_lane_state_at(
740        &self,
741        node_id: NodeId,
742        counter: u64,
743    ) -> crate::ring::cursor::RingRead {
744        let lane = self.lane_index(node_id);
745        self.read_lane_index_state_at(lane, counter)
746    }
747
748    fn read_lane_index_state_at(&self, lane: usize, counter: u64) -> crate::ring::cursor::RingRead {
749        use crate::ring::cursor::RingRead;
750
751        let slot_idx = (counter as usize) & (self.capacity - 1);
752        let slot_ptr = self.slot_ptr(lane, slot_idx);
753        let expected_committed = counter
754            .checked_mul(2)
755            .and_then(|value| value.checked_add(2))
756            .expect("seq overflow");
757
758        for _ in 0..3 {
759            let sequence = unsafe { &*slot_ptr }.seq.load(Ordering::Acquire);
760            if sequence < expected_committed {
761                return if self.topology != RingTopology::Shared {
762                    RingRead::Unavailable
763                } else {
764                    RingRead::Pending
765                };
766            }
767            if sequence > expected_committed {
768                return RingRead::Unavailable;
769            }
770            if let Some(frame) = unsafe { read_committed_frame(slot_ptr, self.payload_capacity) } {
771                return if frame.id.counter() == counter {
772                    RingRead::Ready(frame)
773                } else {
774                    RingRead::Unavailable
775                };
776            }
777        }
778
779        let sequence = unsafe { &*slot_ptr }.seq.load(Ordering::Acquire);
780        if sequence < expected_committed {
781            if self.topology != RingTopology::Shared {
782                RingRead::Unavailable
783            } else {
784                RingRead::Pending
785            }
786        } else {
787            RingRead::Unavailable
788        }
789    }
790
791    /// Clear all slots and reset every lane head to zero.
792    ///
793    /// Intended for owner-controlled boot-time cleanup. Do not call
794    /// while other processes are publishing to this ring: it rewrites
795    /// the shared slot memory in place.
796    pub fn reset(&self) {
797        // SAFETY: the region is mapped and the slot area begins at
798        // HEADER_SIZE. The caller must ensure the ring is quiescent.
799        unsafe {
800            let slots_ptr = self.region.as_ptr().add(self.slots_offset);
801            ptr::write_bytes(slots_ptr, 0, self.lane_count * self.lane_stride);
802        }
803        for lane in 0..self.lane_count {
804            self.lane_header(lane).write_pos.store(0, Ordering::Release);
805        }
806        let header = unsafe { &*(self.region.as_ptr() as *const ShmRingHeader) };
807        header.version_counter.store(0, Ordering::Release);
808    }
809
810    fn lane_index(&self, node_id: NodeId) -> usize {
811        let lane = match self.topology {
812            RingTopology::Shared | RingTopology::SharedOrdered => 0,
813            RingTopology::PerNode => usize::from(node_id.get()),
814        };
815        assert!(
816            lane < self.lane_count,
817            "node {} is outside SHM ring lane count {}",
818            node_id.get(),
819            self.lane_count
820        );
821        lane
822    }
823
824    fn lane_index_for_frame(&self, id: NetId64) -> Option<usize> {
825        let lane = match self.topology {
826            RingTopology::Shared | RingTopology::SharedOrdered => 0,
827            RingTopology::PerNode => usize::from(id.node()),
828        };
829        (lane < self.lane_count).then_some(lane)
830    }
831
832    fn write_slot(
833        &self,
834        lane: usize,
835        node_id: NodeId,
836        counter: u64,
837        frame_kind: u8,
838        ver: u64,
839        payload: &[u8],
840    ) -> NetId64 {
841        let id = NetId64::make(self.kind, node_id.get(), counter);
842        let slot_idx = (counter as usize) & (self.capacity - 1);
843        let slot_ptr = self.slot_ptr(lane, slot_idx);
844
845        // Disruptor-style write: seq goes odd → write content → seq goes even.
846        //
847        // The odd marker is stored `Relaxed` and followed by a `Release`
848        // *fence*, not stored `Release`. A release store orders the accesses
849        // that come *before* it and says nothing about the ones after, so the
850        // content writes below were free to become visible ahead of the marker:
851        // a reader would then see an even seq on both sides of a slot that was
852        // being torn underneath it. The fence is what forbids that, because it
853        // orders everything before it — the marker — ahead of every store after
854        // it.
855        //
856        // The closing store stays `Release`: there it is the content writes
857        // that must be visible first, which is exactly what a release store
858        // gives. It pairs with the reader's `Acquire` fence.
859        //
860        // x86's store-store ordering hides the difference; aarch64 does not.
861        unsafe {
862            let slot = &*slot_ptr;
863            let mid_seq = counter
864                .checked_mul(2)
865                .and_then(|value| value.checked_add(1))
866                .expect("seq overflow");
867            let final_seq = mid_seq.wrapping_add(1);
868
869            slot.seq.store(mid_seq, Ordering::Relaxed);
870            fence(Ordering::Release);
871            slot.id.store(id.raw(), Ordering::Relaxed);
872            slot.ver.store(ver, Ordering::Relaxed);
873            slot.payload_len
874                .store(payload.len() as u32, Ordering::Relaxed);
875            slot.kind.store(frame_kind, Ordering::Relaxed);
876            let bytes = Self::payload_ptr(slot_ptr);
877            for (index, byte) in payload.iter().enumerate() {
878                (*bytes.add(index)).store(*byte, Ordering::Relaxed);
879            }
880            slot.seq.store(final_seq, Ordering::Release);
881        }
882
883        id
884    }
885}
886
887// ─────────────────────────────────────────────────────────────────────
888// Read-only inspection
889// ─────────────────────────────────────────────────────────────────────
890
891/// Persisted geometry of one existing Orbit SHM ring.
892///
893/// The values come from the mapped ring header. An observer does not need the
894/// producer's fleet capacity or compile-time [`RingSpec`] to inspect raw
895/// frames.
896#[derive(Clone, Debug, PartialEq, Eq)]
897pub struct ShmRingMetadata {
898    /// Exact POSIX SHM object name that was opened.
899    pub name: String,
900    /// Number of bytes reported by the existing SHM object.
901    pub segment_size: usize,
902    /// Persisted Orbit ring format version.
903    pub format_version: u32,
904    /// Stable Orbit type kind carried by this ring.
905    pub kind: u8,
906    /// Persisted per-lane capacity, payload limit, and topology.
907    pub spec: RingSpec,
908    /// Number of physical lanes present in the mapping.
909    pub lane_count: usize,
910}
911
912/// Read-only attachment to one existing Orbit SHM ring.
913///
914/// Attaching never creates, initializes, resets, locks, or unlinks a segment.
915/// Dropping the view only unmaps this process's local mapping. Use
916/// [`ShmRingView::lane`] to obtain a
917/// [`RingFrameSource`](crate::ring::cursor::RingFrameSource) for one physical
918/// lane.
919pub struct ShmRingView {
920    region: ShmRegion,
921    metadata: ShmRingMetadata,
922    slot_stride: usize,
923    slots_offset: usize,
924    lane_stride: usize,
925}
926
927impl std::fmt::Debug for ShmRingView {
928    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
929        formatter
930            .debug_struct("ShmRingView")
931            .field("metadata", &self.metadata)
932            .finish()
933    }
934}
935
936impl ShmRingView {
937    /// Attach to an existing ring owned by the effective user.
938    pub fn attach_existing(fleet_name: &str, kind: u8) -> std::io::Result<Self> {
939        // SAFETY: `geteuid` has no error path.
940        let uid = unsafe { libc::geteuid() };
941        Self::attach_existing_for_uid(fleet_name, kind, uid)
942    }
943
944    /// Attach to an existing ring under an explicit uid-scoped name.
945    ///
946    /// POSIX permissions still decide whether the calling process may open
947    /// another user's object.
948    pub fn attach_existing_for_uid(fleet_name: &str, kind: u8, uid: u32) -> std::io::Result<Self> {
949        let name = shm::ring_segment_name_for_uid(fleet_name, kind, uid);
950        let region = ShmRegion::open_existing_read_only(&name, HEADER_SIZE)?;
951
952        // SAFETY: the mapping is at least HEADER_SIZE bytes and the header is
953        // cache-line aligned at offset zero.
954        let header = unsafe { &*(region.as_ptr() as *const ShmRingHeader) };
955        if header.magic != MAGIC {
956            return Err(std::io::Error::new(
957                std::io::ErrorKind::InvalidData,
958                format!(
959                    "SHM segment {} has wrong magic 0x{:08X} (expected 0x{:08X})",
960                    name, header.magic, MAGIC
961                ),
962            ));
963        }
964        if header.version != VERSION {
965            return Err(std::io::Error::new(
966                std::io::ErrorKind::InvalidData,
967                format!(
968                    "SHM segment {} version {} != local {}",
969                    name, header.version, VERSION
970                ),
971            ));
972        }
973        if header.kind != kind {
974            return Err(std::io::Error::new(
975                std::io::ErrorKind::InvalidData,
976                format!(
977                    "SHM segment {} kind {} != requested {}",
978                    name, header.kind, kind
979                ),
980            ));
981        }
982
983        let topology = match header.topology {
984            value if value == RingTopology::Shared as u8 => RingTopology::Shared,
985            value if value == RingTopology::PerNode as u8 => RingTopology::PerNode,
986            value if value == RingTopology::SharedOrdered as u8 => RingTopology::SharedOrdered,
987            value => {
988                return Err(std::io::Error::new(
989                    std::io::ErrorKind::InvalidData,
990                    format!("SHM segment {name} has unknown topology {value}"),
991                ));
992            }
993        };
994        let capacity = usize::try_from(header.capacity).map_err(|_| {
995            std::io::Error::new(
996                std::io::ErrorKind::InvalidData,
997                format!(
998                    "SHM segment {name} capacity {} does not fit this platform",
999                    header.capacity
1000                ),
1001            )
1002        })?;
1003        let lane_count = usize::from(header.lane_count);
1004        if lane_count == 0 {
1005            return Err(std::io::Error::new(
1006                std::io::ErrorKind::InvalidData,
1007                format!("SHM segment {name} has zero lanes"),
1008            ));
1009        }
1010        if topology != RingTopology::PerNode && lane_count != 1 {
1011            return Err(std::io::Error::new(
1012                std::io::ErrorKind::InvalidData,
1013                format!(
1014                    "SHM segment {name} topology {topology:?} requires one lane, found {lane_count}"
1015                ),
1016            ));
1017        }
1018        let spec = RingSpec {
1019            capacity,
1020            payload_capacity: header.payload_capacity as usize,
1021            topology,
1022        };
1023        let fleet_capacity = if topology == RingTopology::PerNode {
1024            header.lane_count
1025        } else {
1026            1
1027        };
1028        let (slot_stride, slots_offset, expected_size) = checked_layout(spec, fleet_capacity)
1029            .map_err(|error| {
1030                std::io::Error::new(
1031                    std::io::ErrorKind::InvalidData,
1032                    format!("SHM segment {name} has invalid geometry: {error}"),
1033                )
1034            })?;
1035        if header.slot_stride as usize != slot_stride {
1036            return Err(std::io::Error::new(
1037                std::io::ErrorKind::InvalidData,
1038                format!(
1039                    "SHM segment {} slot stride {} != geometry {}",
1040                    name, header.slot_stride, slot_stride
1041                ),
1042            ));
1043        }
1044        if region.len() < expected_size {
1045            return Err(std::io::Error::new(
1046                std::io::ErrorKind::InvalidData,
1047                format!(
1048                    "SHM segment {} size {} is smaller than persisted geometry {}",
1049                    name,
1050                    region.len(),
1051                    expected_size
1052                ),
1053            ));
1054        }
1055        let lane_stride = capacity
1056            .checked_mul(slot_stride)
1057            .ok_or_else(|| invalid_input(format!("SHM segment {name} lane stride overflow")))?;
1058
1059        Ok(Self {
1060            metadata: ShmRingMetadata {
1061                name,
1062                segment_size: region.len(),
1063                format_version: header.version,
1064                kind,
1065                spec,
1066                lane_count,
1067            },
1068            region,
1069            slot_stride,
1070            slots_offset,
1071            lane_stride,
1072        })
1073    }
1074
1075    /// Attach and also verify the caller's typed ring contract.
1076    pub fn attach_typed<T: crate::OrbitTyped>(fleet_name: &str) -> std::io::Result<Self> {
1077        let view = Self::attach_existing(fleet_name, T::KIND)?;
1078        if view.metadata.spec != T::RING_SPEC {
1079            return Err(std::io::Error::new(
1080                std::io::ErrorKind::InvalidData,
1081                format!(
1082                    "OrbitTyped KIND {} declares {:?}; existing spec is {:?}",
1083                    T::KIND,
1084                    T::RING_SPEC,
1085                    view.metadata.spec
1086                ),
1087            ));
1088        }
1089        Ok(view)
1090    }
1091
1092    pub fn metadata(&self) -> &ShmRingMetadata {
1093        &self.metadata
1094    }
1095
1096    /// Last semantic version allocated by the ring at observation time.
1097    pub fn current_version(&self) -> u64 {
1098        self.header().version_counter.load(Ordering::Acquire)
1099    }
1100
1101    /// Current notification generation used by native readiness bridges.
1102    pub fn notification_generation(&self) -> u32 {
1103        self.header()
1104            .notification_generation
1105            .load(Ordering::Acquire)
1106    }
1107
1108    /// Select one physical lane as a generic read-only frame source.
1109    pub fn lane(&self, lane: usize) -> std::io::Result<ShmRingLaneView<'_>> {
1110        if lane >= self.metadata.lane_count {
1111            return Err(std::io::Error::new(
1112                std::io::ErrorKind::InvalidInput,
1113                format!(
1114                    "lane {lane} is outside SHM ring lane count {}",
1115                    self.metadata.lane_count
1116                ),
1117            ));
1118        }
1119        Ok(ShmRingLaneView { ring: self, lane })
1120    }
1121
1122    fn header(&self) -> &ShmRingHeader {
1123        // SAFETY: construction validated the mapped header.
1124        unsafe { &*(self.region.as_ptr() as *const ShmRingHeader) }
1125    }
1126
1127    fn lane_header(&self, lane: usize) -> &ShmLaneHeader {
1128        debug_assert!(lane < self.metadata.lane_count);
1129        // SAFETY: construction validated lane count and the complete geometry.
1130        unsafe {
1131            &*(self
1132                .region
1133                .as_ptr()
1134                .add(HEADER_SIZE + lane * LANE_HEADER_SIZE) as *const ShmLaneHeader)
1135        }
1136    }
1137
1138    fn slot_ptr(&self, lane: usize, counter: u64) -> *mut ShmSlotHeader {
1139        debug_assert!(lane < self.metadata.lane_count);
1140        let slot = (counter as usize) & (self.metadata.spec.capacity - 1);
1141        // SAFETY: construction checked the complete persisted geometry and the
1142        // lane/slot indices are bounded above.
1143        unsafe {
1144            self.region
1145                .as_ptr()
1146                .add(self.slots_offset + lane * self.lane_stride + slot * self.slot_stride)
1147                .cast::<ShmSlotHeader>()
1148        }
1149    }
1150
1151    fn read_lane_state_at(&self, lane: usize, counter: u64) -> crate::ring::cursor::RingRead {
1152        use crate::ring::cursor::RingRead;
1153
1154        let slot_ptr = self.slot_ptr(lane, counter);
1155        let Some(expected_committed) = counter
1156            .checked_mul(2)
1157            .and_then(|value| value.checked_add(2))
1158        else {
1159            return RingRead::Unavailable;
1160        };
1161
1162        for _ in 0..3 {
1163            // SAFETY: `slot_ptr` addresses a validated slot in the mapping.
1164            let sequence = unsafe { &*slot_ptr }.seq.load(Ordering::Acquire);
1165            if sequence < expected_committed {
1166                return if self.metadata.spec.topology == RingTopology::Shared {
1167                    RingRead::Pending
1168                } else {
1169                    RingRead::Unavailable
1170                };
1171            }
1172            if sequence > expected_committed {
1173                return RingRead::Unavailable;
1174            }
1175            // SAFETY: `slot_ptr` addresses a validated slot in the mapping.
1176            if let Some(frame) =
1177                unsafe { read_committed_frame(slot_ptr, self.metadata.spec.payload_capacity) }
1178            {
1179                let lane_matches = self.metadata.spec.topology != RingTopology::PerNode
1180                    || usize::from(frame.id.node()) == lane;
1181                return if frame.id.kind() == self.metadata.kind
1182                    && frame.id.counter() == counter
1183                    && lane_matches
1184                {
1185                    RingRead::Ready(frame)
1186                } else {
1187                    RingRead::Unavailable
1188                };
1189            }
1190        }
1191        RingRead::Unavailable
1192    }
1193}
1194
1195/// One physical lane of a [`ShmRingView`].
1196#[derive(Clone, Copy, Debug)]
1197pub struct ShmRingLaneView<'a> {
1198    ring: &'a ShmRingView,
1199    lane: usize,
1200}
1201
1202impl ShmRingLaneView<'_> {
1203    pub fn index(self) -> usize {
1204        self.lane
1205    }
1206
1207    /// Snapshot the currently visible committed/reserved head.
1208    pub fn head(self) -> u64 {
1209        self.ring
1210            .lane_header(self.lane)
1211            .write_pos
1212            .load(Ordering::Acquire)
1213    }
1214
1215    /// Counter range that may still be retained at the observed head.
1216    pub fn retained_range(self) -> std::ops::Range<u64> {
1217        let head = self.head();
1218        head.saturating_sub(self.ring.metadata.spec.capacity as u64)..head
1219    }
1220
1221    /// Read the exact counter if its frame is still committed in this lane.
1222    pub fn read_at(self, counter: u64) -> Option<Frame> {
1223        match self.ring.read_lane_state_at(self.lane, counter) {
1224            crate::ring::cursor::RingRead::Ready(frame) => Some(frame),
1225            crate::ring::cursor::RingRead::Pending | crate::ring::cursor::RingRead::Unavailable => {
1226                None
1227            }
1228        }
1229    }
1230
1231    /// Read the frame immediately below the observed head.
1232    pub fn read_head(self) -> Option<Frame> {
1233        let head = self.head();
1234        head.checked_sub(1)
1235            .and_then(|counter| self.read_at(counter))
1236    }
1237}
1238
1239impl crate::ring::cursor::RingFrameSource for ShmRingLaneView<'_> {
1240    fn kind(&self) -> u8 {
1241        self.ring.metadata.kind
1242    }
1243
1244    fn head(&self) -> u64 {
1245        (*self).head()
1246    }
1247
1248    fn capacity(&self) -> usize {
1249        self.ring.metadata.spec.capacity
1250    }
1251
1252    fn read_at(&self, counter: u64) -> Option<Frame> {
1253        (*self).read_at(counter)
1254    }
1255
1256    fn read_state_at(&self, counter: u64) -> crate::ring::cursor::RingRead {
1257        self.ring.read_lane_state_at(self.lane, counter)
1258    }
1259}
1260
1261// ─────────────────────────────────────────────────────────────────────
1262// ShmRingRegistry — per-fleet, per-KIND map of ShmRings
1263// ─────────────────────────────────────────────────────────────────────
1264
1265/// Type-keyed registry of [`ShmRing`]s held by a fleet. One ring per
1266/// `OrbitTyped::KIND` byte; created on demand the first time a kind
1267/// is published or queried.
1268pub struct ShmRingRegistry {
1269    fleet_name: String,
1270    fleet_capacity: u16,
1271    rings: dashmap::DashMap<u8, std::sync::Arc<ShmRing>>,
1272}
1273
1274impl ShmRingRegistry {
1275    pub fn new(fleet_name: impl Into<String>, fleet_capacity: u16) -> Self {
1276        Self {
1277            fleet_name: fleet_name.into(),
1278            fleet_capacity,
1279            rings: dashmap::DashMap::new(),
1280        }
1281    }
1282
1283    /// Get-or-create the SHM ring for `kind`. Failure here means the
1284    /// SHM open or attach failed (permissions, name too long, etc.)
1285    /// and is propagated as `io::Error`.
1286    pub fn get_or_create_for<T: crate::OrbitTyped>(
1287        &self,
1288    ) -> std::io::Result<std::sync::Arc<ShmRing>> {
1289        if let Some(entry) = self.rings.get(&T::KIND) {
1290            if entry.spec() != T::RING_SPEC {
1291                return Err(std::io::Error::new(
1292                    std::io::ErrorKind::InvalidData,
1293                    format!(
1294                        "OrbitTyped KIND {} was reused with ring spec {:?}; existing spec is {:?}",
1295                        T::KIND,
1296                        T::RING_SPEC,
1297                        entry.spec()
1298                    ),
1299                ));
1300            }
1301            return Ok(entry.clone());
1302        }
1303        let ring = std::sync::Arc::new(ShmRing::open_or_create_for_fleet(
1304            &self.fleet_name,
1305            T::KIND,
1306            T::RING_SPEC,
1307            self.fleet_capacity,
1308        )?);
1309        let entry = self.rings.entry(T::KIND).or_insert_with(|| ring.clone());
1310        if entry.spec() != T::RING_SPEC {
1311            return Err(std::io::Error::new(
1312                std::io::ErrorKind::InvalidData,
1313                format!(
1314                    "OrbitTyped KIND {} raced with ring spec {:?}; installed spec is {:?}",
1315                    T::KIND,
1316                    T::RING_SPEC,
1317                    entry.spec()
1318                ),
1319            ));
1320        }
1321        Ok(entry.clone())
1322    }
1323
1324    /// Look up a ring that has already been created for `kind`.
1325    /// Returns `None` if no such ring exists yet.
1326    pub fn lookup(&self, kind: u8) -> Option<std::sync::Arc<ShmRing>> {
1327        self.rings.get(&kind).map(|e| e.clone())
1328    }
1329}
1330
1331/// Read a slot's content; returns `None` if the seq indicates an
1332/// in-flight write or if pre/post seqs disagree (torn read).
1333///
1334/// # Safety
1335///
1336/// `slot_ptr` must point at a valid `ShmSlot` mapped into our
1337/// address space and aligned per the `repr(C, align(64))` layout.
1338unsafe fn read_committed_frame(
1339    slot_ptr: *mut ShmSlotHeader,
1340    payload_capacity: usize,
1341) -> Option<Frame> {
1342    let slot = unsafe { &*slot_ptr };
1343    let seq_pre = slot.seq.load(Ordering::Acquire);
1344    if seq_pre == 0 {
1345        // never written
1346        return None;
1347    }
1348    if seq_pre & 1 == 1 {
1349        // writer in progress
1350        return None;
1351    }
1352
1353    // Read content fields. `Relaxed` throughout: the seq load above and the
1354    // fence below are what order this, and a value read here is only trusted
1355    // once the two seqs agree.
1356    let id = NetId64::from_raw(slot.id.load(Ordering::Relaxed));
1357    let kind = slot.kind.load(Ordering::Relaxed);
1358    let ver = slot.ver.load(Ordering::Relaxed);
1359    let payload_len = slot.payload_len.load(Ordering::Relaxed) as usize;
1360    if payload_len > payload_capacity {
1361        // corrupt — bail
1362        return None;
1363    }
1364    // `with_capacity` + `push` rather than `vec![0; len]` + overwrite: the
1365    // zeroing would be written over by the very next loop, and a walking reader
1366    // over a 16k-slot ring pays that twice for every byte it reads. The push
1367    // cannot reallocate — the capacity is reserved above.
1368    let payload_src = unsafe { ShmRing::payload_ptr(slot_ptr) };
1369    let mut payload_buf = Vec::with_capacity(payload_len);
1370    for index in 0..payload_len {
1371        payload_buf.push(unsafe { (*payload_src.add(index)).load(Ordering::Relaxed) });
1372    }
1373
1374    // The mirror of the writer's fence. `Acquire` on a load orders the accesses
1375    // that come *after* it, so reading the closing seq with `Acquire` would
1376    // leave the content reads above free to be reordered past it — and the
1377    // comparison below would then be comparing against a slot it never actually
1378    // read. An acquire fence orders those reads ahead of the load that follows,
1379    // which is the guarantee this check is asking for; the load itself needs no
1380    // ordering of its own once the fence is there.
1381    fence(Ordering::Acquire);
1382    let seq_post = slot.seq.load(Ordering::Relaxed);
1383    if seq_pre != seq_post {
1384        // torn write — caller can retry
1385        return None;
1386    }
1387
1388    Some(Frame {
1389        id,
1390        kind,
1391        ver,
1392        payload: Bytes::from(payload_buf),
1393    })
1394}
1395
1396#[cfg(test)]
1397mod tests {
1398    use std::sync::atomic::{AtomicU64, Ordering};
1399
1400    use nix::sys::wait::{WaitStatus, waitpid};
1401    use nix::unistd::{ForkResult, fork};
1402
1403    use super::*;
1404    use crate::ring::cursor::{RingCursor, poll_ring};
1405
1406    #[test]
1407    fn cursor_retries_a_claimed_slot_after_it_commits() {
1408        static TEST_ID: AtomicU64 = AtomicU64::new(0);
1409
1410        let test_id = TEST_ID.fetch_add(1, Ordering::Relaxed);
1411        let fleet_name = format!("p{:x}{test_id:x}", std::process::id());
1412        let ring = ShmRing::open_or_create(&fleet_name, 199, RingSpec::new(4, 16))
1413            .expect("create test ring");
1414        ring.reset();
1415        let slot_ptr = ring.slot_ptr(0, 0);
1416
1417        ring.lane_header(0).write_pos.store(1, Ordering::Release);
1418        unsafe { &*slot_ptr }.seq.store(1, Ordering::Release);
1419
1420        let mut cursor = RingCursor::from_start();
1421        let pending = poll_ring(&ring, &mut cursor);
1422        assert!(pending.is_empty());
1423        assert_eq!(cursor.next_counter(), 0);
1424
1425        let payload = b"ready";
1426        unsafe {
1427            let slot = &*slot_ptr;
1428            slot.id.store(
1429                NetId64::make(199, NodeId::ZERO.get(), 0).raw(),
1430                Ordering::Relaxed,
1431            );
1432            slot.ver.store(7, Ordering::Relaxed);
1433            slot.payload_len
1434                .store(payload.len() as u32, Ordering::Relaxed);
1435            slot.kind.store(1, Ordering::Relaxed);
1436            let bytes = ShmRing::payload_ptr(slot_ptr);
1437            for (index, byte) in payload.iter().enumerate() {
1438                (*bytes.add(index)).store(*byte, Ordering::Relaxed);
1439            }
1440        }
1441        unsafe { &*slot_ptr }.seq.store(2, Ordering::Release);
1442
1443        let committed = poll_ring(&ring, &mut cursor);
1444        assert_eq!(committed.frames.len(), 1);
1445        assert_eq!(&committed.frames[0].payload[..], payload);
1446        assert_eq!(cursor.next_counter(), 1);
1447
1448        ring.unlink().expect("unlink test ring");
1449    }
1450
1451    #[test]
1452    fn per_node_head_ignores_a_writer_that_dies_before_commit() {
1453        static TEST_ID: AtomicU64 = AtomicU64::new(0);
1454
1455        let test_id = TEST_ID.fetch_add(1, Ordering::Relaxed);
1456        let fleet_name = format!("d{:x}{test_id:x}", std::process::id());
1457        let spec = RingSpec::per_node(4, 16);
1458        let abandoned = ShmRing::open_or_create_for_fleet(&fleet_name, 198, spec, 2)
1459            .expect("create per-node test ring");
1460        abandoned.reset();
1461
1462        let slot_ptr = abandoned.slot_ptr(1, 0);
1463        unsafe { &*slot_ptr }.seq.store(1, Ordering::Release);
1464        assert_eq!(abandoned.lane_head(NodeId::new(1)), 0);
1465
1466        let replacement = ShmRing::open_or_create_for_fleet(&fleet_name, 198, spec, 2)
1467            .expect("replacement attaches");
1468        let id = replacement
1469            .write(NodeId::new(1), 1, 9, Bytes::from_static(b"recovered"))
1470            .expect("replacement commits");
1471
1472        assert_eq!(id.counter(), 0);
1473        assert_eq!(replacement.lane_head(NodeId::new(1)), 1);
1474        assert_eq!(
1475            &replacement.read(id).expect("frame visible").payload[..],
1476            b"recovered"
1477        );
1478
1479        replacement.unlink().expect("unlink test ring");
1480    }
1481
1482    #[test]
1483    fn per_node_lane_serializes_concurrent_local_publishers() {
1484        static TEST_ID: AtomicU64 = AtomicU64::new(0);
1485
1486        let test_id = TEST_ID.fetch_add(1, Ordering::Relaxed);
1487        let fleet_name = format!("c{:x}{test_id:x}", std::process::id());
1488        let ring = std::sync::Arc::new(
1489            ShmRing::open_or_create_for_fleet(&fleet_name, 197, RingSpec::per_node(512, 0), 2)
1490                .expect("create concurrent writer ring"),
1491        );
1492        ring.reset();
1493
1494        let mut writers = Vec::new();
1495        for _ in 0..4 {
1496            let ring = ring.clone();
1497            writers.push(std::thread::spawn(move || {
1498                (0..64)
1499                    .map(|_| {
1500                        ring.write(NodeId::new(1), 1, 0, Bytes::new())
1501                            .expect("publish")
1502                            .counter()
1503                    })
1504                    .collect::<Vec<_>>()
1505            }));
1506        }
1507
1508        let mut counters = writers
1509            .into_iter()
1510            .flat_map(|writer| writer.join().expect("writer joins"))
1511            .collect::<Vec<_>>();
1512        counters.sort_unstable();
1513        assert_eq!(counters, (0..256).collect::<Vec<_>>());
1514        assert_eq!(ring.lane_head(NodeId::new(1)), 256);
1515
1516        ring.unlink().expect("unlink test ring");
1517    }
1518
1519    #[test]
1520    fn shared_ordered_recovers_after_a_locked_writer_process_dies() {
1521        static TEST_ID: AtomicU64 = AtomicU64::new(0);
1522
1523        let test_id = TEST_ID.fetch_add(1, Ordering::Relaxed);
1524        let fleet_name = format!("o{:x}{test_id:x}", std::process::id());
1525        let spec = RingSpec::shared_ordered(4, 16);
1526        let ring = ShmRing::open_or_create_for_fleet(&fleet_name, 196, spec, 2)
1527            .expect("create shared-ordered test ring");
1528        ring.reset();
1529
1530        match unsafe { fork() }.expect("fork test writer") {
1531            ForkResult::Child => {
1532                // Use the handle inherited while it was idle. It must not carry
1533                // a persistent lock descriptor across fork; the child opens a
1534                // fresh file description for this critical section.
1535                let _lock = ring.region.lock_exclusive().expect("child locks ring");
1536                let slot_ptr = ring.slot_ptr(0, 0);
1537                unsafe { &*slot_ptr }.seq.store(1, Ordering::Release);
1538
1539                // Exit without destructors: the kernel, not Rust cleanup,
1540                // must release the descriptor lock.
1541                unsafe { libc::_exit(0) };
1542            }
1543            ForkResult::Parent { child } => {
1544                let status = waitpid(child, None).expect("wait for abandoned writer");
1545                assert!(matches!(status, WaitStatus::Exited(_, 0)));
1546
1547                // A separately opened handle must acquire the lock after the
1548                // child dies. If the original region retained an idle flock fd,
1549                // fork would duplicate its open-file-description into the child
1550                // and the dead writer's lock would still be attached here.
1551                let replacement = ShmRing::open_or_create_for_fleet(&fleet_name, 196, spec, 2)
1552                    .expect("replacement attaches");
1553                let recovered_lock = replacement
1554                    .region
1555                    .try_lock_exclusive()
1556                    .expect("kernel released dead writer lock");
1557                drop(recovered_lock);
1558
1559                let id = replacement
1560                    .write(NodeId::new(1), 1, 9, Bytes::from_static(b"recovered"))
1561                    .expect("replacement commits counter zero");
1562                assert_eq!(id.counter(), 0);
1563                assert_eq!(replacement.head(), 1);
1564                assert_eq!(
1565                    &replacement.read(id).expect("recovered frame").payload[..],
1566                    b"recovered"
1567                );
1568
1569                replacement
1570                    .unlink()
1571                    .expect("unlink shared-ordered test ring");
1572            }
1573        }
1574    }
1575}