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
180impl subetha_sidecar::AdaptiveInstance for SharedBroadcastRing {
181    fn header(&self) -> &subetha_core::HandshakeHeader { &self.header_sidecar }
182    fn ring(&self) -> &subetha_core::ObservationRing { &self.ring_sidecar }
183    fn make_policy(&self) -> Box<dyn subetha_sidecar::Policy> {
184        Box::new(subetha_sidecar::NoMigrationPolicy)
185    }
186}
187
188impl SharedBroadcastRing {
189    /// File-backed broadcast ring; cross-process visibility via
190    /// the OS page cache.
191    pub fn create(path: impl AsRef<Path>, capacity: usize) -> Result<Self, BroadcastError> {
192        assert!(capacity >= 2);
193        assert!(capacity <= u32::MAX as usize);
194        let total = broadcast_file_size(capacity);
195        let file = OpenOptions::new()
196            .read(true).write(true).create(true).truncate(true)
197            .open(path.as_ref())?;
198        file.set_len(total as u64)?;
199        let mut mmap = unsafe { MmapOptions::new().len(total).map_mut(&file)? };
200        let raw_ptr = mmap.as_mut_ptr();
201        unsafe { init_broadcast_layout_raw(raw_ptr, capacity); }
202        Ok(Self {
203            _backing: BroadcastBacking::File(file, mmap),
204            raw_ptr, capacity,
205            header_sidecar: subetha_core::HandshakeHeader::new(),
206            ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
207        })
208    }
209
210    /// Anonymous in-process broadcast ring. Fastest construction;
211    /// skips file create + ftruncate + first-page-fault. In-process
212    /// only - subscribers in other processes cannot connect.
213    pub fn create_anon(capacity: usize) -> Result<Self, BroadcastError> {
214        assert!(capacity >= 2);
215        assert!(capacity <= u32::MAX as usize);
216        let total = broadcast_file_size(capacity);
217        let mut mmap = MmapOptions::new().len(total).map_anon()?;
218        let raw_ptr = mmap.as_mut_ptr();
219        unsafe { init_broadcast_layout_raw(raw_ptr, capacity); }
220        Ok(Self {
221            _backing: BroadcastBacking::Anon(mmap),
222            raw_ptr, capacity,
223            header_sidecar: subetha_core::HandshakeHeader::new(),
224            ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
225        })
226    }
227
228    /// Open an existing file-backed broadcast ring. Validates
229    /// magic + capacity.
230    pub fn open(path: impl AsRef<Path>, expected_capacity: usize) -> Result<Self, BroadcastError> {
231        let total = broadcast_file_size(expected_capacity);
232        let file = OpenOptions::new().read(true).write(true).open(path.as_ref())?;
233        if file.metadata()?.len() < total as u64 {
234            return Err(BroadcastError::LayoutMismatch);
235        }
236        let mut mmap = unsafe { MmapOptions::new().len(total).map_mut(&file)? };
237        let raw_ptr = mmap.as_mut_ptr();
238        let hdr = unsafe { &*(raw_ptr as *const BroadcastHeader) };
239        if hdr.magic != BROADCAST_MAGIC || hdr.capacity != expected_capacity as u32 {
240            return Err(BroadcastError::LayoutMismatch);
241        }
242        Ok(Self {
243            _backing: BroadcastBacking::File(file, mmap),
244            raw_ptr, capacity: expected_capacity,
245            header_sidecar: subetha_core::HandshakeHeader::new(),
246            ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
247        })
248    }
249
250    /// Build a fresh broadcast ring on top of a named RAM-resident
251    /// shared-memory backing. Cross-process visible via the
252    /// `logical_name` of the underlying
253    /// [`ShmFile`](crate::shm_file::ShmFile); never touches
254    /// the page cache. The `ShmFile` must be sized to at least
255    /// `broadcast_file_size(capacity)` bytes.
256    pub fn create_from_shm(
257        mut shm: crate::shm_file::ShmFile,
258        capacity: usize,
259    ) -> Result<Self, BroadcastError> {
260        assert!(capacity >= 2);
261        assert!(capacity <= u32::MAX as usize);
262        let total = broadcast_file_size(capacity);
263        if shm.len() < total {
264            return Err(BroadcastError::LayoutMismatch);
265        }
266        let raw_ptr = shm.as_mut_slice().as_mut_ptr();
267        unsafe { init_broadcast_layout_raw(raw_ptr, capacity); }
268        Ok(Self {
269            _backing: BroadcastBacking::Shm(shm),
270            raw_ptr, capacity,
271            header_sidecar: subetha_core::HandshakeHeader::new(),
272            ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
273        })
274    }
275
276    /// Open an existing named ShmFs-backed broadcast ring.
277    /// Validates magic + capacity. Does NOT re-initialise the
278    /// layout - the layout must already be present from a prior
279    /// `create_from_shm` on the same logical name.
280    pub fn open_from_shm(
281        mut shm: crate::shm_file::ShmFile,
282        expected_capacity: usize,
283    ) -> Result<Self, BroadcastError> {
284        let total = broadcast_file_size(expected_capacity);
285        if shm.len() < total {
286            return Err(BroadcastError::LayoutMismatch);
287        }
288        let raw_ptr = shm.as_mut_slice().as_mut_ptr();
289        let hdr = unsafe { &*(raw_ptr as *const BroadcastHeader) };
290        if hdr.magic != BROADCAST_MAGIC || hdr.capacity != expected_capacity as u32 {
291            return Err(BroadcastError::LayoutMismatch);
292        }
293        Ok(Self {
294            _backing: BroadcastBacking::Shm(shm),
295            raw_ptr, capacity: expected_capacity,
296            header_sidecar: subetha_core::HandshakeHeader::new(),
297            ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
298        })
299    }
300
301
302    #[inline]
303    pub fn capacity(&self) -> usize { self.capacity }
304
305    fn header(&self) -> &BroadcastHeader {
306        unsafe { &*(self.raw_ptr as *const BroadcastHeader) }
307    }
308
309    /// Reduce a logical sequence number to a physical slot index. Uses a
310    /// bit-mask when capacity is a power of two (the common ring sizing) -
311    /// removing the `% capacity` hardware DIV - and the modulo otherwise.
312    /// `capacity` is loop-invariant, so the pow2 test folds away in hot
313    /// loops; no cached field needed.
314    #[inline]
315    fn wrap(&self, i: usize) -> usize {
316        if self.capacity.is_power_of_two() {
317            i & (self.capacity - 1)
318        } else {
319            i % self.capacity
320        }
321    }
322
323    fn slot(&self, idx: usize) -> &BroadcastSlot {
324        let physical = self.wrap(idx);
325        let base = unsafe { self.raw_ptr.add(size_of::<BroadcastHeader>()) };
326        unsafe { &*(base.add(physical * size_of::<BroadcastSlot>()) as *const BroadcastSlot) }
327    }
328
329    /// Register as a consumer. Returns a consumer index in
330    /// `0..MAX_CONSUMERS`; that index is used for all subsequent
331    /// recv calls. Initialises the consumer's cursor to the current
332    /// producer_seq (consumer starts reading from "now," not history).
333    pub fn register_consumer(&self) -> Result<usize, BroadcastError> {
334        let hdr = self.header();
335        for i in 0..MAX_CONSUMERS {
336            if hdr.consumer_active[i].compare_exchange(
337                0, 1, Ordering::AcqRel, Ordering::Acquire,
338            ).is_ok() {
339                let cur_producer = hdr.producer_seq.load(Ordering::Acquire);
340                hdr.consumer_seqs[i].store(cur_producer, Ordering::Release);
341                self.ring_sidecar
342                    .push_op(crate::sidecar_ops::broadcast_ring::OP_REGISTER, 0);
343                return Ok(i);
344            }
345        }
346        self.ring_sidecar
347            .push_op(crate::sidecar_ops::broadcast_ring::OP_REGISTER, 1); // no consumer slot available
348        Err(BroadcastError::NoConsumerSlot)
349    }
350
351    /// Unregister a consumer. After this, the producer no longer
352    /// waits for this cursor when reclaiming slots.
353    pub fn unregister_consumer(&self, consumer_idx: usize) {
354        if consumer_idx >= MAX_CONSUMERS { return; }
355        let hdr = self.header();
356        hdr.consumer_active[consumer_idx].store(0, Ordering::Release);
357        self.ring_sidecar
358            .push_op(crate::sidecar_ops::broadcast_ring::OP_UNREGISTER, 0);
359    }
360
361    /// Compute the slot-reclaim safety margin: the producer can push
362    /// when `producer_seq - min_active_consumer_seq < capacity`.
363    /// Returns the smallest cursor among active consumers (or
364    /// producer_seq when there are no active consumers; broadcasts
365    /// to nobody can always push).
366    fn min_consumer_seq(&self) -> u64 {
367        let hdr = self.header();
368        let producer = hdr.producer_seq.load(Ordering::Acquire);
369        let mut min = u64::MAX;
370        let mut any = false;
371        for i in 0..MAX_CONSUMERS {
372            if hdr.consumer_active[i].load(Ordering::Acquire) != 0 {
373                let s = hdr.consumer_seqs[i].load(Ordering::Acquire);
374                if s < min { min = s; }
375                any = true;
376            }
377        }
378        if !any { producer } else { min }
379    }
380
381    /// Push a message. Returns `Err(Full)` when at least one active
382    /// consumer hasn't yet read a previous slot we'd overwrite.
383    pub fn try_push(&self, payload: &[u8]) -> Result<(), BroadcastError> {
384        if payload.len() > BROADCAST_PAYLOAD_BYTES {
385            return Err(BroadcastError::PayloadTooLarge);
386        }
387        let hdr = self.header();
388        let producer = hdr.producer_seq.load(Ordering::Acquire);
389        let min_consumer = self.min_consumer_seq();
390        if producer.saturating_sub(min_consumer) >= self.capacity as u64 {
391            self.ring_sidecar
392                .push_op(crate::sidecar_ops::broadcast_ring::OP_PUSH, 1); // full
393            return Err(BroadcastError::Full);
394        }
395        // SeqLock write the slot at producer_seq % capacity.
396        let slot = self.slot(producer as usize);
397        slot.version.fetch_add(1, Ordering::AcqRel); // odd
398        let dst = unsafe {
399            self.raw_ptr
400                .add(size_of::<BroadcastHeader>())
401                .add(self.wrap(producer as usize) * size_of::<BroadcastSlot>())
402                .add(std::mem::offset_of!(BroadcastSlot, payload))
403        };
404        unsafe {
405            std::ptr::write_bytes(dst, 0, BROADCAST_PAYLOAD_BYTES);
406            std::ptr::copy_nonoverlapping(payload.as_ptr(), dst, payload.len());
407        }
408        slot.version.fetch_add(1, Ordering::AcqRel); // even
409        hdr.producer_seq.fetch_add(1, Ordering::Release);
410        self.ring_sidecar
411            .push_op(crate::sidecar_ops::broadcast_ring::OP_PUSH, 0);
412        Ok(())
413    }
414
415    /// Receive the next unread message for `consumer_idx`. Returns
416    /// the number of bytes filled (always BROADCAST_PAYLOAD_BYTES;
417    /// caller knows the inner-event size from its T contract).
418    pub fn try_recv(&self, consumer_idx: usize, out: &mut [u8]) -> Result<usize, BroadcastError> {
419        if consumer_idx >= MAX_CONSUMERS { return Err(BroadcastError::InvalidConsumer); }
420        let hdr = self.header();
421        if hdr.consumer_active[consumer_idx].load(Ordering::Acquire) == 0 {
422            return Err(BroadcastError::InvalidConsumer);
423        }
424        let my_seq = hdr.consumer_seqs[consumer_idx].load(Ordering::Acquire);
425        let producer = hdr.producer_seq.load(Ordering::Acquire);
426        if my_seq >= producer {
427            self.ring_sidecar
428                .push_op(crate::sidecar_ops::broadcast_ring::OP_RECV, 2); // empty
429            return Err(BroadcastError::Empty);
430        }
431        // SeqLock read the slot at my_seq % capacity.
432        let slot = self.slot(my_seq as usize);
433        loop {
434            let v1 = slot.version.load(Ordering::Acquire);
435            if v1 & 1 != 0 {
436                std::hint::spin_loop();
437                continue;
438            }
439            let src = unsafe {
440                self.raw_ptr
441                    .add(size_of::<BroadcastHeader>())
442                    .add(self.wrap(my_seq as usize) * size_of::<BroadcastSlot>())
443                    .add(std::mem::offset_of!(BroadcastSlot, payload))
444            };
445            let n = out.len().min(BROADCAST_PAYLOAD_BYTES);
446            unsafe { std::ptr::copy_nonoverlapping(src, out.as_mut_ptr(), n); }
447            let v2 = slot.version.load(Ordering::Acquire);
448            if v1 == v2 {
449                hdr.consumer_seqs[consumer_idx].fetch_add(1, Ordering::Release);
450                self.ring_sidecar
451                    .push_op(crate::sidecar_ops::broadcast_ring::OP_RECV, 0);
452                return Ok(n);
453            }
454        }
455    }
456
457    /// Number of messages this consumer has not yet read.
458    pub fn lag(&self, consumer_idx: usize) -> u64 {
459        let hdr = self.header();
460        if consumer_idx >= MAX_CONSUMERS { return 0; }
461        let prod = hdr.producer_seq.load(Ordering::Acquire);
462        let my = hdr.consumer_seqs[consumer_idx].load(Ordering::Acquire);
463        prod.saturating_sub(my)
464    }
465
466    /// Current producer cursor (total messages pushed since creation).
467    pub fn producer_position(&self) -> u64 {
468        self.header().producer_seq.load(Ordering::Acquire)
469    }
470
471    /// Number of currently active consumers.
472    pub fn active_consumer_count(&self) -> usize {
473        let hdr = self.header();
474        (0..MAX_CONSUMERS)
475            .filter(|&i| hdr.consumer_active[i].load(Ordering::Acquire) != 0)
476            .count()
477    }
478
479    pub fn flush(&self) -> Result<(), BroadcastError> {
480        match &self._backing {
481            BroadcastBacking::File(_, mmap) => mmap.flush()?,
482            BroadcastBacking::Anon(mmap) => mmap.flush()?,
483            BroadcastBacking::Shm(_) => {} // RAM-resident; nothing to sync
484        }
485        Ok(())
486    }
487
488    /// Non-blocking flush: schedules a writeback via the OS.
489    /// Note: Windows is only partially async (sync to page cache,
490    /// not to disk).
491    pub fn flush_async(&self) -> Result<(), BroadcastError> {
492        match &self._backing {
493            BroadcastBacking::File(_, mmap) => mmap.flush_async()?,
494            BroadcastBacking::Anon(mmap) => mmap.flush_async()?,
495            BroadcastBacking::Shm(_) => {} // RAM-resident; nothing to sync
496        }
497        Ok(())
498    }
499
500    /// Whether every currently-active consumer has read every
501    /// item the producer has published. Used by the capacity-morph
502    /// wrapper to decide whether a stale broadcast backing can be
503    /// dropped (all subscribers have caught up to the frozen
504    /// producer position).
505    pub fn is_fully_drained(&self) -> bool {
506        let hdr = self.header();
507        let prod = hdr.producer_seq.load(Ordering::Acquire);
508        for i in 0..MAX_CONSUMERS {
509            if hdr.consumer_active[i].load(Ordering::Acquire) != 0
510                && hdr.consumer_seqs[i].load(Ordering::Acquire) < prod
511            {
512                return false;
513            }
514        }
515        true
516    }
517}
518
519#[cfg(test)]
520mod tests {
521    use super::*;
522    use std::sync::Arc;
523    use std::thread;
524    use std::time::Duration;
525
526    fn tmp(name: &str) -> std::path::PathBuf {
527        let mut p = std::env::temp_dir();
528        let pid = std::process::id();
529        p.push(format!("subetha-broadcast-{name}-{pid}.bin"));
530        p
531    }
532
533    fn payload_of(v: u32) -> [u8; BROADCAST_PAYLOAD_BYTES] {
534        let mut b = [0u8; BROADCAST_PAYLOAD_BYTES];
535        b[0..4].copy_from_slice(&v.to_le_bytes());
536        b
537    }
538    fn unpack(b: &[u8]) -> u32 {
539        u32::from_le_bytes(b[0..4].try_into().unwrap())
540    }
541
542    #[test]
543    fn create_initial_state() {
544        let p = tmp("init");
545        let r = SharedBroadcastRing::create(&p, 8).unwrap();
546        assert_eq!(r.capacity(), 8);
547        assert_eq!(r.producer_position(), 0);
548        assert_eq!(r.active_consumer_count(), 0);
549        std::fs::remove_file(&p).ok();
550    }
551
552    #[test]
553    fn one_producer_one_consumer_round_trip() {
554        let p = tmp("1p1c");
555        let r = SharedBroadcastRing::create(&p, 8).unwrap();
556        let c = r.register_consumer().unwrap();
557        r.try_push(&payload_of(42)).unwrap();
558        let mut buf = [0u8; BROADCAST_PAYLOAD_BYTES];
559        let n = r.try_recv(c, &mut buf).unwrap();
560        assert_eq!(n, BROADCAST_PAYLOAD_BYTES);
561        assert_eq!(unpack(&buf), 42);
562        assert_eq!(r.try_recv(c, &mut buf).err(), Some(BroadcastError::Empty));
563        std::fs::remove_file(&p).ok();
564    }
565
566    #[test]
567    fn three_consumers_each_see_all_messages() {
568        let p = tmp("1p3c");
569        let r = SharedBroadcastRing::create(&p, 16).unwrap();
570        let c0 = r.register_consumer().unwrap();
571        let c1 = r.register_consumer().unwrap();
572        let c2 = r.register_consumer().unwrap();
573        for i in 0..5u32 { r.try_push(&payload_of(i * 10)).unwrap(); }
574        let mut buf = [0u8; BROADCAST_PAYLOAD_BYTES];
575        for c in [c0, c1, c2] {
576            for i in 0..5u32 {
577                r.try_recv(c, &mut buf).unwrap();
578                assert_eq!(unpack(&buf), i * 10,
579                    "consumer {c} should see message {i}");
580            }
581            assert_eq!(r.try_recv(c, &mut buf).err(), Some(BroadcastError::Empty));
582        }
583        std::fs::remove_file(&p).ok();
584    }
585
586    #[test]
587    fn lagging_consumer_blocks_producer() {
588        let p = tmp("lag-blocks");
589        let r = SharedBroadcastRing::create(&p, 4).unwrap();
590        let c = r.register_consumer().unwrap();
591        let _c = c;
592        // Fill the ring.
593        for i in 0..4u32 { r.try_push(&payload_of(i)).unwrap(); }
594        // Consumer hasn't read anything; next push should fail Full.
595        assert_eq!(r.try_push(&payload_of(99)).err(), Some(BroadcastError::Full));
596        std::fs::remove_file(&p).ok();
597    }
598
599    #[test]
600    fn unregistering_lagging_consumer_unblocks_producer() {
601        let p = tmp("unreg-unblocks");
602        let r = SharedBroadcastRing::create(&p, 4).unwrap();
603        let c = r.register_consumer().unwrap();
604        for i in 0..4u32 { r.try_push(&payload_of(i)).unwrap(); }
605        assert_eq!(r.try_push(&payload_of(99)).err(), Some(BroadcastError::Full));
606        r.unregister_consumer(c);
607        // With no active consumers, producer can push freely.
608        r.try_push(&payload_of(99)).unwrap();
609        std::fs::remove_file(&p).ok();
610    }
611
612    #[test]
613    fn consumer_registered_late_starts_at_current_producer() {
614        let p = tmp("late-consumer");
615        let r = SharedBroadcastRing::create(&p, 8).unwrap();
616        // Push 3 messages BEFORE registering any consumer.
617        for i in 0..3u32 { r.try_push(&payload_of(i)).unwrap(); }
618        // Now register; this consumer should see only NEW messages.
619        let c = r.register_consumer().unwrap();
620        let mut buf = [0u8; BROADCAST_PAYLOAD_BYTES];
621        assert_eq!(r.try_recv(c, &mut buf).err(), Some(BroadcastError::Empty));
622        r.try_push(&payload_of(100)).unwrap();
623        r.try_recv(c, &mut buf).unwrap();
624        assert_eq!(unpack(&buf), 100);
625        std::fs::remove_file(&p).ok();
626    }
627
628    #[test]
629    fn lag_returns_pending_count() {
630        let p = tmp("lag");
631        let r = SharedBroadcastRing::create(&p, 8).unwrap();
632        let c = r.register_consumer().unwrap();
633        for i in 0..3u32 { r.try_push(&payload_of(i)).unwrap(); }
634        assert_eq!(r.lag(c), 3);
635        let mut buf = [0u8; BROADCAST_PAYLOAD_BYTES];
636        r.try_recv(c, &mut buf).unwrap();
637        assert_eq!(r.lag(c), 2);
638        std::fs::remove_file(&p).ok();
639    }
640
641    #[test]
642    fn max_consumers_returns_no_slot_when_full() {
643        let p = tmp("max-cons");
644        let r = SharedBroadcastRing::create(&p, 4).unwrap();
645        for _ in 0..MAX_CONSUMERS {
646            r.register_consumer().unwrap();
647        }
648        assert_eq!(r.register_consumer().err(), Some(BroadcastError::NoConsumerSlot));
649        std::fs::remove_file(&p).ok();
650    }
651
652    #[test]
653    fn cross_handle_pub_sub() {
654        let p = tmp("cross-handle");
655        let pub_handle = SharedBroadcastRing::create(&p, 8).unwrap();
656        let sub_handle = SharedBroadcastRing::open(&p, 8).unwrap();
657        let c = sub_handle.register_consumer().unwrap();
658        pub_handle.try_push(&payload_of(7777)).unwrap();
659        let mut buf = [0u8; BROADCAST_PAYLOAD_BYTES];
660        sub_handle.try_recv(c, &mut buf).unwrap();
661        assert_eq!(unpack(&buf), 7777);
662        std::fs::remove_file(&p).ok();
663    }
664
665    #[test]
666    fn payload_too_large_rejected_at_push() {
667        let p = tmp("oversized");
668        let r = SharedBroadcastRing::create(&p, 4).unwrap();
669        let _c = r.register_consumer().unwrap();
670        let big = vec![0u8; BROADCAST_PAYLOAD_BYTES + 1];
671        assert_eq!(r.try_push(&big).err(), Some(BroadcastError::PayloadTooLarge));
672        std::fs::remove_file(&p).ok();
673    }
674
675    #[test]
676    fn concurrent_consumers_all_drain_correctly() {
677        let p = tmp("concurrent");
678        let r = Arc::new(SharedBroadcastRing::create(&p, 256).unwrap());
679        let n_consumers = 4;
680        let n_msgs = 100u32;
681        let consumer_ids: Vec<usize> = (0..n_consumers).map(|_| r.register_consumer().unwrap()).collect();
682
683        let r_p = r.clone();
684        let producer = thread::spawn(move || {
685            for i in 0..n_msgs {
686                while r_p.try_push(&payload_of(i)).is_err() {
687                    thread::yield_now();
688                }
689            }
690        });
691
692        let mut handles = vec![];
693        for &c in &consumer_ids {
694            let r = r.clone();
695            handles.push(thread::spawn(move || {
696                let mut received = vec![];
697                let mut buf = [0u8; BROADCAST_PAYLOAD_BYTES];
698                while received.len() < n_msgs as usize {
699                    match r.try_recv(c, &mut buf) {
700                        Ok(_) => received.push(unpack(&buf)),
701                        Err(BroadcastError::Empty) => thread::yield_now(),
702                        Err(e) => panic!("unexpected error: {e:?}"),
703                    }
704                }
705                received
706            }));
707        }
708        producer.join().unwrap();
709        for h in handles {
710            let got = h.join().unwrap();
711            assert_eq!(got.len(), n_msgs as usize);
712            for (i, v) in got.iter().enumerate() {
713                assert_eq!(*v, i as u32, "message order must be preserved");
714            }
715        }
716        std::fs::remove_file(&p).ok();
717    }
718
719    #[test]
720    fn slow_consumer_doesnt_break_fast_consumer() {
721        let p = tmp("slow-fast");
722        let r = Arc::new(SharedBroadcastRing::create(&p, 8).unwrap());
723        let _slow = r.register_consumer().unwrap();
724        let fast = r.register_consumer().unwrap();
725
726        for i in 0..4u32 { r.try_push(&payload_of(i)).unwrap(); }
727        // Fast consumer drains; slow consumer hasn't moved.
728        let mut buf = [0u8; BROADCAST_PAYLOAD_BYTES];
729        for i in 0..4u32 {
730            r.try_recv(fast, &mut buf).unwrap();
731            assert_eq!(unpack(&buf), i);
732        }
733        // Producer is still bounded by slow consumer's cursor though;
734        // try to push 5 more - should hit Full at some point.
735        let mut pushed = 0;
736        for i in 100..200u32 {
737            if r.try_push(&payload_of(i)).is_err() { break; }
738            pushed += 1;
739        }
740        assert!(pushed <= 4, "should be bounded by slow consumer; pushed {pushed}");
741        std::fs::remove_file(&p).ok();
742    }
743
744    #[test]
745    fn observer_can_wait_for_publisher() {
746        let p = tmp("wait");
747        let r = Arc::new(SharedBroadcastRing::create(&p, 8).unwrap());
748        let c = r.register_consumer().unwrap();
749
750        let r_p = r.clone();
751        let pusher = thread::spawn(move || {
752            thread::sleep(Duration::from_millis(20));
753            r_p.try_push(&payload_of(555)).unwrap();
754        });
755        let mut buf = [0u8; BROADCAST_PAYLOAD_BYTES];
756        let r_c = r.clone();
757        let consumer = thread::spawn(move || {
758            loop {
759                match r_c.try_recv(c, &mut buf) {
760                    Ok(_) => break unpack(&buf),
761                    Err(BroadcastError::Empty) => thread::yield_now(),
762                    Err(e) => panic!("{e:?}"),
763                }
764            }
765        });
766        pusher.join().unwrap();
767        assert_eq!(consumer.join().unwrap(), 555);
768        std::fs::remove_file(&p).ok();
769    }
770
771    #[test]
772    fn disk_persistence_survives_reopen() {
773        let p = tmp("disk");
774        {
775            let r = SharedBroadcastRing::create(&p, 8).unwrap();
776            let _c = r.register_consumer().unwrap();
777            r.try_push(&payload_of(1234)).unwrap();
778            r.try_push(&payload_of(5678)).unwrap();
779            r.flush().unwrap();
780        }
781        let r2 = SharedBroadcastRing::open(&p, 8).unwrap();
782        // Producer seq should be 2; consumer 0 should still be at 0.
783        assert_eq!(r2.producer_position(), 2);
784        assert_eq!(r2.lag(0), 2);
785        let mut buf = [0u8; BROADCAST_PAYLOAD_BYTES];
786        r2.try_recv(0, &mut buf).unwrap();
787        assert_eq!(unpack(&buf), 1234);
788        r2.try_recv(0, &mut buf).unwrap();
789        assert_eq!(unpack(&buf), 5678);
790        std::fs::remove_file(&p).ok();
791    }
792}