use bytes::{Buf, BufMut, Bytes, BytesMut};
use std::time::{SystemTime, UNIX_EPOCH};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[repr(u8)]
pub enum RtcpType {
SenderReport = 200,
ReceiverReport = 201,
SourceDescription = 202,
Goodbye = 203,
ApplicationDefined = 204,
TransportFeedback = 205,
PayloadFeedback = 206,
}
impl TryFrom<u8> for RtcpType {
type Error = RtcpParseError;
fn try_from(value: u8) -> Result<Self, Self::Error> {
match value {
200 => Ok(RtcpType::SenderReport),
201 => Ok(RtcpType::ReceiverReport),
202 => Ok(RtcpType::SourceDescription),
203 => Ok(RtcpType::Goodbye),
204 => Ok(RtcpType::ApplicationDefined),
205 => Ok(RtcpType::TransportFeedback),
206 => Ok(RtcpType::PayloadFeedback),
_ => Err(RtcpParseError::UnknownPacketType(value)),
}
}
}
#[derive(Debug, Clone, thiserror::Error)]
pub enum RtcpParseError {
#[error("Packet too short: {0} bytes")]
TooShort(usize),
#[error("Invalid RTCP version: {0}")]
InvalidVersion(u8),
#[error("Unknown packet type: {0}")]
UnknownPacketType(u8),
#[error("Invalid report block count")]
InvalidReportCount,
}
#[derive(Debug, Clone)]
pub struct RtcpHeader {
pub version: u8,
pub padding: bool,
pub count: u8,
pub packet_type: RtcpType,
pub length: u16,
}
impl RtcpHeader {
pub fn parse(data: &[u8]) -> Result<(Self, &[u8]), RtcpParseError> {
if data.len() < 4 {
return Err(RtcpParseError::TooShort(data.len()));
}
let first_byte = data[0];
let version = (first_byte >> 6) & 0x03;
let padding = (first_byte >> 5) & 0x01 == 1;
let count = first_byte & 0x1F;
if version != 2 {
return Err(RtcpParseError::InvalidVersion(version));
}
let packet_type = RtcpType::try_from(data[1])?;
let length = u16::from_be_bytes([data[2], data[3]]);
Ok((
RtcpHeader {
version,
padding,
count,
packet_type,
length,
},
&data[4..],
))
}
pub fn build(&self, buf: &mut BytesMut) {
let first_byte = (self.version << 6) | ((self.padding as u8) << 5) | (self.count & 0x1F);
buf.put_u8(first_byte);
buf.put_u8(self.packet_type as u8);
buf.put_u16(self.length);
}
}
#[derive(Debug, Clone, Copy, Default)]
pub struct NtpTimestamp {
pub seconds: u32,
pub fraction: u32,
}
impl NtpTimestamp {
const NTP_UNIX_OFFSET: u64 = 2_208_988_800;
pub fn now() -> Self {
let duration = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default();
let seconds = duration.as_secs() + Self::NTP_UNIX_OFFSET;
let fraction = ((duration.subsec_nanos() as u64) << 32) / 1_000_000_000;
Self {
seconds: seconds as u32,
fraction: fraction as u32,
}
}
pub fn compact(&self) -> u32 {
((self.seconds & 0xFFFF) << 16) | ((self.fraction >> 16) & 0xFFFF)
}
pub fn from_compact(compact: u32) -> Self {
Self {
seconds: (compact >> 16) & 0xFFFF,
fraction: (compact & 0xFFFF) << 16,
}
}
}
#[derive(Debug, Clone, Default)]
pub struct ReportBlock {
pub ssrc: u32,
pub fraction_lost: u8,
pub cumulative_lost: i32,
pub extended_seq: u32,
pub jitter: u32,
pub last_sr: u32,
pub delay_since_sr: u32,
}
impl ReportBlock {
pub const SIZE: usize = 24;
pub fn parse(data: &[u8]) -> Result<Self, RtcpParseError> {
if data.len() < Self::SIZE {
return Err(RtcpParseError::TooShort(data.len()));
}
let mut buf = data;
let ssrc = buf.get_u32();
let fraction_lost = buf.get_u8();
let lost_bytes = [0, buf.get_u8(), buf.get_u8(), buf.get_u8()];
let cumulative_lost = i32::from_be_bytes(lost_bytes) >> 8; let extended_seq = buf.get_u32();
let jitter = buf.get_u32();
let last_sr = buf.get_u32();
let delay_since_sr = buf.get_u32();
Ok(ReportBlock {
ssrc,
fraction_lost,
cumulative_lost,
extended_seq,
jitter,
last_sr,
delay_since_sr,
})
}
pub fn build(&self, buf: &mut BytesMut) {
buf.put_u32(self.ssrc);
buf.put_u8(self.fraction_lost);
let lost_bytes = self.cumulative_lost.to_be_bytes();
buf.put_u8(lost_bytes[1]);
buf.put_u8(lost_bytes[2]);
buf.put_u8(lost_bytes[3]);
buf.put_u32(self.extended_seq);
buf.put_u32(self.jitter);
buf.put_u32(self.last_sr);
buf.put_u32(self.delay_since_sr);
}
}
#[derive(Debug, Clone)]
pub struct SenderReport {
pub ssrc: u32,
pub ntp_timestamp: NtpTimestamp,
pub rtp_timestamp: u32,
pub sender_packet_count: u32,
pub sender_octet_count: u32,
pub report_blocks: Vec<ReportBlock>,
}
impl SenderReport {
pub const MIN_SIZE: usize = 24;
pub fn parse(data: &[u8]) -> Result<Self, RtcpParseError> {
let (header, rest) = RtcpHeader::parse(data)?;
if header.packet_type != RtcpType::SenderReport {
return Err(RtcpParseError::UnknownPacketType(header.packet_type as u8));
}
if rest.len() < 24 {
return Err(RtcpParseError::TooShort(rest.len()));
}
let mut buf = rest;
let ssrc = buf.get_u32();
let ntp_seconds = buf.get_u32();
let ntp_fraction = buf.get_u32();
let rtp_timestamp = buf.get_u32();
let sender_packet_count = buf.get_u32();
let sender_octet_count = buf.get_u32();
let mut report_blocks = Vec::with_capacity(header.count as usize);
for _ in 0..header.count {
if buf.remaining() < ReportBlock::SIZE {
return Err(RtcpParseError::InvalidReportCount);
}
let block_data = &buf[..ReportBlock::SIZE];
let block = ReportBlock::parse(block_data).expect("report block size checked");
report_blocks.push(block);
buf.advance(ReportBlock::SIZE);
}
Ok(SenderReport {
ssrc,
ntp_timestamp: NtpTimestamp {
seconds: ntp_seconds,
fraction: ntp_fraction,
},
rtp_timestamp,
sender_packet_count,
sender_octet_count,
report_blocks,
})
}
pub fn build(&self) -> Bytes {
let report_count = self.report_blocks.len().min(31) as u8;
let length = 6 + report_count as u16 * 6;
let mut buf = BytesMut::with_capacity(28 + self.report_blocks.len() * ReportBlock::SIZE);
let header = RtcpHeader {
version: 2,
padding: false,
count: report_count,
packet_type: RtcpType::SenderReport,
length,
};
header.build(&mut buf);
buf.put_u32(self.ssrc);
buf.put_u32(self.ntp_timestamp.seconds);
buf.put_u32(self.ntp_timestamp.fraction);
buf.put_u32(self.rtp_timestamp);
buf.put_u32(self.sender_packet_count);
buf.put_u32(self.sender_octet_count);
for block in &self.report_blocks {
block.build(&mut buf);
}
buf.freeze()
}
}
#[derive(Debug, Clone)]
pub struct ReceiverReport {
pub ssrc: u32,
pub report_blocks: Vec<ReportBlock>,
}
impl ReceiverReport {
pub const MIN_SIZE: usize = 4;
pub fn parse(data: &[u8]) -> Result<Self, RtcpParseError> {
let (header, rest) = RtcpHeader::parse(data)?;
if header.packet_type != RtcpType::ReceiverReport {
return Err(RtcpParseError::UnknownPacketType(header.packet_type as u8));
}
if rest.len() < 4 {
return Err(RtcpParseError::TooShort(rest.len()));
}
let mut buf = rest;
let ssrc = buf.get_u32();
let mut report_blocks = Vec::with_capacity(header.count as usize);
for _ in 0..header.count {
if buf.remaining() < ReportBlock::SIZE {
return Err(RtcpParseError::InvalidReportCount);
}
let block_data = &buf[..ReportBlock::SIZE];
let block = ReportBlock::parse(block_data).expect("report block size checked");
report_blocks.push(block);
buf.advance(ReportBlock::SIZE);
}
Ok(ReceiverReport {
ssrc,
report_blocks,
})
}
pub fn build(&self) -> Bytes {
let report_count = self.report_blocks.len().min(31) as u8;
let length = 1 + report_count as u16 * 6;
let mut buf = BytesMut::with_capacity(8 + self.report_blocks.len() * ReportBlock::SIZE);
let header = RtcpHeader {
version: 2,
padding: false,
count: report_count,
packet_type: RtcpType::ReceiverReport,
length,
};
header.build(&mut buf);
buf.put_u32(self.ssrc);
for block in &self.report_blocks {
block.build(&mut buf);
}
buf.freeze()
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[repr(u8)]
pub enum SdesType {
End = 0,
CName = 1,
Name = 2,
Email = 3,
Phone = 4,
Location = 5,
Tool = 6,
Note = 7,
Private = 8,
}
#[derive(Debug, Clone)]
pub struct SdesItem {
pub item_type: SdesType,
pub value: String,
}
#[derive(Debug, Clone)]
pub struct SdesChunk {
pub ssrc: u32,
pub items: Vec<SdesItem>,
}
#[derive(Debug, Clone)]
pub struct SourceDescription {
pub chunks: Vec<SdesChunk>,
}
impl SourceDescription {
pub fn with_cname(ssrc: u32, cname: &str) -> Self {
Self {
chunks: vec![SdesChunk {
ssrc,
items: vec![SdesItem {
item_type: SdesType::CName,
value: cname.to_string(),
}],
}],
}
}
pub fn parse(data: &[u8]) -> Result<Self, RtcpParseError> {
let (header, rest) = RtcpHeader::parse(data)?;
if header.packet_type != RtcpType::SourceDescription {
return Err(RtcpParseError::UnknownPacketType(header.packet_type as u8));
}
let declared = (header.length as usize) * 4;
if rest.len() < declared {
return Err(RtcpParseError::TooShort(rest.len()));
}
let mut buf = &rest[..declared];
let mut chunks = Vec::with_capacity(header.count as usize);
for _ in 0..header.count {
if buf.remaining() < 4 {
return Err(RtcpParseError::TooShort(buf.remaining()));
}
let chunk_start = buf.len();
let ssrc = buf.get_u32();
let mut items = Vec::new();
loop {
if buf.remaining() < 1 {
return Err(RtcpParseError::TooShort(buf.remaining()));
}
let item_type = buf.get_u8();
if item_type == 0 {
break;
}
if buf.remaining() < 1 {
return Err(RtcpParseError::TooShort(buf.remaining()));
}
let len = buf.get_u8() as usize;
if buf.remaining() < len {
return Err(RtcpParseError::TooShort(buf.remaining()));
}
let value_bytes = &buf[..len];
let value = String::from_utf8_lossy(value_bytes).into_owned();
buf.advance(len);
let parsed_type = match item_type {
1 => Some(SdesType::CName),
2 => Some(SdesType::Name),
3 => Some(SdesType::Email),
4 => Some(SdesType::Phone),
5 => Some(SdesType::Location),
6 => Some(SdesType::Tool),
7 => Some(SdesType::Note),
8 => Some(SdesType::Private),
_ => None,
};
if let Some(item_type) = parsed_type {
items.push(SdesItem { item_type, value });
}
}
let consumed = chunk_start - buf.len();
let pad = (4 - (consumed % 4)) % 4;
if buf.remaining() < pad {
return Err(RtcpParseError::TooShort(buf.remaining()));
}
buf.advance(pad);
chunks.push(SdesChunk { ssrc, items });
}
Ok(SourceDescription { chunks })
}
pub fn build(&self) -> Bytes {
let mut buf = BytesMut::with_capacity(256);
let header_pos = buf.len();
buf.put_u32(0);
let chunk_count = self.chunks.len().min(31) as u8;
for chunk in &self.chunks {
buf.put_u32(chunk.ssrc);
for item in &chunk.items {
buf.put_u8(item.item_type as u8);
let value_bytes = item.value.as_bytes();
buf.put_u8(value_bytes.len() as u8);
buf.put_slice(value_bytes);
}
buf.put_u8(0); while !buf.len().is_multiple_of(4) {
buf.put_u8(0);
}
}
let length = ((buf.len() - 4) / 4) as u16;
let header_byte = (2 << 6) | chunk_count;
buf[header_pos] = header_byte;
buf[header_pos + 1] = RtcpType::SourceDescription as u8;
buf[header_pos + 2] = (length >> 8) as u8;
buf[header_pos + 3] = (length & 0xFF) as u8;
buf.freeze()
}
}
#[derive(Debug, Clone)]
pub struct Goodbye {
pub ssrcs: Vec<u32>,
pub reason: Option<String>,
}
impl Goodbye {
pub fn new(ssrc: u32) -> Self {
Self {
ssrcs: vec![ssrc],
reason: None,
}
}
pub fn parse(data: &[u8]) -> Result<Self, RtcpParseError> {
let (header, rest) = RtcpHeader::parse(data)?;
if header.packet_type != RtcpType::Goodbye {
return Err(RtcpParseError::UnknownPacketType(header.packet_type as u8));
}
let declared = (header.length as usize) * 4;
if rest.len() < declared {
return Err(RtcpParseError::TooShort(rest.len()));
}
let mut buf = &rest[..declared];
let ssrc_count = header.count as usize;
if buf.remaining() < ssrc_count * 4 {
return Err(RtcpParseError::TooShort(buf.remaining()));
}
let mut ssrcs = Vec::with_capacity(ssrc_count);
for _ in 0..ssrc_count {
ssrcs.push(buf.get_u32());
}
let reason = if buf.remaining() >= 1 {
let len = buf.get_u8() as usize;
if buf.remaining() < len {
return Err(RtcpParseError::TooShort(buf.remaining()));
}
let value = String::from_utf8_lossy(&buf[..len]).into_owned();
buf.advance(len);
Some(value)
} else {
None
};
Ok(Goodbye { ssrcs, reason })
}
pub fn build(&self) -> Bytes {
let mut buf = BytesMut::with_capacity(32);
let ssrc_count = self.ssrcs.len().min(31) as u8;
let mut content_len = self.ssrcs.len() * 4;
if let Some(ref reason) = self.reason {
content_len += 1 + reason.len();
content_len = (content_len + 3) & !3;
}
let length = (content_len / 4) as u16;
let header = RtcpHeader {
version: 2,
padding: false,
count: ssrc_count,
packet_type: RtcpType::Goodbye,
length,
};
header.build(&mut buf);
for &ssrc in &self.ssrcs {
buf.put_u32(ssrc);
}
if let Some(ref reason) = self.reason {
let reason_bytes = reason.as_bytes();
buf.put_u8(reason_bytes.len() as u8);
buf.put_slice(reason_bytes);
while !buf.len().is_multiple_of(4) {
buf.put_u8(0);
}
}
buf.freeze()
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[repr(u8)]
pub enum TransportFeedbackType {
Nack = 1,
TransportCC = 15,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[repr(u8)]
pub enum PayloadFeedbackType {
Pli = 1,
Sli = 2,
Rpsi = 3,
Fir = 4,
Tstr = 5,
Tstn = 6,
Vbcm = 7,
Afb = 15,
}
#[derive(Debug, Clone)]
pub struct Nack {
pub sender_ssrc: u32,
pub media_ssrc: u32,
pub nacks: Vec<NackEntry>,
}
#[derive(Debug, Clone, Copy)]
pub struct NackEntry {
pub pid: u16,
pub blp: u16,
}
impl NackEntry {
pub fn single(seq: u16) -> Self {
Self { pid: seq, blp: 0 }
}
pub fn from_sequences(seqs: &[u16]) -> Option<Self> {
let pid = *seqs.first()?;
let mut blp = 0u16;
for &seq in seqs.iter().skip(1) {
let diff = seq.wrapping_sub(pid).wrapping_sub(1);
if diff < 16 {
blp |= 1 << diff;
}
}
Some(Self { pid, blp })
}
pub fn lost_sequences(&self) -> Vec<u16> {
let mut seqs = vec![self.pid];
for i in 0..16 {
if (self.blp >> i) & 1 == 1 {
seqs.push(self.pid.wrapping_add(i + 1));
}
}
seqs
}
}
impl Nack {
pub fn new(sender_ssrc: u32, media_ssrc: u32, lost_seq: u16) -> Self {
Self {
sender_ssrc,
media_ssrc,
nacks: vec![NackEntry::single(lost_seq)],
}
}
pub fn from_lost_packets(sender_ssrc: u32, media_ssrc: u32, lost_seqs: &[u16]) -> Self {
let mut nacks = Vec::new();
let mut remaining: Vec<u16> = lost_seqs.to_vec();
remaining.sort();
while !remaining.is_empty() {
let pid = remaining[0];
let mut group = vec![pid];
let mut new_remaining = Vec::new();
for &seq in remaining.iter().skip(1) {
let diff = seq.wrapping_sub(pid);
if diff > 0 && diff <= 16 {
group.push(seq);
} else {
new_remaining.push(seq);
}
}
let entry = NackEntry::from_sequences(&group).expect("nack group empty");
nacks.push(entry);
remaining = new_remaining;
}
Self {
sender_ssrc,
media_ssrc,
nacks,
}
}
pub fn parse(data: &[u8]) -> Result<Self, RtcpParseError> {
let (header, rest) = RtcpHeader::parse(data)?;
if header.packet_type != RtcpType::TransportFeedback || header.count != 1 {
return Err(RtcpParseError::UnknownPacketType(header.packet_type as u8));
}
if rest.len() < 8 {
return Err(RtcpParseError::TooShort(rest.len()));
}
let mut buf = rest;
let sender_ssrc = buf.get_u32();
let media_ssrc = buf.get_u32();
let mut nacks = Vec::new();
while buf.remaining() >= 4 {
let pid = buf.get_u16();
let blp = buf.get_u16();
nacks.push(NackEntry { pid, blp });
}
Ok(Nack {
sender_ssrc,
media_ssrc,
nacks,
})
}
pub fn build(&self) -> Bytes {
let nack_count = self.nacks.len();
let length = (2 + nack_count) as u16;
let mut buf = BytesMut::with_capacity(12 + nack_count * 4);
let header = RtcpHeader {
version: 2,
padding: false,
count: TransportFeedbackType::Nack as u8,
packet_type: RtcpType::TransportFeedback,
length,
};
header.build(&mut buf);
buf.put_u32(self.sender_ssrc);
buf.put_u32(self.media_ssrc);
for nack in &self.nacks {
buf.put_u16(nack.pid);
buf.put_u16(nack.blp);
}
buf.freeze()
}
pub fn all_lost_sequences(&self) -> Vec<u16> {
self.nacks.iter().flat_map(|n| n.lost_sequences()).collect()
}
}
#[derive(Debug, Clone)]
pub struct Pli {
pub sender_ssrc: u32,
pub media_ssrc: u32,
}
impl Pli {
pub fn new(sender_ssrc: u32, media_ssrc: u32) -> Self {
Self {
sender_ssrc,
media_ssrc,
}
}
pub fn parse(data: &[u8]) -> Result<Self, RtcpParseError> {
let (header, rest) = RtcpHeader::parse(data)?;
if header.packet_type != RtcpType::PayloadFeedback || header.count != 1 {
return Err(RtcpParseError::UnknownPacketType(header.packet_type as u8));
}
if rest.len() < 8 {
return Err(RtcpParseError::TooShort(rest.len()));
}
let mut buf = rest;
let sender_ssrc = buf.get_u32();
let media_ssrc = buf.get_u32();
Ok(Pli {
sender_ssrc,
media_ssrc,
})
}
pub fn build(&self) -> Bytes {
let mut buf = BytesMut::with_capacity(12);
let header = RtcpHeader {
version: 2,
padding: false,
count: PayloadFeedbackType::Pli as u8,
packet_type: RtcpType::PayloadFeedback,
length: 2, };
header.build(&mut buf);
buf.put_u32(self.sender_ssrc);
buf.put_u32(self.media_ssrc);
buf.freeze()
}
}
#[derive(Debug, Clone)]
pub struct Fir {
pub sender_ssrc: u32,
pub media_ssrc: u32,
pub entries: Vec<FirEntry>,
}
#[derive(Debug, Clone, Copy)]
pub struct FirEntry {
pub ssrc: u32,
pub seq_nr: u8,
}
impl Fir {
pub fn new(sender_ssrc: u32, target_ssrc: u32, seq_nr: u8) -> Self {
Self {
sender_ssrc,
media_ssrc: 0,
entries: vec![FirEntry {
ssrc: target_ssrc,
seq_nr,
}],
}
}
pub fn parse(data: &[u8]) -> Result<Self, RtcpParseError> {
let (header, rest) = RtcpHeader::parse(data)?;
if header.packet_type != RtcpType::PayloadFeedback || header.count != 4 {
return Err(RtcpParseError::UnknownPacketType(header.packet_type as u8));
}
if rest.len() < 8 {
return Err(RtcpParseError::TooShort(rest.len()));
}
let mut buf = rest;
let sender_ssrc = buf.get_u32();
let media_ssrc = buf.get_u32();
let mut entries = Vec::new();
while buf.remaining() >= 8 {
let ssrc = buf.get_u32();
let seq_nr = buf.get_u8();
let _ = buf.get_u8(); let _ = buf.get_u16(); entries.push(FirEntry { ssrc, seq_nr });
}
Ok(Fir {
sender_ssrc,
media_ssrc,
entries,
})
}
pub fn build(&self) -> Bytes {
let entry_count = self.entries.len();
let length = (2 + entry_count * 2) as u16;
let mut buf = BytesMut::with_capacity(12 + entry_count * 8);
let header = RtcpHeader {
version: 2,
padding: false,
count: PayloadFeedbackType::Fir as u8,
packet_type: RtcpType::PayloadFeedback,
length,
};
header.build(&mut buf);
buf.put_u32(self.sender_ssrc);
buf.put_u32(self.media_ssrc);
for entry in &self.entries {
buf.put_u32(entry.ssrc);
buf.put_u8(entry.seq_nr);
buf.put_u8(0); buf.put_u16(0); }
buf.freeze()
}
}
#[derive(Debug, Clone)]
pub struct Remb {
pub sender_ssrc: u32,
pub media_ssrc: u32,
pub bitrate: u64,
pub ssrcs: Vec<u32>,
}
impl Remb {
const UNIQUE_ID: [u8; 4] = [b'R', b'E', b'M', b'B'];
pub fn new(sender_ssrc: u32, bitrate: u64, ssrcs: Vec<u32>) -> Self {
Self {
sender_ssrc,
media_ssrc: 0,
bitrate,
ssrcs,
}
}
pub fn parse(data: &[u8]) -> Result<Self, RtcpParseError> {
let (header, rest) = RtcpHeader::parse(data)?;
if header.packet_type != RtcpType::PayloadFeedback || header.count != 15 {
return Err(RtcpParseError::UnknownPacketType(header.packet_type as u8));
}
if rest.len() < 16 {
return Err(RtcpParseError::TooShort(rest.len()));
}
let mut buf = rest;
let sender_ssrc = buf.get_u32();
let media_ssrc = buf.get_u32();
let mut unique_id = [0u8; 4];
buf.copy_to_slice(&mut unique_id);
if unique_id != Self::UNIQUE_ID {
return Err(RtcpParseError::UnknownPacketType(206));
}
let num_ssrc = buf.get_u8();
let br_exp = (buf.get_u8() & 0xFC) >> 2;
let mantissa_high = (rest[13] & 0x03) as u32;
let mantissa_mid = rest[14] as u32;
let mantissa_low = rest[15] as u32;
let mantissa = (mantissa_high << 16) | (mantissa_mid << 8) | mantissa_low;
buf.advance(2);
let bitrate = (mantissa as u64) << br_exp;
let mut ssrcs = Vec::with_capacity(num_ssrc as usize);
for _ in 0..num_ssrc {
if buf.remaining() < 4 {
break;
}
ssrcs.push(buf.get_u32());
}
Ok(Remb {
sender_ssrc,
media_ssrc,
bitrate,
ssrcs,
})
}
pub fn build(&self) -> Bytes {
let ssrc_count = self.ssrcs.len().min(255);
let length = (4 + ssrc_count) as u16;
let mut buf = BytesMut::with_capacity(16 + ssrc_count * 4);
let header = RtcpHeader {
version: 2,
padding: false,
count: PayloadFeedbackType::Afb as u8,
packet_type: RtcpType::PayloadFeedback,
length,
};
header.build(&mut buf);
buf.put_u32(self.sender_ssrc);
buf.put_u32(self.media_ssrc);
buf.put_slice(&Self::UNIQUE_ID);
let (mantissa, exp) = Self::encode_bitrate(self.bitrate);
buf.put_u8(ssrc_count as u8);
buf.put_u8((exp << 2) | ((mantissa >> 16) as u8 & 0x03));
buf.put_u16((mantissa & 0xFFFF) as u16);
for &ssrc in self.ssrcs.iter().take(ssrc_count) {
buf.put_u32(ssrc);
}
buf.freeze()
}
fn encode_bitrate(bitrate: u64) -> (u32, u8) {
if bitrate == 0 {
return (0, 0);
}
let bits = 64 - bitrate.leading_zeros();
if bits <= 18 {
(bitrate as u32, 0)
} else {
let exp = (bits - 18) as u8;
let mantissa = (bitrate >> exp) as u32;
(mantissa & 0x3FFFF, exp)
}
}
}
#[derive(Debug, Clone)]
pub struct RtcpCompound {
pub packets: Vec<RtcpPacket>,
}
#[derive(Debug, Clone)]
pub enum RtcpPacket {
SenderReport(SenderReport),
ReceiverReport(ReceiverReport),
SourceDescription(SourceDescription),
Goodbye(Goodbye),
Nack(Nack),
Pli(Pli),
Fir(Fir),
Remb(Remb),
Unknown {
packet_type: u8,
bytes: Bytes,
},
}
impl RtcpCompound {
pub fn sender_compound(sr: SenderReport, cname: &str) -> Self {
let sdes = SourceDescription::with_cname(sr.ssrc, cname);
Self {
packets: vec![
RtcpPacket::SenderReport(sr),
RtcpPacket::SourceDescription(sdes),
],
}
}
pub fn receiver_compound(rr: ReceiverReport, cname: &str) -> Self {
let sdes = SourceDescription::with_cname(rr.ssrc, cname);
Self {
packets: vec![
RtcpPacket::ReceiverReport(rr),
RtcpPacket::SourceDescription(sdes),
],
}
}
pub fn build(&self) -> Bytes {
let mut buf = BytesMut::with_capacity(512);
for packet in &self.packets {
match packet {
RtcpPacket::SenderReport(sr) => buf.extend_from_slice(&sr.build()),
RtcpPacket::ReceiverReport(rr) => buf.extend_from_slice(&rr.build()),
RtcpPacket::SourceDescription(sdes) => buf.extend_from_slice(&sdes.build()),
RtcpPacket::Goodbye(bye) => buf.extend_from_slice(&bye.build()),
RtcpPacket::Nack(nack) => buf.extend_from_slice(&nack.build()),
RtcpPacket::Pli(pli) => buf.extend_from_slice(&pli.build()),
RtcpPacket::Fir(fir) => buf.extend_from_slice(&fir.build()),
RtcpPacket::Remb(remb) => buf.extend_from_slice(&remb.build()),
RtcpPacket::Unknown { bytes, .. } => buf.extend_from_slice(bytes),
}
}
buf.freeze()
}
pub fn add_nack(&mut self, nack: Nack) {
self.packets.push(RtcpPacket::Nack(nack));
}
pub fn add_pli(&mut self, pli: Pli) {
self.packets.push(RtcpPacket::Pli(pli));
}
pub fn add_fir(&mut self, fir: Fir) {
self.packets.push(RtcpPacket::Fir(fir));
}
pub fn add_remb(&mut self, remb: Remb) {
self.packets.push(RtcpPacket::Remb(remb));
}
pub fn parse(data: &[u8]) -> Result<Self, String> {
let mut remaining = data;
let mut packets = Vec::new();
while !remaining.is_empty() {
if remaining.len() < 4 {
return Err(format!(
"truncated RTCP packet: {} bytes left, need at least 4",
remaining.len()
));
}
let version = (remaining[0] >> 6) & 0x03;
if version != 2 {
return Err(format!("invalid RTCP version: {version}"));
}
let pt = remaining[1];
let length_words = u16::from_be_bytes([remaining[2], remaining[3]]) as usize;
let pkt_len = (length_words + 1) * 4;
if pkt_len > remaining.len() {
return Err(format!(
"truncated RTCP packet: declared {pkt_len} bytes, only {} remain",
remaining.len()
));
}
let pkt_bytes = &remaining[..pkt_len];
let subtype = remaining[0] & 0x1F;
let packet = match pt {
200 => RtcpPacket::SenderReport(
SenderReport::parse(pkt_bytes).map_err(|e| format!("{e}"))?,
),
201 => RtcpPacket::ReceiverReport(
ReceiverReport::parse(pkt_bytes).map_err(|e| format!("{e}"))?,
),
202 => RtcpPacket::SourceDescription(
SourceDescription::parse(pkt_bytes).map_err(|e| format!("{e}"))?,
),
203 => RtcpPacket::Goodbye(Goodbye::parse(pkt_bytes).map_err(|e| format!("{e}"))?),
206 if subtype == PayloadFeedbackType::Afb as u8 && pkt_len >= 16 && {
&pkt_bytes[12..16] == b"REMB"
} =>
{
RtcpPacket::Remb(Remb::parse(pkt_bytes).map_err(|e| format!("{e}"))?)
}
_ => RtcpPacket::Unknown {
packet_type: pt,
bytes: Bytes::copy_from_slice(pkt_bytes),
},
};
packets.push(packet);
remaining = &remaining[pkt_len..];
}
Ok(RtcpCompound { packets })
}
}
#[cfg(test)]
mod tests {
use super::*;
fn assert_rtcp_err_contains<T>(result: Result<T, RtcpParseError>, needle: &str) {
assert!(result.is_err());
let err = result.err().unwrap();
assert!(format!("{err:?}").contains(needle));
}
#[test]
fn test_ntp_timestamp() {
let ntp = NtpTimestamp::now();
assert!(ntp.seconds > 0);
let compact = ntp.compact();
assert!(compact > 0);
}
#[test]
fn test_sender_report_build_parse() {
let sr = SenderReport {
ssrc: 12345,
ntp_timestamp: NtpTimestamp::now(),
rtp_timestamp: 160000,
sender_packet_count: 100,
sender_octet_count: 16000,
report_blocks: vec![],
};
let bytes = sr.build();
let parsed = SenderReport::parse(&bytes).unwrap();
assert_eq!(parsed.ssrc, 12345);
assert_eq!(parsed.rtp_timestamp, 160000);
assert_eq!(parsed.sender_packet_count, 100);
assert_eq!(parsed.sender_octet_count, 16000);
}
#[test]
fn test_sender_report_with_report_block() {
let sr = SenderReport {
ssrc: 12345,
ntp_timestamp: NtpTimestamp::now(),
rtp_timestamp: 160000,
sender_packet_count: 100,
sender_octet_count: 16000,
report_blocks: vec![ReportBlock {
ssrc: 67890,
fraction_lost: 25,
cumulative_lost: 10,
extended_seq: 50000,
jitter: 160,
last_sr: 0,
delay_since_sr: 0,
}],
};
let bytes = sr.build();
let parsed = SenderReport::parse(&bytes).unwrap();
assert_eq!(parsed.report_blocks.len(), 1);
assert_eq!(parsed.report_blocks[0].ssrc, 67890);
assert_eq!(parsed.report_blocks[0].fraction_lost, 25);
}
#[test]
fn test_sr_parse_truncated_fuzz_001() {
let crash: [u8; 27] = [
0x80, 0xc8, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00,
0x00, 0x00, 0x00, 0x00, 0x00, 0x12, 0x01, 0x00, 0x00, 0x0a, 0x05, 0x00, 0x80,
];
let result = SenderReport::parse(&crash);
assert!(result.is_err());
assert_rtcp_err_contains(result, "TooShort");
}
#[test]
fn test_receiver_report_build_parse() {
let rr = ReceiverReport {
ssrc: 12345,
report_blocks: vec![ReportBlock {
ssrc: 67890,
fraction_lost: 0,
cumulative_lost: 0,
extended_seq: 1000,
jitter: 80,
last_sr: 0,
delay_since_sr: 0,
}],
};
let bytes = rr.build();
let parsed = ReceiverReport::parse(&bytes).unwrap();
assert_eq!(parsed.ssrc, 12345);
assert_eq!(parsed.report_blocks.len(), 1);
assert_eq!(parsed.report_blocks[0].ssrc, 67890);
}
#[test]
fn test_sdes_build() {
let sdes = SourceDescription::with_cname(12345, "user@example.com");
let bytes = sdes.build();
assert_eq!(bytes[1], RtcpType::SourceDescription as u8);
}
#[test]
fn test_goodbye_build() {
let bye = Goodbye::new(12345);
let bytes = bye.build();
assert_eq!(bytes[1], RtcpType::Goodbye as u8);
}
#[test]
fn test_compound_packet() {
let sr = SenderReport {
ssrc: 12345,
ntp_timestamp: NtpTimestamp::now(),
rtp_timestamp: 160000,
sender_packet_count: 100,
sender_octet_count: 16000,
report_blocks: vec![],
};
let compound = RtcpCompound::sender_compound(sr, "user@example.com");
let bytes = compound.build();
assert!(bytes.len() > 28);
assert_eq!(bytes[1], RtcpType::SenderReport as u8);
}
#[test]
fn test_report_block() {
let block = ReportBlock {
ssrc: 12345,
fraction_lost: 128, cumulative_lost: 1000,
extended_seq: 65536 + 1000,
jitter: 320,
last_sr: 0x12345678,
delay_since_sr: 65536, };
let mut buf = BytesMut::new();
block.build(&mut buf);
let parsed = ReportBlock::parse(&buf).unwrap();
assert_eq!(parsed.ssrc, 12345);
assert_eq!(parsed.fraction_lost, 128);
assert_eq!(parsed.extended_seq, 65536 + 1000);
assert_eq!(parsed.jitter, 320);
}
#[test]
fn test_nack_single_packet() {
let nack = Nack::new(111111, 222222, 1000);
let bytes = nack.build();
assert_eq!(bytes[1], RtcpType::TransportFeedback as u8);
assert_eq!(bytes[0] & 0x1F, TransportFeedbackType::Nack as u8);
let parsed = Nack::parse(&bytes).unwrap();
assert_eq!(parsed.sender_ssrc, 111111);
assert_eq!(parsed.media_ssrc, 222222);
assert_eq!(parsed.nacks.len(), 1);
assert_eq!(parsed.nacks[0].pid, 1000);
assert_eq!(parsed.nacks[0].blp, 0);
}
#[test]
fn test_nack_multiple_packets() {
let lost = vec![100, 101, 103, 105, 200];
let nack = Nack::from_lost_packets(111111, 222222, &lost);
assert_eq!(nack.nacks.len(), 2);
let all_lost = nack.all_lost_sequences();
assert!(all_lost.contains(&100));
assert!(all_lost.contains(&101));
assert!(all_lost.contains(&103));
assert!(all_lost.contains(&105));
assert!(all_lost.contains(&200));
}
#[test]
fn test_nack_entry_blp() {
let entry = NackEntry::from_sequences(&[1000, 1001, 1002, 1016]).unwrap();
assert_eq!(entry.pid, 1000);
assert_eq!(entry.blp & 1, 1); assert_eq!((entry.blp >> 1) & 1, 1); assert_eq!((entry.blp >> 15) & 1, 1);
let seqs = entry.lost_sequences();
assert_eq!(seqs.len(), 4);
assert!(seqs.contains(&1000));
assert!(seqs.contains(&1001));
assert!(seqs.contains(&1002));
assert!(seqs.contains(&1016));
}
#[test]
fn test_pli_build_parse() {
let pli = Pli::new(111111, 222222);
let bytes = pli.build();
assert_eq!(bytes[1], RtcpType::PayloadFeedback as u8);
assert_eq!(bytes[0] & 0x1F, PayloadFeedbackType::Pli as u8);
let parsed = Pli::parse(&bytes).unwrap();
assert_eq!(parsed.sender_ssrc, 111111);
assert_eq!(parsed.media_ssrc, 222222);
}
#[test]
fn test_fir_build_parse() {
let fir = Fir::new(111111, 222222, 5);
let bytes = fir.build();
assert_eq!(bytes[1], RtcpType::PayloadFeedback as u8);
assert_eq!(bytes[0] & 0x1F, PayloadFeedbackType::Fir as u8);
let parsed = Fir::parse(&bytes).unwrap();
assert_eq!(parsed.sender_ssrc, 111111);
assert_eq!(parsed.entries.len(), 1);
assert_eq!(parsed.entries[0].ssrc, 222222);
assert_eq!(parsed.entries[0].seq_nr, 5);
}
#[test]
fn test_remb_build() {
let remb = Remb::new(111111, 1_500_000, vec![222222, 333333]);
let bytes = remb.build();
assert_eq!(bytes[1], RtcpType::PayloadFeedback as u8);
assert_eq!(bytes[0] & 0x1F, PayloadFeedbackType::Afb as u8);
assert_eq!(&bytes[12..16], b"REMB");
assert_eq!(bytes[16], 2);
}
#[test]
fn test_remb_bitrate_encoding() {
for &bitrate in &[0u64, 1000, 100_000, 1_000_000, 10_000_000, 100_000_000] {
let (mantissa, exp) = Remb::encode_bitrate(bitrate);
let decoded = (mantissa as u64) << exp;
if bitrate > 0 {
let error = ((decoded as i64 - bitrate as i64).abs() as f64) / (bitrate as f64);
assert!(error < 0.01);
}
}
}
#[test]
fn test_compound_with_feedback() {
let sr = SenderReport {
ssrc: 12345,
ntp_timestamp: NtpTimestamp::now(),
rtp_timestamp: 160000,
sender_packet_count: 100,
sender_octet_count: 16000,
report_blocks: vec![],
};
let mut compound = RtcpCompound::sender_compound(sr, "user@example.com");
compound.add_nack(Nack::new(12345, 67890, 500));
compound.add_pli(Pli::new(12345, 67890));
let bytes = compound.build();
assert!(bytes.len() > 60);
}
#[test]
fn test_compound_with_goodbye() {
let sr = SenderReport {
ssrc: 12345,
ntp_timestamp: NtpTimestamp::now(),
rtp_timestamp: 160000,
sender_packet_count: 100,
sender_octet_count: 16000,
report_blocks: vec![],
};
let mut compound = RtcpCompound::sender_compound(sr, "user@example.com");
compound
.packets
.push(RtcpPacket::Goodbye(Goodbye::new(12345)));
let bytes = compound.build();
assert!(!bytes.is_empty());
}
#[test]
fn test_rtcp_header_parse_too_short() {
let data = [0u8; 2]; assert_rtcp_err_contains(RtcpHeader::parse(&data), "TooShort");
}
#[test]
fn test_rtcp_header_invalid_version() {
let data = [0x00, 200, 0, 6];
assert_rtcp_err_contains(RtcpHeader::parse(&data), "InvalidVersion");
let data = [0xC0, 200, 0, 6];
assert_rtcp_err_contains(RtcpHeader::parse(&data), "InvalidVersion");
}
#[test]
fn test_rtcp_header_invalid_packet_type() {
let data = [0x80, 199, 0, 1];
assert_rtcp_err_contains(RtcpHeader::parse(&data), "UnknownPacketType");
}
#[test]
fn test_rtcp_type_unknown() {
let result = RtcpType::try_from(199u8);
assert_rtcp_err_contains(result, "UnknownPacketType");
}
#[test]
fn test_rtcp_type_all_values() {
assert_eq!(RtcpType::try_from(200u8).unwrap(), RtcpType::SenderReport);
assert_eq!(RtcpType::try_from(201u8).unwrap(), RtcpType::ReceiverReport);
assert_eq!(
RtcpType::try_from(202u8).unwrap(),
RtcpType::SourceDescription
);
assert_eq!(RtcpType::try_from(203u8).unwrap(), RtcpType::Goodbye);
assert_eq!(
RtcpType::try_from(204u8).unwrap(),
RtcpType::ApplicationDefined
);
assert_eq!(
RtcpType::try_from(205u8).unwrap(),
RtcpType::TransportFeedback
);
assert_eq!(
RtcpType::try_from(206u8).unwrap(),
RtcpType::PayloadFeedback
);
}
#[test]
fn test_ntp_timestamp_compact_roundtrip() {
let ntp = NtpTimestamp {
seconds: 0x12345678,
fraction: 0xABCDE000,
};
let compact = ntp.compact();
let restored = NtpTimestamp::from_compact(compact);
assert_eq!(restored.seconds, ntp.seconds & 0xFFFF);
assert_eq!(restored.fraction & 0xFFFF0000, ntp.fraction & 0xFFFF0000);
}
#[test]
fn test_report_block_too_short() {
let data = [0u8; 10]; assert_rtcp_err_contains(ReportBlock::parse(&data), "TooShort");
}
#[test]
fn test_sender_report_parse_header_too_short() {
let data = [0u8; 2];
assert_rtcp_err_contains(SenderReport::parse(&data), "TooShort");
}
#[test]
fn test_sender_report_wrong_type() {
let rr = ReceiverReport {
ssrc: 12345,
report_blocks: vec![],
};
let bytes = rr.build();
let result = SenderReport::parse(&bytes);
assert!(result.is_err());
}
#[test]
fn test_sender_report_too_short_payload() {
let mut data = vec![0x80, 200, 0, 1]; data.extend_from_slice(&[0u8; 4]); assert_rtcp_err_contains(SenderReport::parse(&data), "TooShort");
}
#[test]
fn test_receiver_report_parse_header_too_short() {
let data = [0u8; 2];
assert_rtcp_err_contains(ReceiverReport::parse(&data), "TooShort");
}
#[test]
fn test_sender_report_invalid_report_count() {
let sr = SenderReport {
ssrc: 12345,
ntp_timestamp: NtpTimestamp::now(),
rtp_timestamp: 160000,
sender_packet_count: 100,
sender_octet_count: 16000,
report_blocks: vec![],
};
let mut bytes = sr.build().to_vec();
bytes[0] = (bytes[0] & 0xE0) | 5;
assert_rtcp_err_contains(SenderReport::parse(&bytes), "InvalidReportCount");
}
#[test]
fn test_receiver_report_wrong_type() {
let sr = SenderReport {
ssrc: 12345,
ntp_timestamp: NtpTimestamp::now(),
rtp_timestamp: 160000,
sender_packet_count: 100,
sender_octet_count: 16000,
report_blocks: vec![],
};
let bytes = sr.build();
let result = ReceiverReport::parse(&bytes);
assert!(result.is_err());
}
#[test]
fn test_receiver_report_too_short_payload() {
let data = vec![0x80, 201, 0, 0]; assert_rtcp_err_contains(ReceiverReport::parse(&data), "TooShort");
}
#[test]
fn test_receiver_report_invalid_report_count() {
let rr = ReceiverReport {
ssrc: 12345,
report_blocks: vec![],
};
let mut bytes = rr.build().to_vec();
bytes[0] = (bytes[0] & 0xE0) | 5;
assert_rtcp_err_contains(ReceiverReport::parse(&bytes), "InvalidReportCount");
}
#[test]
fn test_receiver_report_compound() {
let rr = ReceiverReport {
ssrc: 12345,
report_blocks: vec![ReportBlock {
ssrc: 67890,
fraction_lost: 10,
cumulative_lost: 5,
extended_seq: 1000,
jitter: 50,
last_sr: 0,
delay_since_sr: 0,
}],
};
let compound = RtcpCompound::receiver_compound(rr, "receiver@test.com");
let bytes = compound.build();
assert!(bytes.len() > 8);
assert_eq!(bytes[1], RtcpType::ReceiverReport as u8);
}
#[test]
fn test_goodbye_with_reason() {
let bye = Goodbye {
ssrcs: vec![12345, 67890],
reason: Some("Going offline".to_string()),
};
let bytes = bye.build();
assert_eq!(bytes[1], RtcpType::Goodbye as u8);
assert_eq!(bytes[0] & 0x1F, 2);
}
#[test]
fn test_nack_parse_too_short() {
let data = vec![0x81, 205, 0, 1, 0, 0, 0, 0]; assert_rtcp_err_contains(Nack::parse(&data), "TooShort");
}
#[test]
fn test_nack_parse_header_too_short() {
let data = [0u8; 2];
assert_rtcp_err_contains(Nack::parse(&data), "TooShort");
}
#[test]
fn test_nack_wrong_type() {
let pli = Pli::new(111111, 222222);
let bytes = pli.build();
let result = Nack::parse(&bytes);
assert!(result.is_err());
}
#[test]
fn test_nack_wrong_count() {
let nack = Nack::new(111111, 222222, 1000);
let mut bytes = nack.build().to_vec();
bytes[0] = (bytes[0] & 0xE0) | 2;
let result = Nack::parse(&bytes);
assert!(result.is_err());
}
#[test]
fn test_pli_parse_too_short() {
let data = vec![0x81, 206, 0, 1, 0, 0, 0, 0]; assert_rtcp_err_contains(Pli::parse(&data), "TooShort");
}
#[test]
fn test_pli_parse_header_too_short() {
let data = [0u8; 2];
assert_rtcp_err_contains(Pli::parse(&data), "TooShort");
}
#[test]
fn test_pli_wrong_type() {
let nack = Nack::new(111111, 222222, 1000);
let bytes = nack.build();
let result = Pli::parse(&bytes);
assert!(result.is_err());
}
#[test]
fn test_pli_wrong_count() {
let pli = Pli::new(111111, 222222);
let mut bytes = pli.build().to_vec();
bytes[0] = (bytes[0] & 0xE0) | 2;
let result = Pli::parse(&bytes);
assert!(result.is_err());
}
#[test]
fn test_fir_parse_too_short() {
let data = vec![0x84, 206, 0, 1, 0, 0, 0, 0]; assert_rtcp_err_contains(Fir::parse(&data), "TooShort");
}
#[test]
fn test_fir_parse_header_too_short() {
let data = [0u8; 2];
assert_rtcp_err_contains(Fir::parse(&data), "TooShort");
}
#[test]
fn test_fir_wrong_type() {
let pli = Pli::new(111111, 222222);
let bytes = pli.build();
let result = Fir::parse(&bytes);
assert!(result.is_err());
}
#[test]
fn test_fir_wrong_packet_type() {
let header = RtcpHeader {
version: 2,
padding: false,
count: 0,
packet_type: RtcpType::SenderReport,
length: 0,
};
let mut buf = BytesMut::new();
header.build(&mut buf);
assert_rtcp_err_contains(Fir::parse(&buf), "UnknownPacketType");
}
#[test]
fn test_fir_wrong_count() {
let fir = Fir::new(111111, 222222, 1);
let mut bytes = fir.build().to_vec();
bytes[0] = (bytes[0] & 0xE0) | 5;
let result = Fir::parse(&bytes);
assert!(result.is_err());
}
#[test]
fn test_fir_multiple_entries() {
let fir = Fir {
sender_ssrc: 111111,
media_ssrc: 0,
entries: vec![
FirEntry {
ssrc: 222222,
seq_nr: 1,
},
FirEntry {
ssrc: 333333,
seq_nr: 2,
},
FirEntry {
ssrc: 444444,
seq_nr: 3,
},
],
};
let bytes = fir.build();
let parsed = Fir::parse(&bytes).unwrap();
assert_eq!(parsed.entries.len(), 3);
assert_eq!(parsed.entries[0].ssrc, 222222);
assert_eq!(parsed.entries[1].ssrc, 333333);
assert_eq!(parsed.entries[2].ssrc, 444444);
}
#[test]
fn test_remb_parse_too_short() {
let data = vec![0x8F, 206, 0, 3, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0];
let result = Remb::parse(&data);
assert!(result.is_err());
}
#[test]
fn test_remb_wrong_unique_id() {
let remb = Remb::new(111111, 1_000_000, vec![222222]);
let mut bytes = remb.build().to_vec();
bytes[12] = b'X';
let result = Remb::parse(&bytes);
assert!(result.is_err());
}
#[test]
fn test_remb_wrong_type() {
let fir = Fir::new(111111, 222222, 1);
let bytes = fir.build();
let result = Remb::parse(&bytes);
assert!(result.is_err());
}
#[test]
fn test_remb_wrong_packet_type() {
let header = RtcpHeader {
version: 2,
padding: false,
count: 0,
packet_type: RtcpType::SenderReport,
length: 0,
};
let mut buf = BytesMut::new();
header.build(&mut buf);
assert_rtcp_err_contains(Remb::parse(&buf), "UnknownPacketType");
}
#[test]
fn test_remb_parse_header_too_short() {
let data = [0u8; 2];
assert_rtcp_err_contains(Remb::parse(&data), "TooShort");
}
#[test]
fn test_remb_wrong_count() {
let remb = Remb::new(111111, 1_000_000, vec![222222]);
let mut bytes = remb.build().to_vec();
bytes[0] = (bytes[0] & 0xE0) | 1;
let result = Remb::parse(&bytes);
assert!(result.is_err());
}
#[test]
fn test_remb_parse_truncated_ssrcs() {
let remb = Remb::new(111111, 1_000_000, vec![222222]);
let mut bytes = remb.build().to_vec();
bytes[16] = 2;
let parsed = Remb::parse(&bytes).unwrap();
assert_eq!(parsed.ssrcs.len(), 1);
}
#[test]
fn test_remb_zero_bitrate() {
let remb = Remb::new(111111, 0, vec![222222]);
let bytes = remb.build();
let parsed = Remb::parse(&bytes).unwrap();
assert_eq!(parsed.bitrate, 0);
}
#[test]
fn test_remb_large_bitrate() {
let remb = Remb::new(111111, 1_000_000_000, vec![222222]); let bytes = remb.build();
let parsed = Remb::parse(&bytes).unwrap();
let error = ((parsed.bitrate as i64 - 1_000_000_000i64).abs() as f64) / 1_000_000_000.0;
assert!(error < 0.01);
}
#[test]
fn test_remb_multiple_ssrcs() {
let remb = Remb::new(111111, 2_000_000, vec![222222, 333333, 444444, 555555]);
let bytes = remb.build();
let parsed = Remb::parse(&bytes).unwrap();
assert_eq!(parsed.ssrcs.len(), 4);
assert!(parsed.ssrcs.contains(&222222));
assert!(parsed.ssrcs.contains(&333333));
assert!(parsed.ssrcs.contains(&444444));
assert!(parsed.ssrcs.contains(&555555));
}
#[test]
fn test_nack_entry_single() {
let entry = NackEntry::single(1000);
assert_eq!(entry.pid, 1000);
assert_eq!(entry.blp, 0);
let seqs = entry.lost_sequences();
assert_eq!(seqs.len(), 1);
assert_eq!(seqs[0], 1000);
}
#[test]
fn test_nack_entry_from_sequences_empty() {
let result = NackEntry::from_sequences(&[]);
assert!(result.is_none());
}
#[test]
fn test_nack_from_lost_packets_large_gap() {
let nack = Nack::from_lost_packets(111111, 222222, &[1000, 2000]);
assert_eq!(nack.nacks.len(), 2);
assert_eq!(nack.nacks[0].pid, 1000);
assert_eq!(nack.nacks[1].pid, 2000);
}
#[test]
fn test_nack_from_lost_packets_in_range() {
let nack = Nack::from_lost_packets(111111, 222222, &[1000, 1002, 1010]);
assert_eq!(nack.nacks.len(), 1);
let seqs = nack.nacks[0].lost_sequences();
assert!(seqs.contains(&1000));
assert!(seqs.contains(&1002));
assert!(seqs.contains(&1010));
}
#[test]
fn test_nack_from_lost_packets_with_duplicates() {
let nack = Nack::from_lost_packets(111111, 222222, &[1000, 1000, 1001]);
assert_eq!(nack.nacks.len(), 2);
assert_eq!(nack.nacks[0].pid, 1000);
assert_eq!(nack.nacks[1].pid, 1000);
}
#[test]
fn test_nack_entry_out_of_range() {
let seqs = vec![1000, 1001, 1100]; let entry = NackEntry::from_sequences(&seqs).unwrap();
let lost = entry.lost_sequences();
assert!(lost.contains(&1000));
assert!(lost.contains(&1001));
assert!(!lost.contains(&1100)); }
#[test]
fn test_nack_all_blp_bits() {
let mut seqs = vec![1000u16];
for i in 1..=16 {
seqs.push(1000 + i);
}
let entry = NackEntry::from_sequences(&seqs).unwrap();
assert_eq!(entry.blp, 0xFFFF);
let lost = entry.lost_sequences();
assert_eq!(lost.len(), 17); }
#[test]
fn test_compound_add_fir() {
let sr = SenderReport {
ssrc: 12345,
ntp_timestamp: NtpTimestamp::now(),
rtp_timestamp: 160000,
sender_packet_count: 100,
sender_octet_count: 16000,
report_blocks: vec![],
};
let mut compound = RtcpCompound::sender_compound(sr, "test@example.com");
compound.add_fir(Fir::new(12345, 67890, 1));
let bytes = compound.build();
assert!(bytes.len() > 28);
}
#[test]
fn test_compound_add_remb() {
let sr = SenderReport {
ssrc: 12345,
ntp_timestamp: NtpTimestamp::now(),
rtp_timestamp: 160000,
sender_packet_count: 100,
sender_octet_count: 16000,
report_blocks: vec![],
};
let mut compound = RtcpCompound::sender_compound(sr, "test@example.com");
compound.add_remb(Remb::new(12345, 5_000_000, vec![67890]));
let bytes = compound.build();
assert!(bytes.len() > 28);
}
#[test]
fn test_sdes_type_enum() {
assert_eq!(SdesType::End as u8, 0);
assert_eq!(SdesType::CName as u8, 1);
assert_eq!(SdesType::Name as u8, 2);
assert_eq!(SdesType::Email as u8, 3);
assert_eq!(SdesType::Phone as u8, 4);
assert_eq!(SdesType::Location as u8, 5);
assert_eq!(SdesType::Tool as u8, 6);
assert_eq!(SdesType::Note as u8, 7);
assert_eq!(SdesType::Private as u8, 8);
}
#[test]
fn test_sdes_multiple_items() {
let sdes = SourceDescription {
chunks: vec![SdesChunk {
ssrc: 12345,
items: vec![
SdesItem {
item_type: SdesType::CName,
value: "user@host".to_string(),
},
SdesItem {
item_type: SdesType::Name,
value: "Test User".to_string(),
},
SdesItem {
item_type: SdesType::Email,
value: "test@example.com".to_string(),
},
],
}],
};
let bytes = sdes.build();
assert_eq!(bytes[1], RtcpType::SourceDescription as u8);
}
#[test]
fn test_multiple_report_blocks_sr() {
let sr = SenderReport {
ssrc: 12345,
ntp_timestamp: NtpTimestamp::now(),
rtp_timestamp: 160000,
sender_packet_count: 100,
sender_octet_count: 16000,
report_blocks: vec![
ReportBlock {
ssrc: 67890,
fraction_lost: 10,
cumulative_lost: 5,
extended_seq: 1000,
jitter: 50,
last_sr: 0x11111111,
delay_since_sr: 32768,
},
ReportBlock {
ssrc: 11111,
fraction_lost: 20,
cumulative_lost: 10,
extended_seq: 2000,
jitter: 100,
last_sr: 0x22222222,
delay_since_sr: 65536,
},
],
};
let bytes = sr.build();
let parsed = SenderReport::parse(&bytes).unwrap();
assert_eq!(parsed.report_blocks.len(), 2);
assert_eq!(parsed.report_blocks[0].ssrc, 67890);
assert_eq!(parsed.report_blocks[1].ssrc, 11111);
}
#[test]
fn test_rtcp_header_build() {
let header = RtcpHeader {
version: 2,
padding: true,
count: 5,
packet_type: RtcpType::SenderReport,
length: 10,
};
let mut buf = BytesMut::new();
header.build(&mut buf);
assert_eq!(buf.len(), 4);
assert_eq!(buf[0], 0b10100101); assert_eq!(buf[1], 200); assert_eq!(u16::from_be_bytes([buf[2], buf[3]]), 10);
}
#[test]
fn test_report_block_negative_cumulative_lost() {
let block = ReportBlock {
ssrc: 12345,
fraction_lost: 0,
cumulative_lost: -100, extended_seq: 1000,
jitter: 50,
last_sr: 0,
delay_since_sr: 0,
};
let mut buf = BytesMut::new();
block.build(&mut buf);
let parsed = ReportBlock::parse(&buf).unwrap();
assert!(parsed.ssrc == 12345);
assert!(parsed.extended_seq == 1000);
}
#[test]
fn test_transport_feedback_type_enum() {
assert_eq!(TransportFeedbackType::Nack as u8, 1);
assert_eq!(TransportFeedbackType::TransportCC as u8, 15);
}
#[test]
fn test_payload_feedback_type_enum() {
assert_eq!(PayloadFeedbackType::Pli as u8, 1);
assert_eq!(PayloadFeedbackType::Sli as u8, 2);
assert_eq!(PayloadFeedbackType::Rpsi as u8, 3);
assert_eq!(PayloadFeedbackType::Fir as u8, 4);
assert_eq!(PayloadFeedbackType::Tstr as u8, 5);
assert_eq!(PayloadFeedbackType::Tstn as u8, 6);
assert_eq!(PayloadFeedbackType::Vbcm as u8, 7);
assert_eq!(PayloadFeedbackType::Afb as u8, 15);
}
#[test]
fn roundtrip_sr_sdes() {
let sr = SenderReport {
ssrc: 0xDEAD_BEEF,
ntp_timestamp: NtpTimestamp {
seconds: 0x1234_5678,
fraction: 0x9ABC_DEF0,
},
rtp_timestamp: 160_000,
sender_packet_count: 42,
sender_octet_count: 6720,
report_blocks: vec![],
};
let compound = RtcpCompound::sender_compound(sr, "user@example.com");
let bytes = compound.build();
let parsed = RtcpCompound::parse(&bytes).expect("parse compound");
assert_eq!(parsed.packets.len(), 2);
match &parsed.packets[0] {
RtcpPacket::SenderReport(sr) => {
assert_eq!(sr.ssrc, 0xDEAD_BEEF);
assert_eq!(sr.rtp_timestamp, 160_000);
assert_eq!(sr.sender_packet_count, 42);
assert_eq!(sr.sender_octet_count, 6720);
}
other => panic!("expected SenderReport, got {other:?}"),
}
match &parsed.packets[1] {
RtcpPacket::SourceDescription(sdes) => {
assert_eq!(sdes.chunks.len(), 1);
assert_eq!(sdes.chunks[0].ssrc, 0xDEAD_BEEF);
assert_eq!(sdes.chunks[0].items.len(), 1);
assert_eq!(sdes.chunks[0].items[0].item_type, SdesType::CName);
assert_eq!(sdes.chunks[0].items[0].value, "user@example.com");
}
other => panic!("expected SourceDescription, got {other:?}"),
}
}
#[test]
fn roundtrip_rr_remb() {
let rr = ReceiverReport {
ssrc: 0x1111_2222,
report_blocks: vec![ReportBlock {
ssrc: 0x3333_4444,
fraction_lost: 0,
cumulative_lost: 0,
extended_seq: 100,
jitter: 25,
last_sr: 0,
delay_since_sr: 0,
}],
};
let mut buf = BytesMut::new();
buf.extend_from_slice(&rr.build());
let remb = Remb::new(0x1111_2222, 64_000, vec![0x3333_4444]);
buf.extend_from_slice(&remb.build());
let parsed = RtcpCompound::parse(&buf).expect("parse compound");
assert_eq!(parsed.packets.len(), 2);
match &parsed.packets[0] {
RtcpPacket::ReceiverReport(rr) => {
assert_eq!(rr.ssrc, 0x1111_2222);
assert_eq!(rr.report_blocks.len(), 1);
assert_eq!(rr.report_blocks[0].ssrc, 0x3333_4444);
}
other => panic!("expected ReceiverReport, got {other:?}"),
}
match &parsed.packets[1] {
RtcpPacket::Remb(remb) => {
assert_eq!(remb.bitrate, 64_000);
assert_eq!(remb.sender_ssrc, 0x1111_2222);
}
other => panic!("expected Remb, got {other:?}"),
}
}
#[test]
fn parse_rejects_truncated() {
let sr = SenderReport {
ssrc: 1,
ntp_timestamp: NtpTimestamp::now(),
rtp_timestamp: 0,
sender_packet_count: 0,
sender_octet_count: 0,
report_blocks: vec![],
};
let bytes = sr.build();
assert!(bytes.len() > 8);
let truncated = &bytes[..bytes.len() - 4];
let result = RtcpCompound::parse(truncated);
assert!(result.is_err(), "expected error on truncated input");
let err = result.err().unwrap();
assert!(
err.contains("truncated") || err.contains("TooShort"),
"expected truncation error, got: {err}"
);
}
#[test]
fn parse_preserves_unknown_types() {
let raw = [0x80u8, 230, 0x00, 0x00];
let parsed = RtcpCompound::parse(&raw).expect("parse must accept unknown PT");
assert_eq!(parsed.packets.len(), 1);
match &parsed.packets[0] {
RtcpPacket::Unknown { packet_type, bytes } => {
assert_eq!(*packet_type, 230);
assert_eq!(&bytes[..], &raw[..]);
}
other => panic!("expected Unknown variant, got {other:?}"),
}
}
#[test]
fn parse_walks_past_unknown_in_middle() {
let sr = SenderReport {
ssrc: 0xAAAA_AAAA,
ntp_timestamp: NtpTimestamp {
seconds: 0x0102_0304,
fraction: 0x0506_0708,
},
rtp_timestamp: 8_000,
sender_packet_count: 1,
sender_octet_count: 160,
report_blocks: vec![],
};
let rr = ReceiverReport {
ssrc: 0xBBBB_BBBB,
report_blocks: vec![],
};
let mut buf = BytesMut::new();
buf.extend_from_slice(&sr.build());
buf.extend_from_slice(&[0x80u8, 230, 0x00, 0x00]);
buf.extend_from_slice(&rr.build());
let parsed = RtcpCompound::parse(&buf).expect("parse compound");
assert_eq!(parsed.packets.len(), 3);
match &parsed.packets[0] {
RtcpPacket::SenderReport(sr) => {
assert_eq!(sr.ssrc, 0xAAAA_AAAA);
}
other => panic!("expected SenderReport, got {other:?}"),
}
match &parsed.packets[1] {
RtcpPacket::Unknown { packet_type, .. } => {
assert_eq!(*packet_type, 230);
}
other => panic!("expected Unknown variant, got {other:?}"),
}
match &parsed.packets[2] {
RtcpPacket::ReceiverReport(rr) => {
assert_eq!(rr.ssrc, 0xBBBB_BBBB);
}
other => panic!("expected ReceiverReport, got {other:?}"),
}
}
#[test]
fn psfb_fmt15_without_remb_magic_is_unknown() {
let mut buf = BytesMut::new();
buf.put_u8(0x80 | 15);
buf.put_u8(206);
buf.put_u16(4);
buf.put_u32(0xCAFE_BABE);
buf.put_u32(0x0000_0000);
buf.put_slice(b"AFBX");
buf.put_u32(0xDEAD_BEEF);
assert_eq!(buf.len(), 20);
let parsed = RtcpCompound::parse(&buf).expect("parse compound");
assert_eq!(parsed.packets.len(), 1);
match &parsed.packets[0] {
RtcpPacket::Unknown { packet_type, .. } => {
assert_eq!(*packet_type, 206);
}
other => panic!("expected Unknown variant for non-REMB AFB, got {other:?}"),
}
}
}