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
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 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 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 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 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 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 #[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 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); Err(BroadcastError::NoConsumerSlot)
349 }
350
351 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 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 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); return Err(BroadcastError::Full);
394 }
395 let slot = self.slot(producer as usize);
397 slot.version.fetch_add(1, Ordering::AcqRel); 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); 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 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); return Err(BroadcastError::Empty);
430 }
431 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 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 pub fn producer_position(&self) -> u64 {
468 self.header().producer_seq.load(Ordering::Acquire)
469 }
470
471 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(_) => {} }
485 Ok(())
486 }
487
488 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(_) => {} }
497 Ok(())
498 }
499
500 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 for i in 0..4u32 { r.try_push(&payload_of(i)).unwrap(); }
594 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 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 for i in 0..3u32 { r.try_push(&payload_of(i)).unwrap(); }
618 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 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 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 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}