use bytes::Bytes;
use media_plane::{ByteTap, TapItem, TapPoint};
const FIXTURE: &str = concat!(env!("CARGO_MANIFEST_DIR"), "/../fixtures/ts/h264_aac.ts");
const TS_PACKET_SIZE: usize = 188;
const TAP_CAPACITY: usize = 32;
const CONSUMER_POLL_STRIDE: usize = 64;
use broadcast_common::stage::Timestamp;
fn main() {
let bytes = std::fs::read(FIXTURE).expect("read the committed real-capture TS fixture");
assert_eq!(
bytes.len() % TS_PACKET_SIZE,
0,
"fixture must be a whole number of 188-byte TS packets"
);
let packet_count = bytes.len() / TS_PACKET_SIZE;
println!(
"byte_tap_wire_observer: read {} bytes ({packet_count} TS packets) from {FIXTURE}",
bytes.len()
);
let mut tap = ByteTap::new(TapPoint::Wire, TAP_CAPACITY);
let mut delivered = 0u64;
let mut lagged_total = 0u64;
let mut lagged_reports = 0u64;
for (i, packet) in bytes.chunks(TS_PACKET_SIZE).enumerate() {
tap.record(
Bytes::copy_from_slice(packet),
Timestamp::from_nanos(i as u64),
);
if i % CONSUMER_POLL_STRIDE == 0 {
while let Some(item) = tap.poll() {
match item {
TapItem::Data(_, _) => delivered += 1,
TapItem::Lagged { skipped } => {
lagged_total += skipped;
lagged_reports += 1;
}
other => panic!("unhandled TapItem variant: {other:?}"),
}
}
}
}
while let Some(item) = tap.poll() {
match item {
TapItem::Data(_, _) => delivered += 1,
TapItem::Lagged { skipped } => {
lagged_total += skipped;
lagged_reports += 1;
}
other => panic!("unhandled TapItem variant: {other:?}"),
}
}
println!(
"byte_tap_wire_observer: {delivered} packet(s) delivered, {lagged_reports} Lagged report(s) totalling {lagged_total} skipped packet(s)"
);
assert_eq!(
delivered + lagged_total,
packet_count as u64,
"every real packet must be accounted for as delivered or skipped"
);
assert!(
lagged_total > 0,
"this run's ring/stride are chosen so a real capture this size must lag at least once"
);
}