const PT_PSFB: u8 = 206;
pub fn build_pli(sender_ssrc: u32, media_ssrc: u32) -> Vec<u8> {
let mut p = Vec::with_capacity(12);
p.push(0x80 | 1);
p.push(PT_PSFB);
p.extend_from_slice(&2u16.to_be_bytes());
p.extend_from_slice(&sender_ssrc.to_be_bytes());
p.extend_from_slice(&media_ssrc.to_be_bytes());
p
}
pub fn build_fir(sender_ssrc: u32, media_ssrc: u32, seq_nr: u8) -> Vec<u8> {
let mut p = Vec::with_capacity(20);
p.push(0x80 | 4);
p.push(PT_PSFB);
p.extend_from_slice(&4u16.to_be_bytes());
p.extend_from_slice(&sender_ssrc.to_be_bytes());
p.extend_from_slice(&0u32.to_be_bytes()); p.extend_from_slice(&media_ssrc.to_be_bytes()); p.push(seq_nr);
p.extend_from_slice(&[0, 0, 0]); p
}
pub fn build_sr(
sender_ssrc: u32,
ntp: u64,
rtp_ts: u32,
packet_count: u32,
octet_count: u32,
) -> Vec<u8> {
let mut p = Vec::with_capacity(28);
p.push(0x80); p.push(PT_SR);
p.extend_from_slice(&6u16.to_be_bytes()); p.extend_from_slice(&sender_ssrc.to_be_bytes());
p.extend_from_slice(&ntp.to_be_bytes()); p.extend_from_slice(&rtp_ts.to_be_bytes());
p.extend_from_slice(&packet_count.to_be_bytes());
p.extend_from_slice(&octet_count.to_be_bytes());
p
}
pub fn ntp_now() -> u64 {
const NTP_UNIX_OFFSET: u64 = 2_208_988_800;
let d = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default();
let secs = d.as_secs() + NTP_UNIX_OFFSET;
let frac = ((d.subsec_nanos() as u64) << 32) / 1_000_000_000;
(secs << 32) | frac
}
const PT_RTPFB: u8 = 205;
const PT_RR: u8 = 201;
const PT_SR: u8 = 200;
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum RtcpFeedback {
Pli {
sender_ssrc: u32,
media_ssrc: u32,
},
Fir {
sender_ssrc: u32,
media_ssrc: u32,
},
Nack {
sender_ssrc: u32,
media_ssrc: u32,
lost: Vec<u16>,
},
Remb {
sender_ssrc: u32,
bitrate_bps: u64,
},
ReceiverReport {
ssrc: u32,
fraction_lost: u8,
cumulative_lost: u32,
jitter: u32,
},
}
pub fn parse_compound(mut buf: &[u8]) -> Vec<RtcpFeedback> {
let mut out = Vec::new();
while buf.len() >= 4 {
let version = buf[0] >> 6;
if version != 2 {
break;
}
let fmt = buf[0] & 0x1f;
let pt = buf[1];
let len_words = u16::from_be_bytes([buf[2], buf[3]]) as usize;
let pkt_len = (len_words + 1) * 4;
if pkt_len == 0 || pkt_len > buf.len() {
break;
}
let pkt = &buf[..pkt_len];
match (pt, fmt) {
(PT_PSFB, 1) if pkt_len >= 12 => out.push(RtcpFeedback::Pli {
sender_ssrc: be32(&pkt[4..]),
media_ssrc: be32(&pkt[8..]),
}),
(PT_PSFB, 4) if pkt_len >= 16 => out.push(RtcpFeedback::Fir {
sender_ssrc: be32(&pkt[4..]),
media_ssrc: be32(&pkt[12..]),
}),
(PT_PSFB, 15) if pkt_len >= 20 && &pkt[12..16] == b"REMB" => out.push(parse_remb(pkt)),
(PT_RTPFB, 1) if pkt_len >= 12 => out.push(parse_nack(pkt)),
(PT_RR, _) if pkt_len >= 8 => parse_report_blocks(&pkt[8..], fmt, &mut out),
(PT_SR, _) if pkt_len >= 28 => parse_report_blocks(&pkt[28..], fmt, &mut out),
_ => {}
}
buf = &buf[pkt_len..];
}
out
}
fn parse_nack(pkt: &[u8]) -> RtcpFeedback {
let sender_ssrc = be32(&pkt[4..]);
let media_ssrc = be32(&pkt[8..]);
let mut lost = Vec::new();
let mut off = 12;
while off + 4 <= pkt.len() {
let pid = u16::from_be_bytes([pkt[off], pkt[off + 1]]);
let blp = u16::from_be_bytes([pkt[off + 2], pkt[off + 3]]);
lost.push(pid);
for i in 0..16 {
if blp & (1 << i) != 0 {
lost.push(pid.wrapping_add(i + 1));
}
}
off += 4;
}
RtcpFeedback::Nack {
sender_ssrc,
media_ssrc,
lost,
}
}
fn parse_remb(pkt: &[u8]) -> RtcpFeedback {
let sender_ssrc = be32(&pkt[4..]);
let exp = (pkt[17] >> 2) as u32;
let mantissa = (((pkt[17] & 0x03) as u64) << 16) | ((pkt[18] as u64) << 8) | pkt[19] as u64;
let bitrate_bps = mantissa.checked_shl(exp).unwrap_or(u64::MAX);
RtcpFeedback::Remb {
sender_ssrc,
bitrate_bps,
}
}
fn parse_report_blocks(blocks: &[u8], count: u8, out: &mut Vec<RtcpFeedback>) {
for i in 0..count as usize {
let off = i * 24;
if off + 24 > blocks.len() {
break;
}
let b = &blocks[off..off + 24];
out.push(RtcpFeedback::ReceiverReport {
ssrc: be32(b),
fraction_lost: b[4],
cumulative_lost: u32::from_be_bytes([0, b[5], b[6], b[7]]),
jitter: be32(&b[12..]),
});
}
}
#[inline]
fn be32(b: &[u8]) -> u32 {
u32::from_be_bytes([b[0], b[1], b[2], b[3]])
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn pli_has_correct_header_and_ssrcs() {
let p = build_pli(0x1111_1111, 0x2222_2222);
assert_eq!(p.len(), 12);
assert_eq!(p[0], 0x81); assert_eq!(p[1], 206); assert_eq!(u16::from_be_bytes([p[2], p[3]]), 2); assert_eq!(&p[4..8], &0x1111_1111u32.to_be_bytes());
assert_eq!(&p[8..12], &0x2222_2222u32.to_be_bytes());
}
#[test]
fn fir_carries_seq_and_target_ssrc() {
let p = build_fir(1, 0xDEAD_BEEF, 7);
assert_eq!(p.len(), 20);
assert_eq!(p[0], 0x84); assert_eq!(&p[12..16], &0xDEAD_BEEFu32.to_be_bytes());
assert_eq!(p[16], 7); }
#[test]
fn parses_our_own_pli_round_trip() {
let p = build_pli(0xAAAA_AAAA, 0xBBBB_BBBB);
assert_eq!(
parse_compound(&p),
vec![RtcpFeedback::Pli {
sender_ssrc: 0xAAAA_AAAA,
media_ssrc: 0xBBBB_BBBB,
}]
);
}
#[test]
fn parses_fir_round_trip() {
let p = build_fir(1, 0xDEAD_BEEF, 3);
assert_eq!(
parse_compound(&p),
vec![RtcpFeedback::Fir {
sender_ssrc: 1,
media_ssrc: 0xDEAD_BEEF,
}]
);
}
#[test]
fn parses_generic_nack_with_bitmask() {
let mut p = Vec::new();
p.push(0x80 | 1);
p.push(PT_RTPFB);
p.extend_from_slice(&3u16.to_be_bytes()); p.extend_from_slice(&1u32.to_be_bytes());
p.extend_from_slice(&2u32.to_be_bytes());
p.extend_from_slice(&100u16.to_be_bytes());
p.extend_from_slice(&0b0000_0000_0000_0101u16.to_be_bytes());
assert_eq!(
parse_compound(&p),
vec![RtcpFeedback::Nack {
sender_ssrc: 1,
media_ssrc: 2,
lost: vec![100, 101, 103],
}]
);
}
#[test]
fn parses_receiver_report_block() {
let mut p = Vec::new();
p.push(0x80 | 1); p.push(PT_RR);
p.extend_from_slice(&7u16.to_be_bytes()); p.extend_from_slice(&0xCAFEu32.to_be_bytes()); p.extend_from_slice(&0x1234_5678u32.to_be_bytes()); p.push(64); p.extend_from_slice(&[0, 0, 10]); p.extend_from_slice(&55u32.to_be_bytes()); p.extend_from_slice(&99u32.to_be_bytes()); p.extend_from_slice(&0u32.to_be_bytes()); p.extend_from_slice(&0u32.to_be_bytes()); assert_eq!(
parse_compound(&p),
vec![RtcpFeedback::ReceiverReport {
ssrc: 0x1234_5678,
fraction_lost: 64,
cumulative_lost: 10,
jitter: 99,
}]
);
}
#[test]
fn sr_has_sender_info_layout() {
let p = build_sr(0xABCD_1234, 0x1122_3344_5566_7788, 90_000, 5, 1000);
assert_eq!(p.len(), 28);
assert_eq!(p[0], 0x80); assert_eq!(p[1], 200); assert_eq!(u16::from_be_bytes([p[2], p[3]]), 6);
assert_eq!(&p[4..8], &0xABCD_1234u32.to_be_bytes());
assert_eq!(&p[8..16], &0x1122_3344_5566_7788u64.to_be_bytes());
assert_eq!(&p[16..20], &90_000u32.to_be_bytes());
assert_eq!(&p[20..24], &5u32.to_be_bytes());
assert_eq!(&p[24..28], &1000u32.to_be_bytes());
}
#[test]
fn ntp_now_is_after_2020() {
assert!((ntp_now() >> 32) > 3_786_825_600);
}
#[test]
fn parses_remb_bitrate() {
let mut p = Vec::new();
p.push(0x80 | 15);
p.push(PT_PSFB);
p.extend_from_slice(&5u16.to_be_bytes()); p.extend_from_slice(&9u32.to_be_bytes());
p.extend_from_slice(&0u32.to_be_bytes());
p.extend_from_slice(b"REMB");
p.push(1); let mantissa: u32 = 1000;
let exp: u32 = 3;
p.push(((exp << 2) as u8) | ((mantissa >> 16) as u8 & 0x03));
p.push((mantissa >> 8) as u8);
p.push(mantissa as u8);
p.extend_from_slice(&0x1234u32.to_be_bytes()); assert_eq!(
parse_compound(&p),
vec![RtcpFeedback::Remb {
sender_ssrc: 9,
bitrate_bps: 8000,
}]
);
}
#[test]
fn walks_compound_and_skips_unknown() {
let mut p = build_pli(1, 2);
p.push(0x80);
p.push(202); p.extend_from_slice(&1u16.to_be_bytes()); p.extend_from_slice(&0u32.to_be_bytes());
p.extend_from_slice(&build_pli(3, 4));
assert_eq!(
parse_compound(&p),
vec![
RtcpFeedback::Pli {
sender_ssrc: 1,
media_ssrc: 2
},
RtcpFeedback::Pli {
sender_ssrc: 3,
media_ssrc: 4
},
]
);
}
#[test]
fn tolerates_truncated_tail() {
let mut p = build_pli(1, 2);
p.extend_from_slice(&[0x80, 206]); assert_eq!(parse_compound(&p).len(), 1);
}
}