1use 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#[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#[allow(dead_code)]
143enum BroadcastBacking {
144 Anon(MmapMut),
146 File(File, MmapMut),
148 Shm(crate::shm_file::ShmFile),
151}
152
153unsafe 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
180unsafe 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 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 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 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 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 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 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 #[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 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); Err(BroadcastError::NoConsumerSlot)
392 }
393
394 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 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 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); return Err(BroadcastError::Full);
437 }
438 let slot = self.slot(producer as usize);
440 slot.version.fetch_add(1, Ordering::AcqRel); 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); 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 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); return Err(BroadcastError::Empty);
473 }
474 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 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 pub fn producer_position(&self) -> u64 {
511 self.header().producer_seq.load(Ordering::Acquire)
512 }
513
514 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(_) => {} }
528 Ok(())
529 }
530
531 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(_) => {} }
540 Ok(())
541 }
542
543 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 #[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 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 for i in 0..4u32 { r.try_push(&payload_of(i)).unwrap(); }
663 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 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 for i in 0..3u32 { r.try_push(&payload_of(i)).unwrap(); }
687 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 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 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 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}