use rtc_interceptor::{
Interceptor, JitterBufferBuilder, NackGeneratorBuilder, NoopInterceptor, Packet, RTCPFeedback,
Registry, StreamInfo, TaggedPacket, interceptor,
};
use sansio::Protocol;
use shared::TransportContext;
use shared::error::Error;
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
#[derive(Interceptor)]
struct Marker<P> {
#[next]
inner: P,
released: Arc<Mutex<Vec<u16>>>,
}
#[interceptor]
impl<P: Interceptor> Marker<P> {
#[overrides]
fn handle_read(&mut self, msg: TaggedPacket) -> Result<(), Self::Error> {
if let Packet::Rtp(rtp) = &msg.message {
self.released
.lock()
.unwrap()
.push(rtp.header.sequence_number);
}
self.inner.handle_read(msg)
}
}
const SSRC: u32 = 1;
const CLOCK: u32 = 90_000;
const NACK_INTERVAL: Duration = Duration::from_millis(100);
const ROUND_TRIP: Duration = Duration::from_millis(60);
const SPACING: Duration = Duration::from_millis(20);
fn ms(milliseconds: u64) -> Duration {
Duration::from_millis(milliseconds)
}
fn ticks(duration: Duration) -> u32 {
(duration.as_secs_f64() * f64::from(CLOCK)) as u32
}
struct Harness {
chain: Box<dyn Interceptor>,
released: Arc<Mutex<Vec<u16>>>,
epoch: Instant,
}
impl Harness {
fn new(depth: Duration) -> Self {
let epoch = Instant::now();
let released = Arc::new(Mutex::new(Vec::new()));
let marker_released = Arc::clone(&released);
let chain = Registry::new()
.with(move |inner: NoopInterceptor| Marker {
inner,
released: marker_released,
})
.with(JitterBufferBuilder::new().with_depth(depth).build())
.with(
NackGeneratorBuilder::new()
.with_interval(NACK_INTERVAL)
.with_skip_last_n(0)
.build(),
)
.boxed()
.build();
let mut harness = Self {
chain,
released,
epoch,
};
harness.bind();
harness
}
fn bind(&mut self) {
self.chain.bind_remote_stream(&StreamInfo {
ssrc: SSRC,
clock_rate: CLOCK,
rtcp_feedback: vec![RTCPFeedback {
typ: "nack".to_owned(),
parameter: String::new(),
}],
..Default::default()
});
}
fn arrive(&mut self, at: Duration, sequence_number: u16, media_time: Duration) {
self.advance_to(at);
self.chain
.handle_read(TaggedPacket {
now: self.epoch + at,
transport: TransportContext::default(),
message: Packet::Rtp(rtp::Packet {
header: rtp::header::Header {
ssrc: SSRC,
sequence_number,
timestamp: ticks(media_time),
..Default::default()
},
..Default::default()
}),
})
.expect("handle_read");
}
fn tick(&mut self, at: Duration) {
self.chain
.handle_timeout(self.epoch + at)
.expect("handle_timeout");
}
fn advance_to(&mut self, at: Duration) {
let deadline = self.epoch + at;
while let Some(next) = self.chain.poll_timeout() {
if next > deadline {
break;
}
self.chain.handle_timeout(next).expect("handle_timeout");
}
self.tick(at);
}
fn drain_nacks(&mut self) -> Vec<u16> {
let mut requested = Vec::new();
while let Some(packet) = self.chain.poll_write() {
if let Packet::Rtcp(rtcp_packets) = &packet.message {
for rtcp_packet in rtcp_packets {
if let Some(nack) = rtcp_packet
.as_any()
.downcast_ref::<rtcp::transport_feedbacks::transport_layer_nack::TransportLayerNack>(
)
{
for pair in &nack.nacks {
requested.extend(pair.packet_list());
}
}
}
}
}
requested
}
fn released(&self) -> Vec<u16> {
self.released.lock().unwrap().clone()
}
}
fn run_until_retransmission(harness: &mut Harness) -> Duration {
harness.arrive(ms(0), 0, Duration::ZERO);
harness.arrive(ms(20), 1, SPACING * 1);
harness.arrive(ms(60), 3, SPACING * 3);
harness.arrive(ms(80), 4, SPACING * 4);
let nack_at = NACK_INTERVAL;
harness.advance_to(nack_at);
let requested = harness.drain_nacks();
assert!(
requested.contains(&2),
"precondition: the generator asked for the lost packet, got {requested:?}"
);
nack_at + ROUND_TRIP
}
#[test]
fn a_deep_enough_buffer_still_has_a_place_for_the_retransmission() {
let depth = NACK_INTERVAL + ROUND_TRIP + ms(40);
let mut harness = Harness::new(depth);
let retransmission_at = run_until_retransmission(&mut harness);
assert!(
harness.released().is_empty(),
"nothing has been played out yet, so 2's position is still open"
);
harness.arrive(retransmission_at, 2, SPACING * 2);
harness.advance_to(ms(1000));
assert_eq!(
vec![0, 1, 2, 3, 4],
harness.released(),
"the recovered packet took its place in order"
);
}
#[test]
fn a_buffer_shallower_than_the_recovery_window_cannot_use_the_retransmission() {
let depth = ms(20);
assert!(
depth < NACK_INTERVAL + ROUND_TRIP,
"the point of this test is a depth below the recovery window"
);
let mut harness = Harness::new(depth);
let retransmission_at = run_until_retransmission(&mut harness);
assert_eq!(
vec![0, 1, 3, 4],
harness.released(),
"playout has already moved past 2's position with the gap unfilled"
);
harness.arrive(retransmission_at, 2, SPACING * 2);
harness.advance_to(ms(1000));
assert_eq!(
vec![0, 1, 3, 4],
harness.released(),
"the retransmission is dropped: emitting 2 after 3 and 4 would break ordering, which is \
the one thing the buffer exists to preserve"
);
}
#[test]
fn the_boundary_is_detection_plus_the_round_trip() {
let window = NACK_INTERVAL + ROUND_TRIP;
for (depth, expected, label) in [
(window + ms(40), vec![0u16, 1, 2, 3, 4], "comfortably above"),
(window / 2, vec![0, 1, 3, 4], "half the window"),
] {
let mut harness = Harness::new(depth);
let retransmission_at = run_until_retransmission(&mut harness);
harness.arrive(retransmission_at, 2, SPACING * 2);
harness.advance_to(ms(1000));
assert_eq!(
expected,
harness.released(),
"depth {depth:?} ({label}) against a {window:?} recovery window"
);
}
}
#[test]
fn a_generator_below_the_buffer_detects_loss_a_whole_depth_late() {
let depth = ms(200);
let epoch = Instant::now();
let mut chain = Registry::new()
.with(|inner: NoopInterceptor| {
NackGeneratorBuilder::new()
.with_interval(NACK_INTERVAL)
.with_skip_last_n(0)
.build()(inner)
})
.with(JitterBufferBuilder::new().with_depth(depth).build())
.boxed()
.build();
chain.bind_remote_stream(&StreamInfo {
ssrc: SSRC,
clock_rate: CLOCK,
rtcp_feedback: vec![RTCPFeedback {
typ: "nack".to_owned(),
parameter: String::new(),
}],
..Default::default()
});
let arrive = |chain: &mut Box<dyn Interceptor>, at: Duration, sequence_number: u16| {
chain
.handle_read(TaggedPacket {
now: epoch + at,
transport: TransportContext::default(),
message: Packet::Rtp(rtp::Packet {
header: rtp::header::Header {
ssrc: SSRC,
sequence_number,
timestamp: ticks(SPACING * u32::from(sequence_number)),
..Default::default()
},
..Default::default()
}),
})
.expect("handle_read");
};
arrive(&mut chain, ms(0), 0);
arrive(&mut chain, ms(20), 1);
arrive(&mut chain, ms(60), 3);
chain
.handle_timeout(epoch + NACK_INTERVAL)
.expect("handle_timeout");
let mut asked = false;
while let Some(packet) = chain.poll_write() {
if let Packet::Rtcp(rtcp_packets) = &packet.message {
asked |= !rtcp_packets.is_empty();
}
}
assert!(
!asked,
"the generator below the buffer has not seen a single packet yet, so it cannot have \
noticed the gap — every NACK it eventually sends is late by the depth"
);
}