use frust_devtools_protocol::FrameStats;
use tokio::sync::broadcast;
pub(crate) const DEFAULT_CAPACITY: usize = 64;
pub(crate) struct FrameStatsBus {
tx: broadcast::Sender<FrameStats>,
}
impl FrameStatsBus {
pub(crate) fn new(capacity: usize) -> Self {
let (tx, _) = broadcast::channel(capacity.max(1));
Self { tx }
}
pub(crate) fn publish(&self, stats: FrameStats) {
let _ = self.tx.send(stats);
}
pub(crate) fn subscribe(&self) -> broadcast::Receiver<FrameStats> {
self.tx.subscribe()
}
}
#[cfg(test)]
mod tests {
use super::*;
fn stats(n: u64) -> FrameStats {
FrameStats {
n,
total_us: 1_000,
rebuild_us: 100,
layout_us: 200,
paint_us: 300,
encode_us: 150,
acquire_us: 50,
submit_us: 200,
skipped: false,
}
}
#[test]
fn publish_with_no_subscribers_is_a_no_op() {
let bus = FrameStatsBus::new(4);
for n in 0..100 {
bus.publish(stats(n));
}
}
#[test]
fn a_full_queue_drops_the_oldest_frames_and_keeps_the_newest() {
let bus = FrameStatsBus::new(4);
let mut rx = bus.subscribe();
for n in 0..10 {
bus.publish(stats(n));
}
match rx.try_recv() {
Err(broadcast::error::TryRecvError::Lagged(missed)) => assert_eq!(missed, 6),
other => panic!("expected a Lagged report, got {other:?}"),
}
for n in 6..10 {
assert_eq!(rx.try_recv().map(|s| s.n), Ok(n));
}
assert!(matches!(
rx.try_recv(),
Err(broadcast::error::TryRecvError::Empty)
));
}
#[test]
fn every_subscriber_gets_every_frame_while_it_keeps_up() {
let bus = FrameStatsBus::new(4);
let mut a = bus.subscribe();
let mut b = bus.subscribe();
bus.publish(stats(1));
bus.publish(stats(2));
assert_eq!(a.try_recv().map(|s| s.n), Ok(1));
assert_eq!(b.try_recv().map(|s| s.n), Ok(1));
assert_eq!(a.try_recv().map(|s| s.n), Ok(2));
assert_eq!(b.try_recv().map(|s| s.n), Ok(2));
}
#[test]
fn a_late_subscriber_sees_only_frames_published_after_it_joined() {
let bus = FrameStatsBus::new(4);
bus.publish(stats(1));
let mut rx = bus.subscribe();
bus.publish(stats(2));
assert_eq!(rx.try_recv().map(|s| s.n), Ok(2));
}
#[test]
fn zero_capacity_is_clamped_rather_than_panicking() {
let bus = FrameStatsBus::new(0);
let mut rx = bus.subscribe();
bus.publish(stats(1));
assert_eq!(rx.try_recv().map(|s| s.n), Ok(1));
}
}