Skip to main content

subetha_cxc/
shared_broadcast_ring.rs

1//! `SharedBroadcastRing` - single-producer, multi-consumer pub/sub
2//! ring backed by an MMF.
3//!
4//! Distinct from [`SharedRing`](crate::SharedRing) (MPMC; each slot
5//! consumed exactly once). In a broadcast ring, EVERY registered
6//! consumer sees EVERY message independently with its own cursor.
7//! This is the pub/sub / kafka-topic / log-tail shape.
8//!
9//! # Layout
10//!
11//! ```text
12//! +---------------------------+
13//! | BroadcastHeader (128B)    |
14//! |   magic, capacity         |
15//! |   producer_seq: AtomicU64 |
16//! |   consumer_seqs[MAX]:     |
17//! |     [AtomicU64; MAX_CONS] |
18//! |   consumer_active[MAX]:   |
19//! |     [AtomicU32; MAX_CONS] |
20//! +---------------------------+
21//! | Slot[0]    (64B)          |  version + payload
22//! | ...                       |
23//! | Slot[capacity - 1]        |
24//! +---------------------------+
25//! ```
26//!
27//! Each slot is a SeqLock cell; producer writes under odd version
28//! and bumps even on completion. Consumers read with version-spin.
29//!
30//! # Protocol
31//!
32//! ## Producer
33//!
34//! `try_push(payload)`:
35//! 1. Find oldest consumer cursor: `min_consumer = min(active
36//!    consumer_seqs)`.
37//! 2. If `producer_seq - min_consumer >= capacity`, the ring is
38//!    full from at least one consumer's perspective: return
39//!    `BroadcastFull`. (A slot in `[min_consumer, producer_seq)`
40//!    can't be overwritten without making that consumer miss it.)
41//! 3. Write payload into `slot[producer_seq % capacity]` under
42//!    SeqLock.
43//! 4. `producer_seq.fetch_add(1, Release)`.
44//!
45//! ## Consumer
46//!
47//! `register()` -> consumer_idx in 0..MAX_CONSUMERS:
48//! - CAS the first inactive slot to active; initialise
49//!   `consumer_seqs[i]` to current `producer_seq` (so the consumer
50//!   starts from "now," not from the beginning of history).
51//!
52//! `try_recv(consumer_idx, buf)`:
53//! 1. Load `my_seq = consumer_seqs[consumer_idx]`.
54//! 2. If `my_seq >= producer_seq`, no new messages: return Empty.
55//! 3. SeqLock-read `slot[my_seq % capacity]` into `buf`.
56//! 4. `consumer_seqs[consumer_idx].fetch_add(1, Release)`.
57//!
58//! `unregister(consumer_idx)`: mark slot inactive; producer no
59//! longer waits for this consumer's cursor.
60//!
61//! # Why single-producer?
62//!
63//! Multi-producer broadcast adds complexity (producers must
64//! coordinate slot claim AND ordering must be preserved per topic).
65//! The single-producer case covers most pub/sub use cases: one
66//! source, many subscribers. For multi-producer fan-in, the
67//! producers should fan into a SharedRing first, then a single
68//! relay process re-emits into a SharedBroadcastRing.
69
70use std::fs::{File, OpenOptions};
71use std::mem::size_of;
72use std::path::Path;
73use std::sync::atomic::{AtomicU32, AtomicU64, Ordering};
74
75use memmap2::{MmapMut, MmapOptions};
76
77pub const BROADCAST_MAGIC: u32 = 0x4150_4243;
78pub const BROADCAST_PAYLOAD_BYTES: usize = 52;
79pub const MAX_CONSUMERS: usize = 16;
80
81/// Header is two cache lines so the producer_seq and the consumer
82/// seq array can sit on separate lines, reducing false-sharing
83/// between producer-side and consumer-side hot atomics.
84#[repr(C, align(64))]
85pub struct BroadcastHeader {
86    pub magic: u32,
87    pub capacity: u32,
88    pub producer_seq: AtomicU64,
89    _pad1: [u8; 48],
90
91    pub consumer_seqs: [AtomicU64; MAX_CONSUMERS],
92    pub consumer_active: [AtomicU32; MAX_CONSUMERS],
93}
94
95#[repr(C, align(64))]
96pub struct BroadcastSlot {
97    pub version: AtomicU32,
98    _pad: [u8; 4],
99    pub payload: [u8; BROADCAST_PAYLOAD_BYTES],
100}
101
102const _: () = {
103    assert!(size_of::<BroadcastSlot>() == 64);
104};
105
106pub const fn broadcast_file_size(capacity: usize) -> usize {
107    size_of::<BroadcastHeader>() + capacity * size_of::<BroadcastSlot>()
108}
109
110#[derive(Debug, Clone, Copy, PartialEq, Eq)]
111pub enum BroadcastError {
112    Full,
113    Empty,
114    NoConsumerSlot,
115    InvalidConsumer,
116    PayloadTooLarge,
117    LayoutMismatch,
118    IoError(std::io::ErrorKind),
119}
120
121impl From<std::io::Error> for BroadcastError {
122    fn from(e: std::io::Error) -> Self { Self::IoError(e.kind()) }
123}
124
125pub struct SharedBroadcastRing {
126    _backing: BroadcastBacking,
127    raw_ptr: *mut u8,
128    capacity: usize,
129    header_sidecar: subetha_core::HandshakeHeader,
130    ring_sidecar: Box<subetha_core::ObservationRing>,
131}
132
133unsafe impl Send for SharedBroadcastRing {}
134unsafe impl Sync for SharedBroadcastRing {}
135
136/// Storage backing for a `SharedBroadcastRing`. The hot-path
137/// pointer (`raw_ptr` on `SharedBroadcastRing`) is cached at
138/// construction; this enum holds the owning resource alive for
139/// the lifetime of the ring. The held values are intentionally
140/// only kept for ownership / flush-target purposes; the
141/// `File` and `ShmFile` payloads are not read at runtime.
142#[allow(dead_code)]
143enum BroadcastBacking {
144    /// In-process anonymous mmap. No file, no shm-name.
145    Anon(MmapMut),
146    /// File-backed mmap. Cross-process via the OS page cache.
147    File(File, MmapMut),
148    /// Named shared-memory mmap (ShmFs locale). Cross-process,
149    /// RAM-resident, never touches the page cache.
150    Shm(crate::shm_file::ShmFile),
151}
152
153/// Backing-agnostic layout init. Writes the header magic + capacity
154/// and zeroes the slot array at the given raw pointer. Caller
155/// guarantees that `ptr` points to at least
156/// `broadcast_file_size(capacity)` bytes of mutable, suitably-
157/// aligned memory.
158unsafe fn init_broadcast_layout_raw(ptr: *mut u8, capacity: usize) {
159    let hdr_ptr = ptr as *mut BroadcastHeader;
160    unsafe {
161        std::ptr::write_bytes(hdr_ptr as *mut u8, 0, size_of::<BroadcastHeader>());
162        (*hdr_ptr).magic = BROADCAST_MAGIC;
163        (*hdr_ptr).capacity = capacity as u32;
164    }
165    for i in 0..capacity {
166        let slot_ptr = unsafe {
167            ptr.add(size_of::<BroadcastHeader>())
168                .add(i * size_of::<BroadcastSlot>())
169        } as *mut BroadcastSlot;
170        unsafe {
171            std::ptr::write(slot_ptr, BroadcastSlot {
172                version: AtomicU32::new(0),
173                _pad: [0; 4],
174                payload: [0u8; BROADCAST_PAYLOAD_BYTES],
175            });
176        }
177    }
178}
179
180/// Layout init for the elected creator in the attach flow: capacity
181/// first, magic last via a volatile store, because attachers spin on
182/// the magic. The zeroed region is already the empty slot array
183/// (version 0), so the slots are not touched.
184///
185/// # Safety
186/// `ptr` addresses at least `broadcast_file_size(capacity)` writable
187/// zeroed bytes.
188unsafe fn init_broadcast_attach_raw(ptr: *mut u8, capacity: usize) {
189    let hdr_ptr = ptr as *mut BroadcastHeader;
190    unsafe {
191        (*hdr_ptr).capacity = capacity as u32;
192        std::ptr::write_volatile(&raw mut (*hdr_ptr).magic, BROADCAST_MAGIC);
193    }
194}
195
196impl subetha_sidecar::AdaptiveInstance for SharedBroadcastRing {
197    fn header(&self) -> &subetha_core::HandshakeHeader { &self.header_sidecar }
198    fn ring(&self) -> &subetha_core::ObservationRing { &self.ring_sidecar }
199    fn make_policy(&self) -> Box<dyn subetha_sidecar::Policy> {
200        Box::new(subetha_sidecar::NoMigrationPolicy)
201    }
202}
203
204impl SharedBroadcastRing {
205    /// File-backed broadcast ring; cross-process visibility via the OS
206    /// page cache. Obtains the ring at `path`: it initializes an empty
207    /// one only when the path does not yet exist, and otherwise
208    /// attaches with published slots and versions intact. A region
209    /// built with a different capacity is a `LayoutMismatch`.
210    /// [`reset`](Self::reset) reinitializes.
211    pub fn create(path: impl AsRef<Path>, capacity: usize) -> Result<Self, BroadcastError> {
212        assert!(capacity >= 2);
213        assert!(capacity <= u32::MAX as usize);
214        let total = broadcast_file_size(capacity);
215        let (file, mut mmap) = crate::mmf_attach::create_or_attach(
216            path.as_ref(),
217            total,
218            |ptr| unsafe { init_broadcast_attach_raw(ptr, capacity) },
219            |ptr| unsafe { (*(ptr as *const BroadcastHeader)).magic == BROADCAST_MAGIC },
220        )?;
221        let raw_ptr = mmap.as_mut_ptr();
222        let hdr = unsafe { &*(raw_ptr as *const BroadcastHeader) };
223        if hdr.capacity != capacity as u32 {
224            return Err(BroadcastError::LayoutMismatch);
225        }
226        Ok(Self {
227            _backing: BroadcastBacking::File(file, mmap),
228            raw_ptr, capacity,
229            header_sidecar: subetha_core::HandshakeHeader::new(),
230            ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
231        })
232    }
233
234    /// Truncate the ring at `path` and initialize an empty one,
235    /// discarding every published slot live peers share. For a caller
236    /// that knows it owns the path.
237    pub fn reset(path: impl AsRef<Path>, capacity: usize) -> Result<Self, BroadcastError> {
238        assert!(capacity >= 2);
239        assert!(capacity <= u32::MAX as usize);
240        let total = broadcast_file_size(capacity);
241        let (file, mut mmap) = crate::mmf_attach::reset(path.as_ref(), total, |ptr| unsafe {
242            init_broadcast_attach_raw(ptr, capacity)
243        })?;
244        let raw_ptr = mmap.as_mut_ptr();
245        Ok(Self {
246            _backing: BroadcastBacking::File(file, mmap),
247            raw_ptr, capacity,
248            header_sidecar: subetha_core::HandshakeHeader::new(),
249            ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
250        })
251    }
252
253    /// Anonymous in-process broadcast ring. Fastest construction;
254    /// skips file create + ftruncate + first-page-fault. In-process
255    /// only - subscribers in other processes cannot connect.
256    pub fn create_anon(capacity: usize) -> Result<Self, BroadcastError> {
257        assert!(capacity >= 2);
258        assert!(capacity <= u32::MAX as usize);
259        let total = broadcast_file_size(capacity);
260        let mut mmap = MmapOptions::new().len(total).map_anon()?;
261        let raw_ptr = mmap.as_mut_ptr();
262        unsafe { init_broadcast_layout_raw(raw_ptr, capacity); }
263        Ok(Self {
264            _backing: BroadcastBacking::Anon(mmap),
265            raw_ptr, capacity,
266            header_sidecar: subetha_core::HandshakeHeader::new(),
267            ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
268        })
269    }
270
271    /// Open an existing file-backed broadcast ring. Validates
272    /// magic + capacity.
273    pub fn open(path: impl AsRef<Path>, expected_capacity: usize) -> Result<Self, BroadcastError> {
274        let total = broadcast_file_size(expected_capacity);
275        let file = OpenOptions::new().read(true).write(true).open(path.as_ref())?;
276        if file.metadata()?.len() < total as u64 {
277            return Err(BroadcastError::LayoutMismatch);
278        }
279        let mut mmap = unsafe { MmapOptions::new().len(total).map_mut(&file)? };
280        let raw_ptr = mmap.as_mut_ptr();
281        let hdr = unsafe { &*(raw_ptr as *const BroadcastHeader) };
282        if hdr.magic != BROADCAST_MAGIC || hdr.capacity != expected_capacity as u32 {
283            return Err(BroadcastError::LayoutMismatch);
284        }
285        Ok(Self {
286            _backing: BroadcastBacking::File(file, mmap),
287            raw_ptr, capacity: expected_capacity,
288            header_sidecar: subetha_core::HandshakeHeader::new(),
289            ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
290        })
291    }
292
293    /// Build a fresh broadcast ring on top of a named RAM-resident
294    /// shared-memory backing. Cross-process visible via the
295    /// `logical_name` of the underlying
296    /// [`ShmFile`](crate::shm_file::ShmFile); never touches
297    /// the page cache. The `ShmFile` must be sized to at least
298    /// `broadcast_file_size(capacity)` bytes.
299    pub fn create_from_shm(
300        mut shm: crate::shm_file::ShmFile,
301        capacity: usize,
302    ) -> Result<Self, BroadcastError> {
303        assert!(capacity >= 2);
304        assert!(capacity <= u32::MAX as usize);
305        let total = broadcast_file_size(capacity);
306        if shm.len() < total {
307            return Err(BroadcastError::LayoutMismatch);
308        }
309        let raw_ptr = shm.as_mut_slice().as_mut_ptr();
310        unsafe { init_broadcast_layout_raw(raw_ptr, capacity); }
311        Ok(Self {
312            _backing: BroadcastBacking::Shm(shm),
313            raw_ptr, capacity,
314            header_sidecar: subetha_core::HandshakeHeader::new(),
315            ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
316        })
317    }
318
319    /// Open an existing named ShmFs-backed broadcast ring.
320    /// Validates magic + capacity. Does NOT re-initialise the
321    /// layout - the layout must already be present from a prior
322    /// `create_from_shm` on the same logical name.
323    pub fn open_from_shm(
324        mut shm: crate::shm_file::ShmFile,
325        expected_capacity: usize,
326    ) -> Result<Self, BroadcastError> {
327        let total = broadcast_file_size(expected_capacity);
328        if shm.len() < total {
329            return Err(BroadcastError::LayoutMismatch);
330        }
331        let raw_ptr = shm.as_mut_slice().as_mut_ptr();
332        let hdr = unsafe { &*(raw_ptr as *const BroadcastHeader) };
333        if hdr.magic != BROADCAST_MAGIC || hdr.capacity != expected_capacity as u32 {
334            return Err(BroadcastError::LayoutMismatch);
335        }
336        Ok(Self {
337            _backing: BroadcastBacking::Shm(shm),
338            raw_ptr, capacity: expected_capacity,
339            header_sidecar: subetha_core::HandshakeHeader::new(),
340            ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
341        })
342    }
343
344
345    #[inline]
346    pub fn capacity(&self) -> usize { self.capacity }
347
348    fn header(&self) -> &BroadcastHeader {
349        unsafe { &*(self.raw_ptr as *const BroadcastHeader) }
350    }
351
352    /// Reduce a logical sequence number to a physical slot index. Uses a
353    /// bit-mask when capacity is a power of two (the common ring sizing) -
354    /// removing the `% capacity` hardware DIV - and the modulo otherwise.
355    /// `capacity` is loop-invariant, so the pow2 test folds away in hot
356    /// loops; no cached field needed.
357    #[inline]
358    fn wrap(&self, i: usize) -> usize {
359        if self.capacity.is_power_of_two() {
360            i & (self.capacity - 1)
361        } else {
362            i % self.capacity
363        }
364    }
365
366    fn slot(&self, idx: usize) -> &BroadcastSlot {
367        let physical = self.wrap(idx);
368        let base = unsafe { self.raw_ptr.add(size_of::<BroadcastHeader>()) };
369        unsafe { &*(base.add(physical * size_of::<BroadcastSlot>()) as *const BroadcastSlot) }
370    }
371
372    /// Register as a consumer. Returns a consumer index in
373    /// `0..MAX_CONSUMERS`; that index is used for all subsequent
374    /// recv calls. Initialises the consumer's cursor to the current
375    /// producer_seq (consumer starts reading from "now," not history).
376    pub fn register_consumer(&self) -> Result<usize, BroadcastError> {
377        let hdr = self.header();
378        for i in 0..MAX_CONSUMERS {
379            if hdr.consumer_active[i].compare_exchange(
380                0, 1, Ordering::AcqRel, Ordering::Acquire,
381            ).is_ok() {
382                let cur_producer = hdr.producer_seq.load(Ordering::Acquire);
383                hdr.consumer_seqs[i].store(cur_producer, Ordering::Release);
384                self.ring_sidecar
385                    .push_op(crate::sidecar_ops::broadcast_ring::OP_REGISTER, 0);
386                return Ok(i);
387            }
388        }
389        self.ring_sidecar
390            .push_op(crate::sidecar_ops::broadcast_ring::OP_REGISTER, 1); // no consumer slot available
391        Err(BroadcastError::NoConsumerSlot)
392    }
393
394    /// Unregister a consumer. After this, the producer no longer
395    /// waits for this cursor when reclaiming slots.
396    pub fn unregister_consumer(&self, consumer_idx: usize) {
397        if consumer_idx >= MAX_CONSUMERS { return; }
398        let hdr = self.header();
399        hdr.consumer_active[consumer_idx].store(0, Ordering::Release);
400        self.ring_sidecar
401            .push_op(crate::sidecar_ops::broadcast_ring::OP_UNREGISTER, 0);
402    }
403
404    /// Compute the slot-reclaim safety margin: the producer can push
405    /// when `producer_seq - min_active_consumer_seq < capacity`.
406    /// Returns the smallest cursor among active consumers (or
407    /// producer_seq when there are no active consumers; broadcasts
408    /// to nobody can always push).
409    fn min_consumer_seq(&self) -> u64 {
410        let hdr = self.header();
411        let producer = hdr.producer_seq.load(Ordering::Acquire);
412        let mut min = u64::MAX;
413        let mut any = false;
414        for i in 0..MAX_CONSUMERS {
415            if hdr.consumer_active[i].load(Ordering::Acquire) != 0 {
416                let s = hdr.consumer_seqs[i].load(Ordering::Acquire);
417                if s < min { min = s; }
418                any = true;
419            }
420        }
421        if !any { producer } else { min }
422    }
423
424    /// Push a message. Returns `Err(Full)` when at least one active
425    /// consumer hasn't yet read a previous slot we'd overwrite.
426    pub fn try_push(&self, payload: &[u8]) -> Result<(), BroadcastError> {
427        if payload.len() > BROADCAST_PAYLOAD_BYTES {
428            return Err(BroadcastError::PayloadTooLarge);
429        }
430        let hdr = self.header();
431        let producer = hdr.producer_seq.load(Ordering::Acquire);
432        let min_consumer = self.min_consumer_seq();
433        if producer.saturating_sub(min_consumer) >= self.capacity as u64 {
434            self.ring_sidecar
435                .push_op(crate::sidecar_ops::broadcast_ring::OP_PUSH, 1); // full
436            return Err(BroadcastError::Full);
437        }
438        // SeqLock write the slot at producer_seq % capacity.
439        let slot = self.slot(producer as usize);
440        slot.version.fetch_add(1, Ordering::AcqRel); // odd
441        let dst = unsafe {
442            self.raw_ptr
443                .add(size_of::<BroadcastHeader>())
444                .add(self.wrap(producer as usize) * size_of::<BroadcastSlot>())
445                .add(std::mem::offset_of!(BroadcastSlot, payload))
446        };
447        unsafe {
448            std::ptr::write_bytes(dst, 0, BROADCAST_PAYLOAD_BYTES);
449            std::ptr::copy_nonoverlapping(payload.as_ptr(), dst, payload.len());
450        }
451        slot.version.fetch_add(1, Ordering::AcqRel); // even
452        hdr.producer_seq.fetch_add(1, Ordering::Release);
453        self.ring_sidecar
454            .push_op(crate::sidecar_ops::broadcast_ring::OP_PUSH, 0);
455        Ok(())
456    }
457
458    /// Receive the next unread message for `consumer_idx`. Returns
459    /// the number of bytes filled (always BROADCAST_PAYLOAD_BYTES;
460    /// caller knows the inner-event size from its T contract).
461    pub fn try_recv(&self, consumer_idx: usize, out: &mut [u8]) -> Result<usize, BroadcastError> {
462        if consumer_idx >= MAX_CONSUMERS { return Err(BroadcastError::InvalidConsumer); }
463        let hdr = self.header();
464        if hdr.consumer_active[consumer_idx].load(Ordering::Acquire) == 0 {
465            return Err(BroadcastError::InvalidConsumer);
466        }
467        let my_seq = hdr.consumer_seqs[consumer_idx].load(Ordering::Acquire);
468        let producer = hdr.producer_seq.load(Ordering::Acquire);
469        if my_seq >= producer {
470            self.ring_sidecar
471                .push_op(crate::sidecar_ops::broadcast_ring::OP_RECV, 2); // empty
472            return Err(BroadcastError::Empty);
473        }
474        // SeqLock read the slot at my_seq % capacity.
475        let slot = self.slot(my_seq as usize);
476        loop {
477            let v1 = slot.version.load(Ordering::Acquire);
478            if v1 & 1 != 0 {
479                std::hint::spin_loop();
480                continue;
481            }
482            let src = unsafe {
483                self.raw_ptr
484                    .add(size_of::<BroadcastHeader>())
485                    .add(self.wrap(my_seq as usize) * size_of::<BroadcastSlot>())
486                    .add(std::mem::offset_of!(BroadcastSlot, payload))
487            };
488            let n = out.len().min(BROADCAST_PAYLOAD_BYTES);
489            unsafe { std::ptr::copy_nonoverlapping(src, out.as_mut_ptr(), n); }
490            let v2 = slot.version.load(Ordering::Acquire);
491            if v1 == v2 {
492                hdr.consumer_seqs[consumer_idx].fetch_add(1, Ordering::Release);
493                self.ring_sidecar
494                    .push_op(crate::sidecar_ops::broadcast_ring::OP_RECV, 0);
495                return Ok(n);
496            }
497        }
498    }
499
500    /// Number of messages this consumer has not yet read.
501    pub fn lag(&self, consumer_idx: usize) -> u64 {
502        let hdr = self.header();
503        if consumer_idx >= MAX_CONSUMERS { return 0; }
504        let prod = hdr.producer_seq.load(Ordering::Acquire);
505        let my = hdr.consumer_seqs[consumer_idx].load(Ordering::Acquire);
506        prod.saturating_sub(my)
507    }
508
509    /// Current producer cursor (total messages pushed since creation).
510    pub fn producer_position(&self) -> u64 {
511        self.header().producer_seq.load(Ordering::Acquire)
512    }
513
514    /// Number of currently active consumers.
515    pub fn active_consumer_count(&self) -> usize {
516        let hdr = self.header();
517        (0..MAX_CONSUMERS)
518            .filter(|&i| hdr.consumer_active[i].load(Ordering::Acquire) != 0)
519            .count()
520    }
521
522    pub fn flush(&self) -> Result<(), BroadcastError> {
523        match &self._backing {
524            BroadcastBacking::File(_, mmap) => mmap.flush()?,
525            BroadcastBacking::Anon(mmap) => mmap.flush()?,
526            BroadcastBacking::Shm(_) => {} // RAM-resident; nothing to sync
527        }
528        Ok(())
529    }
530
531    /// Non-blocking flush: schedules a writeback via the OS.
532    /// Note: Windows is only partially async (sync to page cache,
533    /// not to disk).
534    pub fn flush_async(&self) -> Result<(), BroadcastError> {
535        match &self._backing {
536            BroadcastBacking::File(_, mmap) => mmap.flush_async()?,
537            BroadcastBacking::Anon(mmap) => mmap.flush_async()?,
538            BroadcastBacking::Shm(_) => {} // RAM-resident; nothing to sync
539        }
540        Ok(())
541    }
542
543    /// Whether every currently-active consumer has read every
544    /// item the producer has published. Used by the capacity-morph
545    /// wrapper to decide whether a stale broadcast backing can be
546    /// dropped (all subscribers have caught up to the frozen
547    /// producer position).
548    pub fn is_fully_drained(&self) -> bool {
549        let hdr = self.header();
550        let prod = hdr.producer_seq.load(Ordering::Acquire);
551        for i in 0..MAX_CONSUMERS {
552            if hdr.consumer_active[i].load(Ordering::Acquire) != 0
553                && hdr.consumer_seqs[i].load(Ordering::Acquire) < prod
554            {
555                return false;
556            }
557        }
558        true
559    }
560}
561
562#[cfg(test)]
563mod tests {
564    use super::*;
565    use std::sync::Arc;
566    use std::thread;
567    use std::time::Duration;
568
569    fn tmp(name: &str) -> std::path::PathBuf {
570        let mut p = std::env::temp_dir();
571        let pid = std::process::id();
572        p.push(format!("subetha-broadcast-{name}-{pid}.bin"));
573        p
574    }
575
576    fn payload_of(v: u32) -> [u8; BROADCAST_PAYLOAD_BYTES] {
577        let mut b = [0u8; BROADCAST_PAYLOAD_BYTES];
578        b[0..4].copy_from_slice(&v.to_le_bytes());
579        b
580    }
581    fn unpack(b: &[u8]) -> u32 {
582        u32::from_le_bytes(b[0..4].try_into().unwrap())
583    }
584
585    #[test]
586    fn create_initial_state() {
587        let p = tmp("init");
588        let r = SharedBroadcastRing::create(&p, 8).unwrap();
589        assert_eq!(r.capacity(), 8);
590        assert_eq!(r.producer_position(), 0);
591        assert_eq!(r.active_consumer_count(), 0);
592        std::fs::remove_file(&p).ok();
593    }
594
595    /// A second create attaches with published slots in place; reset
596    /// is what strips them.
597    #[test]
598    fn second_create_attaches_and_keeps_slots() {
599        let p = tmp("attach");
600        std::fs::remove_file(&p).ok();
601        let r = SharedBroadcastRing::create(&p, 8).unwrap();
602        r.try_push(&payload_of(42)).unwrap();
603
604        let r2 = SharedBroadcastRing::create(&p, 8).unwrap();
605        assert_eq!(r2.producer_position(), 1, "attach rewound the producer");
606        assert!(matches!(
607            SharedBroadcastRing::create(&p, 4),
608            Err(BroadcastError::LayoutMismatch),
609        ));
610
611        // Windows refuses to truncate a mapped file, so every handle goes
612        // before the reset.
613        drop(r);
614        drop(r2);
615        let fresh = SharedBroadcastRing::reset(&p, 8).unwrap();
616        assert_eq!(fresh.producer_position(), 0, "reset kept a published slot");
617        drop(fresh);
618        std::fs::remove_file(&p).ok();
619    }
620
621    #[test]
622    fn one_producer_one_consumer_round_trip() {
623        let p = tmp("1p1c");
624        let r = SharedBroadcastRing::create(&p, 8).unwrap();
625        let c = r.register_consumer().unwrap();
626        r.try_push(&payload_of(42)).unwrap();
627        let mut buf = [0u8; BROADCAST_PAYLOAD_BYTES];
628        let n = r.try_recv(c, &mut buf).unwrap();
629        assert_eq!(n, BROADCAST_PAYLOAD_BYTES);
630        assert_eq!(unpack(&buf), 42);
631        assert_eq!(r.try_recv(c, &mut buf).err(), Some(BroadcastError::Empty));
632        std::fs::remove_file(&p).ok();
633    }
634
635    #[test]
636    fn three_consumers_each_see_all_messages() {
637        let p = tmp("1p3c");
638        let r = SharedBroadcastRing::create(&p, 16).unwrap();
639        let c0 = r.register_consumer().unwrap();
640        let c1 = r.register_consumer().unwrap();
641        let c2 = r.register_consumer().unwrap();
642        for i in 0..5u32 { r.try_push(&payload_of(i * 10)).unwrap(); }
643        let mut buf = [0u8; BROADCAST_PAYLOAD_BYTES];
644        for c in [c0, c1, c2] {
645            for i in 0..5u32 {
646                r.try_recv(c, &mut buf).unwrap();
647                assert_eq!(unpack(&buf), i * 10,
648                    "consumer {c} should see message {i}");
649            }
650            assert_eq!(r.try_recv(c, &mut buf).err(), Some(BroadcastError::Empty));
651        }
652        std::fs::remove_file(&p).ok();
653    }
654
655    #[test]
656    fn lagging_consumer_blocks_producer() {
657        let p = tmp("lag-blocks");
658        let r = SharedBroadcastRing::create(&p, 4).unwrap();
659        let c = r.register_consumer().unwrap();
660        let _c = c;
661        // Fill the ring.
662        for i in 0..4u32 { r.try_push(&payload_of(i)).unwrap(); }
663        // Consumer hasn't read anything; next push should fail Full.
664        assert_eq!(r.try_push(&payload_of(99)).err(), Some(BroadcastError::Full));
665        std::fs::remove_file(&p).ok();
666    }
667
668    #[test]
669    fn unregistering_lagging_consumer_unblocks_producer() {
670        let p = tmp("unreg-unblocks");
671        let r = SharedBroadcastRing::create(&p, 4).unwrap();
672        let c = r.register_consumer().unwrap();
673        for i in 0..4u32 { r.try_push(&payload_of(i)).unwrap(); }
674        assert_eq!(r.try_push(&payload_of(99)).err(), Some(BroadcastError::Full));
675        r.unregister_consumer(c);
676        // With no active consumers, producer can push freely.
677        r.try_push(&payload_of(99)).unwrap();
678        std::fs::remove_file(&p).ok();
679    }
680
681    #[test]
682    fn consumer_registered_late_starts_at_current_producer() {
683        let p = tmp("late-consumer");
684        let r = SharedBroadcastRing::create(&p, 8).unwrap();
685        // Push 3 messages BEFORE registering any consumer.
686        for i in 0..3u32 { r.try_push(&payload_of(i)).unwrap(); }
687        // Now register; this consumer should see only NEW messages.
688        let c = r.register_consumer().unwrap();
689        let mut buf = [0u8; BROADCAST_PAYLOAD_BYTES];
690        assert_eq!(r.try_recv(c, &mut buf).err(), Some(BroadcastError::Empty));
691        r.try_push(&payload_of(100)).unwrap();
692        r.try_recv(c, &mut buf).unwrap();
693        assert_eq!(unpack(&buf), 100);
694        std::fs::remove_file(&p).ok();
695    }
696
697    #[test]
698    fn lag_returns_pending_count() {
699        let p = tmp("lag");
700        let r = SharedBroadcastRing::create(&p, 8).unwrap();
701        let c = r.register_consumer().unwrap();
702        for i in 0..3u32 { r.try_push(&payload_of(i)).unwrap(); }
703        assert_eq!(r.lag(c), 3);
704        let mut buf = [0u8; BROADCAST_PAYLOAD_BYTES];
705        r.try_recv(c, &mut buf).unwrap();
706        assert_eq!(r.lag(c), 2);
707        std::fs::remove_file(&p).ok();
708    }
709
710    #[test]
711    fn max_consumers_returns_no_slot_when_full() {
712        let p = tmp("max-cons");
713        let r = SharedBroadcastRing::create(&p, 4).unwrap();
714        for _ in 0..MAX_CONSUMERS {
715            r.register_consumer().unwrap();
716        }
717        assert_eq!(r.register_consumer().err(), Some(BroadcastError::NoConsumerSlot));
718        std::fs::remove_file(&p).ok();
719    }
720
721    #[test]
722    fn cross_handle_pub_sub() {
723        let p = tmp("cross-handle");
724        let pub_handle = SharedBroadcastRing::create(&p, 8).unwrap();
725        let sub_handle = SharedBroadcastRing::open(&p, 8).unwrap();
726        let c = sub_handle.register_consumer().unwrap();
727        pub_handle.try_push(&payload_of(7777)).unwrap();
728        let mut buf = [0u8; BROADCAST_PAYLOAD_BYTES];
729        sub_handle.try_recv(c, &mut buf).unwrap();
730        assert_eq!(unpack(&buf), 7777);
731        std::fs::remove_file(&p).ok();
732    }
733
734    #[test]
735    fn payload_too_large_rejected_at_push() {
736        let p = tmp("oversized");
737        let r = SharedBroadcastRing::create(&p, 4).unwrap();
738        let _c = r.register_consumer().unwrap();
739        let big = vec![0u8; BROADCAST_PAYLOAD_BYTES + 1];
740        assert_eq!(r.try_push(&big).err(), Some(BroadcastError::PayloadTooLarge));
741        std::fs::remove_file(&p).ok();
742    }
743
744    #[test]
745    fn concurrent_consumers_all_drain_correctly() {
746        let p = tmp("concurrent");
747        let r = Arc::new(SharedBroadcastRing::create(&p, 256).unwrap());
748        let n_consumers = 4;
749        let n_msgs = 100u32;
750        let consumer_ids: Vec<usize> = (0..n_consumers).map(|_| r.register_consumer().unwrap()).collect();
751
752        let r_p = r.clone();
753        let producer = thread::spawn(move || {
754            for i in 0..n_msgs {
755                while r_p.try_push(&payload_of(i)).is_err() {
756                    thread::yield_now();
757                }
758            }
759        });
760
761        let mut handles = vec![];
762        for &c in &consumer_ids {
763            let r = r.clone();
764            handles.push(thread::spawn(move || {
765                let mut received = vec![];
766                let mut buf = [0u8; BROADCAST_PAYLOAD_BYTES];
767                while received.len() < n_msgs as usize {
768                    match r.try_recv(c, &mut buf) {
769                        Ok(_) => received.push(unpack(&buf)),
770                        Err(BroadcastError::Empty) => thread::yield_now(),
771                        Err(e) => panic!("unexpected error: {e:?}"),
772                    }
773                }
774                received
775            }));
776        }
777        producer.join().unwrap();
778        for h in handles {
779            let got = h.join().unwrap();
780            assert_eq!(got.len(), n_msgs as usize);
781            for (i, v) in got.iter().enumerate() {
782                assert_eq!(*v, i as u32, "message order must be preserved");
783            }
784        }
785        std::fs::remove_file(&p).ok();
786    }
787
788    #[test]
789    fn slow_consumer_doesnt_break_fast_consumer() {
790        let p = tmp("slow-fast");
791        let r = Arc::new(SharedBroadcastRing::create(&p, 8).unwrap());
792        let _slow = r.register_consumer().unwrap();
793        let fast = r.register_consumer().unwrap();
794
795        for i in 0..4u32 { r.try_push(&payload_of(i)).unwrap(); }
796        // Fast consumer drains; slow consumer hasn't moved.
797        let mut buf = [0u8; BROADCAST_PAYLOAD_BYTES];
798        for i in 0..4u32 {
799            r.try_recv(fast, &mut buf).unwrap();
800            assert_eq!(unpack(&buf), i);
801        }
802        // Producer is still bounded by slow consumer's cursor though;
803        // try to push 5 more - should hit Full at some point.
804        let mut pushed = 0;
805        for i in 100..200u32 {
806            if r.try_push(&payload_of(i)).is_err() { break; }
807            pushed += 1;
808        }
809        assert!(pushed <= 4, "should be bounded by slow consumer; pushed {pushed}");
810        std::fs::remove_file(&p).ok();
811    }
812
813    #[test]
814    fn observer_can_wait_for_publisher() {
815        let p = tmp("wait");
816        let r = Arc::new(SharedBroadcastRing::create(&p, 8).unwrap());
817        let c = r.register_consumer().unwrap();
818
819        let r_p = r.clone();
820        let pusher = thread::spawn(move || {
821            thread::sleep(Duration::from_millis(20));
822            r_p.try_push(&payload_of(555)).unwrap();
823        });
824        let mut buf = [0u8; BROADCAST_PAYLOAD_BYTES];
825        let r_c = r.clone();
826        let consumer = thread::spawn(move || {
827            loop {
828                match r_c.try_recv(c, &mut buf) {
829                    Ok(_) => break unpack(&buf),
830                    Err(BroadcastError::Empty) => thread::yield_now(),
831                    Err(e) => panic!("{e:?}"),
832                }
833            }
834        });
835        pusher.join().unwrap();
836        assert_eq!(consumer.join().unwrap(), 555);
837        std::fs::remove_file(&p).ok();
838    }
839
840    #[test]
841    fn disk_persistence_survives_reopen() {
842        let p = tmp("disk");
843        {
844            let r = SharedBroadcastRing::create(&p, 8).unwrap();
845            let _c = r.register_consumer().unwrap();
846            r.try_push(&payload_of(1234)).unwrap();
847            r.try_push(&payload_of(5678)).unwrap();
848            r.flush().unwrap();
849        }
850        let r2 = SharedBroadcastRing::open(&p, 8).unwrap();
851        // Producer seq should be 2; consumer 0 should still be at 0.
852        assert_eq!(r2.producer_position(), 2);
853        assert_eq!(r2.lag(0), 2);
854        let mut buf = [0u8; BROADCAST_PAYLOAD_BYTES];
855        r2.try_recv(0, &mut buf).unwrap();
856        assert_eq!(unpack(&buf), 1234);
857        r2.try_recv(0, &mut buf).unwrap();
858        assert_eq!(unpack(&buf), 5678);
859        std::fs::remove_file(&p).ok();
860    }
861}