Skip to content

Commit c055b68

Browse files
committed
feat: add PLI/FIR RTCP feedback packets
1 parent 0b17954 commit c055b68

17 files changed

Lines changed: 1249 additions & 212 deletions

File tree

media/rtc/src/rtp_session/inbound/mod.rs

Lines changed: 76 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,9 @@ mod stats;
1313

1414
pub use stats::{RtpInboundRemoteStats, RtpInboundStats};
1515

16+
/// Minimum interval in which FIR/PLI requests can be sent
17+
const RTCP_FEEDBACK_COOLDOWN: Duration = Duration::from_millis(250);
18+
1619
/// RTP receive stream
1720
pub struct RtpInboundStream {
1821
ssrc: Ssrc,
@@ -22,6 +25,15 @@ pub struct RtpInboundStream {
2225
last_received_sender_report: Option<NtpTimestamp>,
2326

2427
remote_stats: Option<RtpInboundRemoteStats>,
28+
29+
// RTCP feedback NACK PLI
30+
want_nack_pli: bool,
31+
last_nack_pli: Option<Instant>,
32+
33+
// RTCP feedback CCM FIR
34+
want_ccm_fir: bool,
35+
next_fir_seq: u8,
36+
last_ccm_fir: Option<Instant>,
2537
}
2638

2739
impl RtpInboundStream {
@@ -33,32 +45,87 @@ impl RtpInboundStream {
3345
last_report_sent: None,
3446
last_received_sender_report: None,
3547
remote_stats: None,
48+
49+
want_nack_pli: false,
50+
last_nack_pli: None,
51+
52+
want_ccm_fir: false,
53+
next_fir_seq: rand::random(),
54+
last_ccm_fir: None,
3655
}
3756
}
3857

58+
pub fn request_nack_pli(&mut self) {
59+
self.want_nack_pli = true
60+
}
61+
62+
pub fn request_ccm_fir(&mut self) {
63+
self.want_ccm_fir = true
64+
}
65+
3966
pub(crate) fn timeout(&self, now: Instant) -> Option<Duration> {
4067
let queue = self.queue.timeout(now);
4168

4269
let report = if self.queue.highest_sequence_number_received().is_some() {
43-
Some(
44-
self.last_report_sent
45-
.and_then(|(last_report_sent, _)| {
46-
(last_report_sent + self.report_interval).checked_duration_since(now)
47-
})
48-
.unwrap_or_default(),
49-
)
70+
let report_interval = self
71+
.last_report_sent
72+
.and_then(|(last_report_sent, _)| {
73+
(last_report_sent + self.report_interval).checked_duration_since(now)
74+
})
75+
.unwrap_or_default();
76+
77+
let nack_pli = self
78+
.last_nack_pli
79+
.map(|ts| (ts + RTCP_FEEDBACK_COOLDOWN).saturating_duration_since(now));
80+
81+
let ccm_fir = self
82+
.last_ccm_fir
83+
.map(|ts| (ts + RTCP_FEEDBACK_COOLDOWN).saturating_duration_since(now));
84+
85+
opt_min(Some(report_interval), opt_min(nack_pli, ccm_fir))
5086
} else {
5187
None
5288
};
5389

5490
opt_min(queue, report)
5591
}
5692

57-
pub(crate) fn collect_reports(&mut self, now: Instant, reports: &mut ReportsQueue) {
58-
let make_report = self
93+
pub(super) fn collect_reports(&mut self, now: Instant, reports: &mut ReportsQueue) {
94+
if self.want_nack_pli {
95+
let cooldown_elapsed = self
96+
.last_nack_pli
97+
.is_none_or(|i| i + RTCP_FEEDBACK_COOLDOWN <= now);
98+
99+
if cooldown_elapsed {
100+
self.want_nack_pli = false;
101+
self.last_nack_pli = Some(now);
102+
reports.add_nack_pli(self.ssrc);
103+
}
104+
}
105+
106+
if self.want_ccm_fir {
107+
let cooldown_elapsed = self
108+
.last_ccm_fir
109+
.is_none_or(|i| i + RTCP_FEEDBACK_COOLDOWN <= now);
110+
111+
if cooldown_elapsed {
112+
self.want_ccm_fir = false;
113+
self.last_ccm_fir = Some(now);
114+
reports.add_ccm_fir(self.ssrc, self.next_fir_seq);
115+
self.next_fir_seq = self.next_fir_seq.wrapping_add(1);
116+
}
117+
}
118+
119+
// When emitting feedback packets & reduced size RTCP is not supported:
120+
// Emit a receiver report so the ReportsQueue can generate a valid RTCP packet
121+
let receiver_report_for_feedback = reports.has_feedback() && !reports.rtcp_rsize();
122+
123+
let report_interval_elapsed = self
59124
.last_report_sent
60125
.is_none_or(|(instant, _)| now > instant + self.report_interval);
61126

127+
let make_report = receiver_report_for_feedback || report_interval_elapsed;
128+
62129
if !make_report {
63130
return;
64131
}

media/rtc/src/rtp_session/mod.rs

Lines changed: 52 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,7 @@ use crate::Mtu;
99
use report::ReportsQueue;
1010
use rtp::{
1111
RtpPacket, Ssrc,
12-
rtcp_types::{Compound, Packet as RtcpPacket},
12+
rtcp_types::{Compound, Fir, Packet as RtcpPacket, Pli, RtcpPacketParser},
1313
};
1414
use ssrc_hasher::SsrcHasher;
1515
use std::{
@@ -44,15 +44,20 @@ pub struct RtpSession {
4444

4545
impl RtpSession {
4646
/// Create a new RTP session
47-
pub fn new() -> Self {
47+
pub fn new(rtcp_rsize: bool) -> Self {
4848
Self {
4949
next_tx_ssrc: Ssrc(rand::random()),
5050
tx: HashMap::default(),
5151
rx: HashMap::default(),
52-
reports: ReportsQueue::default(),
52+
reports: ReportsQueue::new(rtcp_rsize),
5353
}
5454
}
5555

56+
/// Is reduced size RTCP allowed
57+
pub fn rtcp_rsize(&self) -> bool {
58+
self.reports.rtcp_rsize()
59+
}
60+
5661
/// Create a new outbound RTP stream with the given parameters
5762
pub fn new_tx_stream(&mut self, clock_rate: u32) -> &mut RtpOutboundStream {
5863
let ssrc = self.next_tx_ssrc;
@@ -100,13 +105,20 @@ impl RtpSession {
100105
}
101106

102107
/// Hand of the RTCP packet to the RTP session
103-
pub fn receive_rtcp(&mut self, now: Instant, rtcp_packet: Compound<'_>) {
108+
#[must_use]
109+
pub fn receive_rtcp(
110+
&mut self,
111+
now: Instant,
112+
rtcp_packet: Compound<'_>,
113+
) -> Vec<RtpSessionReceiveRtcpEvent> {
114+
let mut events = vec![];
115+
104116
for rtcp_packet in rtcp_packet {
105117
let rtcp_packet = match rtcp_packet {
106118
Ok(rtcp_packet) => rtcp_packet,
107119
Err(e) => {
108120
log::warn!("Failed to parse RTCP packet in compound packet, {e}");
109-
return;
121+
return events;
110122
}
111123
};
112124

@@ -139,14 +151,36 @@ impl RtpSession {
139151
RtcpPacket::TransportFeedback(_transport_feedback) => {
140152
// TODO: handle feedback
141153
}
142-
RtcpPacket::PayloadFeedback(_payload_feedback) => {
143-
// TODO: handle feedback
154+
RtcpPacket::PayloadFeedback(payload_feedback) => {
155+
if payload_feedback.parse_fci::<Pli>().is_ok() {
156+
events.push(RtpSessionReceiveRtcpEvent::NackPliReceived(Ssrc(
157+
payload_feedback.media_ssrc(),
158+
)));
159+
} else if let Ok(fir) = payload_feedback.parse_fci::<Fir>() {
160+
for entry in fir.entries() {
161+
events.push(RtpSessionReceiveRtcpEvent::CcmFirReceived(Ssrc(
162+
entry.ssrc(),
163+
)));
164+
}
165+
} else {
166+
log::warn!(
167+
"Received unknown RTCP payload feedback packet header={:02X?} sender_ssrc={} media_ssrc={}",
168+
payload_feedback.header_data(),
169+
payload_feedback.sender_ssrc(),
170+
payload_feedback.media_ssrc(),
171+
)
172+
}
173+
}
174+
RtcpPacket::Xr(_xr) => {
175+
// ignore
144176
}
145177
RtcpPacket::Unknown(..) => {
146178
// ignore
147179
}
148180
}
149181
}
182+
183+
events
150184
}
151185

152186
/// Returns the duration to wait from the given Instant before polling again
@@ -165,7 +199,7 @@ impl RtpSession {
165199
}
166200

167201
/// Poll the session for any new events
168-
pub fn poll(&mut self, now: Instant, mtu: Mtu) -> Option<RtpSessionEvent> {
202+
pub fn poll(&mut self, now: Instant, mtu: Mtu) -> Option<RtpSessionPollEvent> {
169203
let fallback_sender_ssrc = *self.tx.keys().next().unwrap_or(&self.next_tx_ssrc);
170204

171205
for tx in self.tx.values_mut() {
@@ -177,34 +211,34 @@ impl RtpSession {
177211
}
178212

179213
if let Some(report) = self.reports.make_report(fallback_sender_ssrc, mtu) {
180-
return Some(RtpSessionEvent::SendRtcp(report));
214+
return Some(RtpSessionPollEvent::SendRtcp(report));
181215
}
182216

183217
for tx in self.tx.values_mut() {
184218
if let Some(rtp_packet) = tx.pop(now) {
185-
return Some(RtpSessionEvent::SendRtp(rtp_packet));
219+
return Some(RtpSessionPollEvent::SendRtp(rtp_packet));
186220
}
187221
}
188222

189223
for rx in self.rx.values_mut() {
190224
if let Some(rtp_packet) = rx.pop(now) {
191-
return Some(RtpSessionEvent::ReceiveRtp(rtp_packet));
225+
return Some(RtpSessionPollEvent::ReceiveRtp(rtp_packet));
192226
}
193227
}
194228

195229
None
196230
}
197231
}
198232

199-
impl Default for RtpSession {
200-
fn default() -> Self {
201-
Self::new()
202-
}
203-
}
204-
205233
/// Event returned by [`RtpSession::poll`]
206-
pub enum RtpSessionEvent {
234+
pub enum RtpSessionPollEvent {
207235
ReceiveRtp(RtpPacket),
208236
SendRtp(RtpPacket),
209237
SendRtcp(Vec<u8>),
210238
}
239+
240+
/// Event returned by [`RtpSession::receive_rtcp`]
241+
pub enum RtpSessionReceiveRtcpEvent {
242+
NackPliReceived(Ssrc),
243+
CcmFirReceived(Ssrc),
244+
}

media/rtc/src/rtp_session/outbound/mod.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -61,7 +61,7 @@ impl RtpOutboundStream {
6161
opt_min(queue, report)
6262
}
6363

64-
pub(crate) fn collect_reports(&mut self, now: Instant, reports: &mut ReportsQueue) {
64+
pub(super) fn collect_reports(&mut self, now: Instant, reports: &mut ReportsQueue) {
6565
let make_report = self
6666
.last_report_sent
6767
.is_none_or(|instant| now > instant + self.report_interval);

0 commit comments

Comments
 (0)