use alloc::collections::VecDeque;
use broadcast_common::stage::Timestamp;
use bytes::Bytes;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
#[non_exhaustive]
pub enum TapPoint {
Wire,
PostTransform,
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum TapItem {
Data(Bytes, Timestamp),
Lagged {
skipped: u64,
},
}
pub struct ByteTap {
point: TapPoint,
capacity: usize,
ring: VecDeque<(Bytes, Timestamp)>,
skipped_since_last_poll: u64,
}
impl ByteTap {
pub fn new(point: TapPoint, capacity: usize) -> Self {
assert!(capacity > 0, "ByteTap capacity must be > 0");
ByteTap {
point,
capacity,
ring: VecDeque::with_capacity(capacity),
skipped_since_last_poll: 0,
}
}
pub fn point(&self) -> TapPoint {
self.point
}
pub fn capacity(&self) -> usize {
self.capacity
}
pub fn len(&self) -> usize {
self.ring.len()
}
pub fn is_empty(&self) -> bool {
self.ring.is_empty()
}
pub fn record(&mut self, bytes: Bytes, at: Timestamp) {
if self.ring.len() >= self.capacity {
self.ring.pop_front();
self.skipped_since_last_poll = self.skipped_since_last_poll.saturating_add(1);
}
self.ring.push_back((bytes, at));
}
pub fn poll(&mut self) -> Option<TapItem> {
if self.skipped_since_last_poll > 0 {
let skipped = self.skipped_since_last_poll;
self.skipped_since_last_poll = 0;
return Some(TapItem::Lagged { skipped });
}
self.ring.pop_front().map(|(b, t)| TapItem::Data(b, t))
}
}
#[cfg(test)]
mod tests {
use super::*;
use alloc::vec::Vec;
#[test]
fn tap_yields_bytes_a_demuxer_would_reject() {
let mut tap = ByteTap::new(TapPoint::Wire, 4);
let malformed = Bytes::from_static(&[0x00, 0x80, 0xFF, 0xFF, 0xFF]);
tap.record(malformed.clone(), Timestamp::from_nanos(42));
assert_eq!(
tap.poll(),
Some(TapItem::Data(malformed, Timestamp::from_nanos(42)))
);
assert_eq!(tap.poll(), None);
}
#[test]
fn slow_consumer_gets_accurate_lagged_and_producer_is_never_blocked() {
let capacity = 4;
let mut tap = ByteTap::new(TapPoint::Wire, capacity);
let total_records: u64 = 1_000;
for i in 0..total_records {
tap.record(
Bytes::copy_from_slice(&i.to_be_bytes()),
Timestamp::from_nanos(i),
);
}
assert_eq!(tap.len(), capacity);
let expected_skipped = total_records - capacity as u64;
assert_eq!(
tap.poll(),
Some(TapItem::Lagged {
skipped: expected_skipped
})
);
let mut drained = Vec::new();
while let Some(item) = tap.poll() {
drained.push(item);
}
assert_eq!(drained.len(), capacity);
for (offset, item) in drained.into_iter().enumerate() {
let expected_index = total_records - capacity as u64 + offset as u64;
assert_eq!(
item,
TapItem::Data(
Bytes::copy_from_slice(&expected_index.to_be_bytes()),
Timestamp::from_nanos(expected_index)
)
);
}
}
#[test]
fn tap_ring_is_bounded_under_flood() {
let capacity = 8;
let mut tap = ByteTap::new(TapPoint::PostTransform, capacity);
for i in 0..200_000u64 {
tap.record(Bytes::from_static(b"x"), Timestamp::from_nanos(i));
assert!(tap.len() <= capacity, "ring exceeded capacity mid-flood");
}
assert_eq!(tap.len(), capacity);
assert_eq!(tap.capacity(), capacity);
}
#[test]
fn empty_tap_polls_none() {
let mut tap = ByteTap::new(TapPoint::Wire, 2);
assert_eq!(tap.poll(), None);
assert!(tap.is_empty());
}
#[test]
#[should_panic(expected = "capacity must be > 0")]
fn zero_capacity_panics() {
let _ = ByteTap::new(TapPoint::Wire, 0);
}
}