Skip to content
File

Blob: firmware/vendor/str0m/src/session.rs

rust1330 lines
1use std::collections::{HashMap, VecDeque};
2use std::ops::RangeInclusive;
3use std::sync::Arc;
4use std::time::{Duration, Instant};
5 
6use crate::Event;
7use crate::bwe::BweKind;
8use crate::bwe_::Bwe;
9use crate::config::KeyingMaterial;
10use crate::config_mod::RtcpReportIntervals;
11use crate::crypto::CryptoProvider;
12use crate::crypto::dtls::SrtpProfile;
13use crate::format::CodecConfig;
14use crate::format::PayloadParams;
15use crate::format::Vp9PacketizerMode;
16use crate::io::DatagramSend;
17use crate::media::AppSpecificFeedback;
18use crate::media::Media;
19use crate::media::{KeyframeRequestKind, MID_PROBE};
20use crate::media::{MediaAdded, MediaChanged};
21use crate::pacer::PacerControl;
22use crate::pacer::{Pacer, PacerImpl};
23use crate::rtp::{Extension, RawPacket};
24use crate::rtp_::Direction;
25use crate::rtp_::MidRid;
26use crate::rtp_::Pli;
27use crate::rtp_::Pt;
28use crate::rtp_::SRTCP_OVERHEAD;
29use crate::rtp_::SeqNo;
30use crate::rtp_::{Bitrate, ExtensionMap, Goodbye, Mid, ReportList, Rtcp, RtcpFb};
31use crate::rtp_::{Dlrr, DlrrItem, ExtendedReport, ReportBlock};
32use crate::rtp_::{RtpHeader, SessionId, TwccPacketId, extend_u16};
33use crate::rtp_::{SrtpContext, Ssrc};
34use crate::rtp_::{TwccRecvRegister, TwccSendRegister};
35use crate::stats::StatsSnapshot;
36use crate::streams::{RtpPacket, StreamTimeoutConfig, Streams};
37use crate::util::{Soonest, SystemTimeExt, already_happened, not_happening};
38use crate::{Reason, net};
39use crate::{RtcConfig, RtcError};
40 
41/// Minimum time we delay between sending nacks. This should be
42/// set high enough to not cause additional problems in very bad
43/// network conditions.
44const NACK_MIN_INTERVAL: Duration = Duration::from_millis(33);
45 
46/// Delay between reports of TWCC. This is deliberately very low.
47const TWCC_INTERVAL: Duration = Duration::from_millis(50);
48 
49/// Maximum number of pending RRTR entries before oldest are dropped.
50const MAX_PENDING_RRTRS: usize = 300;
51 
52/// Maximum number of DLRR items to include in a single ExtendedReport.
53const MAX_DLRR_PER_REPORT: usize = 50;
54 
55pub(crate) struct Session {
56 id: SessionId,
57 
58 // These fields are pub to allow session_sdp.rs modify them.
59 // Notice the fields are maybe not in m-line index order since the app
60 // might be spliced in somewhere.
61 pub medias: Vec<Media>,
62 
63 // The actual RTP encoded streams.
64 pub streams: Streams,
65 
66 /// The app m-line. Spliced into medias above.
67 app: Option<(Mid, usize)>,
68 
69 reordering_size_audio: usize,
70 reordering_size_video: usize,
71 intervals: RtcpReportIntervals,
72 pub send_buffer_audio: usize,
73 pub send_buffer_video: usize,
74 
75 /// Extension mappings are _per BUNDLE_, but we can only have one a=group BUNDLE
76 /// in WebRTC (one ice connection), so they are effectively per session.
77 pub exts: ExtensionMap,
78 
79 // Configuration of how we are sending/receiving media.
80 pub codec_config: CodecConfig,
81 
82 // Payload params for SSRC 0 BWE probes.
83 // Lazy init when we see the PT.
84 probe_payload_params: Option<PayloadParams>,
85 
86 srtp_rx: Option<SrtpContext>,
87 srtp_tx: Option<SrtpContext>,
88 last_nack: Instant,
89 last_twcc: Instant,
90 twcc: u64,
91 twcc_rx_register: TwccRecvRegister,
92 twcc_tx_register: TwccSendRegister,
93 max_rx_seq_lookup: HashMap<Ssrc, SeqNo>,
94 
95 bwe: Option<Bwe>,
96 
97 enable_twcc_feedback: bool,
98 
99 /// A pacer for sending RTP at specific rate.
100 pacer: PacerImpl,
101 pacer_control: PacerControl,
102 
103 // temporary buffer when getting the next (unencrypted) RTP packet from Media line.
104 poll_packet_buf: Vec<u8>,
105 
106 // Next packet for RtpPacket event.
107 pending_packet: Option<RtpPacket>,
108 
109 // Whether we sent a single outgoing RTP packet.
110 packet_first_sent: bool,
111 
112 // Total RTP payload bytes for media, including retransmissions.
113 media_bytes_rx: u64,
114 media_bytes_tx: u64,
115 
116 pub ice_lite: bool,
117 
118 /// Whether we are running in RTP-mode.
119 pub rtp_mode: bool,
120 
121 /// VP9 packetizer mode.
122 vp9_packetizer_mode: Vp9PacketizerMode,
123 
124 feedback_tx: VecDeque<Rtcp>,
125 feedback_rx: VecDeque<Rtcp>,
126 rtp_closing: bool,
127 
128 /// Pending RRTRs from remote SSRCs in FIFO order, deduplicated by SSRC.
129 /// A new RRTR for an already-pending SSRC updates its value in place.
130 /// Capped at 300 entries; at most 50 are consumed per DLRR report.
131 pending_rrtrs: VecDeque<(Ssrc, LastRrtr)>,
132 
133 raw_packets: Option<VecDeque<Box<RawPacket>>>,
134 
135 // Pending application-specific feedback (PSFB FMT=15) messages to emit as events.
136 pending_app_feedback: Option<AppSpecificFeedback>,
137 
138 /// Target MTU (start) and warn threshold (end). Buffer sizing uses the
139 /// target; oversized outgoing datagrams above the warn threshold log a warning.
140 mtu: RangeInclusive<usize>,
141 
142 #[cfg(feature = "_internal_test_exports")]
143 pending_probe: Option<crate::bwe_::ProbeClusterConfig>,
144}
145 
146/// Last RRTR (Receiver Reference Time Report, RFC 3611 §4.4) received from a remote SSRC.
147/// Stored at session level so we can respond with a DLRR regardless of which local stream
148/// the remote is observing.
149#[derive(Debug, Clone, Copy)]
150struct LastRrtr {
151 /// Middle 32 bits of the RRTR NTP timestamp (compact NTP / LRR).
152 lrr: u32,
153 /// When the RRTR arrived locally.
154 received_at: Instant,
155}
156 
157impl Session {
158 pub fn new(config: &RtcConfig) -> Self {
159 let mut id = SessionId::new();
160 // Max 2^62 - 1: https://bugzilla.mozilla.org/show_bug.cgi?id=861895
161 const MAX_ID: u64 = 2_u64.pow(62) - 1;
162 while *id > MAX_ID {
163 id = (*id >> 1).into();
164 }
165 let (pacer, bwe) = if let Some(config) = &config.bwe_config {
166 let rate = config.initial_bitrate;
167 let pacer = PacerImpl::leaky_bucket(rate * 2.0);
168 let bwe = Bwe::new(rate);
169 (pacer, Some(bwe))
170 } else {
171 (PacerImpl::null(), None)
172 };
173 
174 let enable_stats = config.stats_interval.is_some();
175 
176 Session {
177 id,
178 medias: vec![],
179 streams: Streams::new(enable_stats, *config.mtu.end()),
180 app: None,
181 reordering_size_audio: config.reordering_size_audio,
182 reordering_size_video: config.reordering_size_video,
183 intervals: config.intervals,
184 send_buffer_audio: config.send_buffer_audio,
185 send_buffer_video: config.send_buffer_video,
186 exts: config.exts.clone(),
187 
188 // Both sending and receiving starts from the configured codecs.
189 // These can then be changed in the SDP OFFER/ANSWER dance.
190 codec_config: config.codec_config.clone(),
191 probe_payload_params: None,
192 
193 srtp_rx: None,
194 srtp_tx: None,
195 last_nack: already_happened(),
196 last_twcc: already_happened(),
197 twcc: 0,
198 twcc_rx_register: TwccRecvRegister::new(100),
199 twcc_tx_register: TwccSendRegister::new(1000),
200 max_rx_seq_lookup: HashMap::new(),
201 bwe,
202 enable_twcc_feedback: false,
203 pacer,
204 pacer_control: PacerControl::new(),
205 poll_packet_buf: vec![0; 2000],
206 pending_packet: None,
207 packet_first_sent: false,
208 media_bytes_rx: 0,
209 media_bytes_tx: 0,
210 ice_lite: config.ice_lite,
211 rtp_mode: config.rtp_mode,
212 vp9_packetizer_mode: config.vp9_packetizer_mode,
213 feedback_tx: VecDeque::new(),
214 feedback_rx: VecDeque::new(),
215 rtp_closing: false,
216 pending_rrtrs: VecDeque::new(),
217 raw_packets: if config.enable_raw_packets {
218 Some(VecDeque::new())
219 } else {
220 None
221 },
222 pending_app_feedback: None,
223 mtu: config.mtu.clone(),
224 #[cfg(feature = "_internal_test_exports")]
225 pending_probe: None,
226 }
227 }
228 
229 fn mtu(&self) -> usize {
230 *self.mtu.start()
231 }
232 
233 fn mtu_warn(&self) -> usize {
234 *self.mtu.end()
235 }
236 
237 pub fn id(&self) -> SessionId {
238 self.id
239 }
240 
241 pub fn set_app(&mut self, mid: Mid, index: usize) -> Result<(), String> {
242 if let Some((mid_existing, index_existing)) = self.app {
243 if mid_existing != mid {
244 return Err(format!("App mid changed {} != {}", mid, mid_existing,));
245 }
246 if index_existing != index {
247 return Err(format!("App index changed {} != {}", index, index_existing,));
248 }
249 } else {
250 self.app = Some((mid, index));
251 }
252 Ok(())
253 }
254 
255 pub fn app(&self) -> &Option<(Mid, usize)> {
256 &self.app
257 }
258 
259 pub fn set_keying_material(
260 &mut self,
261 mat: KeyingMaterial,
262 crypto: &CryptoProvider,
263 srtp_profile: SrtpProfile,
264 active: bool,
265 ) {
266 // TODO: rename this to `initialise_srtp_context`?
267 // Whether we're active or passive determines if we use the left or right
268 // hand side of the key material to derive input/output.
269 let left = active;
270 
271 self.srtp_rx = Some(SrtpContext::new(crypto, srtp_profile, &mat, !left));
272 self.srtp_tx = Some(SrtpContext::new(crypto, srtp_profile, &mat, left));
273 }
274 
275 pub fn handle_timeout(&mut self, now: Instant) -> Result<(), RtcError> {
276 if self.rtp_closing {
277 return Ok(());
278 }
279 
280 // Payload any waiting frames
281 self.do_payload()?;
282 
283 let sender_ssrc = self.streams.first_ssrc_local();
284 
285 let do_nack = now >= self.nack_at().unwrap_or(not_happening());
286 
287 self.streams.handle_timeout(
288 now,
289 sender_ssrc,
290 do_nack,
291 &self.medias,
292 StreamTimeoutConfig {
293 codecs: &self.codec_config,
294 intervals: self.intervals,
295 },
296 &mut self.feedback_tx,
297 );
298 
299 // Flush any pending RRTRs as a DLRR alongside the Sender Report compound packet.
300 // We piggyback on the SR cadence so the DLRR doesn't generate extra standalone
301 // RTCP traffic that could perturb bandwidth estimation.
302 // Consume at most 50 entries per report; the rest remain for subsequent SRs.
303 if !self.pending_rrtrs.is_empty()
304 && self
305 .feedback_tx
306 .iter()
307 .any(|r| matches!(r, Rtcp::SenderReport(_)))
308 {
309 let n = self.pending_rrtrs.len().min(MAX_DLRR_PER_REPORT);
310 let items: Vec<DlrrItem> = self
311 .pending_rrtrs
312 .drain(..n)
313 .map(|(ssrc, rrtr)| {
314 let delay = now.saturating_duration_since(rrtr.received_at);
315 let last_rr_delay = ((delay.as_micros() * 65_536) / 1_000_000) as u32;
316 DlrrItem {
317 ssrc,
318 last_rr_time: rrtr.lrr,
319 last_rr_delay,
320 }
321 })
322 .collect();
323 
324 self.feedback_tx
325 .push_back(Rtcp::ExtendedReport(ExtendedReport {
326 ssrc: sender_ssrc,
327 blocks: vec![ReportBlock::Dlrr(Dlrr { items })],
328 }));
329 }
330 
331 if do_nack {
332 self.last_nack = now;
333 }
334 
335 self.update_queue_state(now);
336 
337 if let Some(twcc_at) = self.twcc_at() {
338 if now >= twcc_at {
339 self.create_twcc_feedback(sender_ssrc, now);
340 }
341 }
342 
343 self.handle_timeout_bwe(now);
344 
345 Ok(())
346 }
347 
348 fn handle_timeout_bwe(&mut self, now: Instant) {
349 let Some(bwe) = self.bwe.as_mut() else {
350 return;
351 };
352 
353 // We can only run probes after first packet is sent and there
354 // are any queues that can handle padding requests.
355 let do_probe = self.packet_first_sent && self.pacer.has_padding_queue();
356 
357 if let Some(probe_config) = bwe.handle_timeout(now, do_probe) {
358 // Only start the probe in the pacer if the estimator accepted it.
359 if bwe.start_probe(probe_config, now) {
360 #[cfg(feature = "_internal_test_exports")]
361 {
362 self.pending_probe = Some(probe_config);
363 }
364 self.pacer.start_probe(probe_config);
365 }
366 }
367 
368 // Check if active probe just completed
369 if let Some(cluster_id) = self.pacer.check_probe_complete(now) {
370 bwe.end_probe(now, cluster_id);
371 }
372 }
373 
374 fn update_queue_state(&mut self, now: Instant) {
375 let iter = self.streams.streams_tx().map(|m| m.queue_state(now));
376 
377 let Some(padding_request) = self.pacer.handle_timeout(now, iter) else {
378 return;
379 };
380 
381 let stream = self
382 .streams
383 .stream_tx_by_midrid(padding_request.midrid)
384 .expect("pacer to use an existing stream");
385 
386 stream.generate_padding(padding_request.padding);
387 }
388 
389 fn create_twcc_feedback(&mut self, sender_ssrc: Ssrc, now: Instant) -> Option<()> {
390 self.last_twcc = now;
391 let mut twcc = self.twcc_rx_register.build_report(self.mtu() - 100)?;
392 
393 // These SSRC are on media level, but twcc is on session level,
394 // we fill in the first discovered media SSRC in each direction.
395 twcc.sender_ssrc = sender_ssrc;
396 twcc.ssrc = self.streams.first_ssrc_remote();
397 
398 trace!("Created feedback TWCC: {:?}", twcc);
399 self.feedback_tx.push_front(Rtcp::Twcc(twcc));
400 Some(())
401 }
402 
403 /// Enqueue an application-specific feedback message (PSFB FMT=15, PT=206) for
404 /// transmission as part of the regular RTCP compound packet.
405 pub(crate) fn send_app_specific_feedback(
406 &mut self,
407 sender_ssrc: Ssrc,
408 media_ssrc: Ssrc,
409 payload: impl Into<Arc<[u8]>>,
410 ) {
411 use crate::rtp_::AppSpecificFeedback as RtcpAppFeedback;
412 let feedback = RtcpAppFeedback {
413 sender_ssrc,
414 media_ssrc,
415 payload: payload.into(),
416 };
417 self.feedback_tx
418 .push_back(Rtcp::AppSpecificFeedback(feedback));
419 }
420 
421 /// Enqueue a Picture Loss Indication (PLI, PT=206 FMT=1) for `media_ssrc`, to be
422 /// sent as part of the regular RTCP compound packet. The feedback packet's sender
423 /// SSRC is set to `sender_ssrc`.
424 pub(crate) fn send_pli_feedback(&mut self, sender_ssrc: Ssrc, media_ssrc: Ssrc) {
425 self.feedback_tx.push_back(Rtcp::Pli(Pli {
426 sender_ssrc,
427 ssrc: media_ssrc,
428 }));
429 }
430 
431 pub fn handle_rtp_receive(&mut self, now: Instant, message: &[u8]) {
432 let Some(header) = RtpHeader::parse(message, &self.exts) else {
433 trace!("Failed to parse RTP header");
434 return;
435 };
436 
437 self.handle_rtp(now, header, message);
438 }
439 
440 pub fn handle_rtcp_receive(&mut self, now: Instant, message: &[u8]) {
441 // According to spec, the outer enclosing SRTCP packet should always be a SR or RR,
442 // even if it's irrelevant and empty.
443 // In practice I'm not sure that is happening, because libWebRTC hates empty packets.
444 self.handle_rtcp(now, message);
445 }
446 
447 fn mid_and_ssrc_for_header(&mut self, now: Instant, header: &RtpHeader) -> Option<(Mid, Ssrc)> {
448 let ssrc_header = header.ssrc;
449 
450 if let Some(r) = self.streams.mid_ssrc_rx_by_ssrc_or_rtx(now, ssrc_header) {
451 return Some(r);
452 }
453 
454 // Attempt to dynamically map this header to some Media/ReceiveStream.
455 self.map_dynamic(header);
456 
457 // The dynamic mapping might have added an entry by now.
458 if let Some(r) = self.streams.mid_ssrc_rx_by_ssrc_or_rtx(now, ssrc_header) {
459 return Some(r);
460 }
461 
462 // SSRC 0 is used for non-media BWE probes from libwebrtc.
463 // These probes need TWCC feedback but don't carry actual media.
464 if ssrc_header.is_probe() {
465 return self.ensure_probe_stream(header.payload_type);
466 }
467 
468 None
469 }
470 
471 /// Creates the probe stream on-demand for handling SSRC 0 BWE probes.
472 ///
473 /// No Media is created since probes don't carry real media - they only need
474 /// SRTP decryption and TWCC feedback.
475 fn ensure_probe_stream(&mut self, pt: Pt) -> Option<(Mid, Ssrc)> {
476 let ssrc: Ssrc = 0.into();
477 let midrid = MidRid(MID_PROBE, None);
478 
479 // Add PayloadParams for this PT if not already configured.
480 // Probes may use different PTs, so we add them as we see them.
481 self.probe_payload_params = Some(PayloadParams::new_probe(pt));
482 
483 // Create the stream with NACK suppressed (probes don't need retransmission)
484 self.streams.expect_stream_rx(ssrc, None, midrid, true);
485 
486 Some((MID_PROBE, ssrc))
487 }
488 
489 fn map_dynamic(&mut self, header: &RtpHeader) {
490 // There are two strategies for dynamically mapping SSRC. Both use the RTP "mid"
491 // header extension.
492 // A) Mid+Rid - used when doing simulcast. Rid points out which
493 // simulcast layer is in use. There is a separate header
494 // to indicate repair (RTX) stream.
495 // B) Mid+PT - when not doing simulcast, the PT identifies whether
496 // this is a repair stream.
497 
498 let Some(mid) = header.ext_vals.mid else {
499 return;
500 };
501 let rid = header.ext_vals.rid.or(header.ext_vals.rid_repair);
502 
503 // The media the mid points out. Bail if the mid points to something
504 // we don't know about.
505 let Some(media) = self.medias.iter_mut().find(|m| m.mid() == mid) else {
506 return;
507 };
508 
509 // Figure out which payload the PT maps to. Either main or RTX.
510 let maybe_payload = self
511 .codec_config
512 .iter()
513 .find(|p| p.pt() == header.payload_type || p.resend() == Some(header.payload_type));
514 
515 // If we don't find it, bail out.
516 let Some(payload) = maybe_payload else {
517 return;
518 };
519 
520 if let Some(rid) = rid {
521 // Case A - use the rid_repair header to identify RTX.
522 let is_main = header.ext_vals.rid.is_some();
523 
524 let midrid = MidRid(mid, Some(rid));
525 
526 self.streams
527 .map_dynamic_by_rid(header.ssrc, midrid, media, *payload, is_main);
528 } else {
529 // Case B - the payload type identifies RTX.
530 let is_main = payload.pt() == header.payload_type;
531 
532 let midrid = MidRid(mid, None);
533 
534 self.streams
535 .map_dynamic_by_pt(header.ssrc, midrid, media, *payload, is_main);
536 }
537 }
538 
539 pub(crate) fn handle_rtp(&mut self, now: Instant, mut header: RtpHeader, buf: &[u8]) {
540 // Rewrite absolute-send-time (if present) to be relative to now.
541 header.ext_vals.update_absolute_send_time(now);
542 
543 trace!("Handle RTP: {:?}", header);
544 
545 // The ssrc is the _main_ ssrc (no the rtx, that might be in the header).
546 let Some((mid, ssrc)) = self.mid_and_ssrc_for_header(now, &header) else {
547 debug!("No mid/SSRC for header: {:?}", header);
548 return;
549 };
550 
551 let srtp = match self.srtp_rx.as_mut() {
552 Some(v) => v,
553 None => {
554 trace!("Rejecting SRTP while missing SrtpContext");
555 return;
556 }
557 };
558 
559 // mid_and_ssrc_for_header guarantees stream exists.
560 // Media may not exist for internal probe streams (SSRC 0).
561 let stream = self.streams.stream_rx(&ssrc).unwrap();
562 
563 let maybe_params = if ssrc.is_probe() {
564 self.probe_payload_params.as_ref()
565 } else {
566 main_payload_params(&self.codec_config, header.payload_type)
567 };
568 
569 let params = match maybe_params {
570 Some(p) => p,
571 None => {
572 trace!(
573 "No payload params could be found (main or RTX) for {:?}",
574 header.payload_type
575 );
576 return;
577 }
578 };
579 let clock_rate = params.spec().rtp_clock_rate();
580 let pt = params.pt();
581 let is_repair = pt != header.payload_type;
582 
583 let max_seq_lookup = make_max_seq_lookup(&self.max_rx_seq_lookup);
584 
585 // is_repair controls whether update is updating the main register or the RTX register.
586 // Either way we get a seq_no_outer which is used to decrypt the SRTP.
587 let mut seq_no = stream.extend_seq(&header, is_repair, max_seq_lookup);
588 
589 if !stream.is_new_packet(is_repair, seq_no) {
590 // Dupe packet. This could be a potential SRTP replay attack, which means
591 // we should not spend any CPU cycles towards decrypting it.
592 trace!(
593 "Ignoring dupe packet mid: {} seq_no: {} is_repair: {}",
594 mid, seq_no, is_repair
595 );
596 return;
597 }
598 
599 // The decrypted plaintext borrows from the SRTP scratch buffer. We
600 // narrow it down to the actual payload via slicing only, then build
601 // the final `Arc<[u8]>` in a single allocation at the end.
602 let data = match srtp.unprotect_rtp(buf, &header, *seq_no) {
603 Some(d) => d,
604 None => {
605 trace!(
606 "Failed to unprotect SRTP for SSRC: {} pt: {} mid: {} \
607 rid: {:?} seq_no: {} is_repair: {}",
608 header.ssrc,
609 pt,
610 stream.mid(),
611 stream.rid(),
612 seq_no,
613 is_repair
614 );
615 return;
616 }
617 };
618 
619 let data = if header.has_padding {
620 match RtpHeader::unpad_payload(data) {
621 Some(d) => d,
622 None => {
623 trace!("unpadding of unprotected payload failed");
624 return;
625 }
626 }
627 } else {
628 data
629 };
630 
631 if let Some(raw_packets) = &mut self.raw_packets {
632 raw_packets.push_back(Box::new(RawPacket::RtpRx(header.clone(), data.to_vec())));
633 }
634 
635 // Mark as received for TWCC purposes
636 if let Some(transport_cc) = header.ext_vals.transport_cc {
637 let prev = self.twcc_rx_register.max_seq();
638 let extended = extend_u16(prev.map(|s| *s), transport_cc);
639 self.twcc_rx_register.update_seq(extended.into(), now);
640 }
641 
642 // Store largest seen seq_no for the SSRC. This is used in case we get SSRC changes
643 // like A -> B -> A. When we go back to A, we must keep the ROC.
644 update_max_seq(&mut self.max_rx_seq_lookup, header.ssrc, seq_no);
645 
646 // Register reception in nack registers.
647 let receipt_outer = stream.update_register(now, &header, clock_rate, is_repair, seq_no);
648 
649 // RTX packets must be rewritten to be a normal packet. This only changes the
650 // the seq_no, however MediaTime might be different when interpreted against the
651 // the "main" register.
652 let (receipt, data) = if is_repair {
653 // Drop RTX packets that are just empty padding. The payload here
654 // is empty because we would have done RtpHeader::unpad_payload above.
655 // For unpausing, it's enough with the stream.update() already done above.
656 if data.is_empty() {
657 return;
658 }
659 
660 // Rewrite the header, and strip the resent seq_no prefix from the body.
661 let data = stream.un_rtx(&mut header, data, pt);
662 
663 let max_seq_lookup = make_max_seq_lookup(&self.max_rx_seq_lookup);
664 
665 // Header has changed, which means we extend a new seq_no. This time
666 // without is_repair since this is the wrapped resend. This is the
667 // extended number of the main stream.
668 seq_no = stream.extend_seq(&header, false, max_seq_lookup);
669 
670 // header.ssrc is changed by un_rtx() above to be main SSRC.
671 update_max_seq(&mut self.max_rx_seq_lookup, header.ssrc, seq_no);
672 
673 // Now update the "main" register with the repaired packet info.
674 let receipt = stream.update_register(now, &header, clock_rate, false, seq_no);
675 (receipt, data)
676 } else {
677 // This is not RTX, the outer seq and time is what we use. The first
678 // stream.update will have updated the main register.
679 (receipt_outer, data)
680 };
681 
682 // Probe packets (SSRC 0) contain only padding, no real media to process
683 if ssrc.is_probe() {
684 return;
685 }
686 
687 self.media_bytes_rx += data.len() as u64;
688 
689 // One-shot conversion to Arc<[u8]>: a single allocation that copies
690 // only the trimmed payload bytes out of the SRTP scratch buffer.
691 let payload: Arc<[u8]> = Arc::from(data);
692 
693 let packet = stream.handle_rtp(now, header, payload, seq_no, receipt.time);
694 
695 if self.rtp_mode {
696 // In RTP mode, we store the packet temporarily here for the next poll_output().
697 // However only if this is a packet not seen before. This filters out spurious
698 // resends for padding.
699 if receipt.is_new_packet {
700 self.pending_packet = Some(packet);
701 }
702 } else {
703 // In non-RTP mode, we let the Media use a Depayloader.
704 // unwrap is fine because mid_and_ssrc_for_header guarantees it.
705 let media = self.medias.iter_mut().find(|m| m.mid() == mid).unwrap();
706 media.depayload(
707 stream.rid(),
708 packet,
709 self.reordering_size_audio,
710 self.reordering_size_video,
711 &self.codec_config,
712 );
713 }
714 }
715 
716 fn handle_rtcp(&mut self, now: Instant, buf: &[u8]) -> Option<()> {
717 let srtp: &mut SrtpContext = self.srtp_rx.as_mut()?;
718 let unprotected = srtp.unprotect_rtcp(buf)?;
719 
720 Rtcp::read_packet(&unprotected, &mut self.feedback_rx);
721 let mut need_configure_pacer = false;
722 
723 if let Some(raw_packets) = &mut self.raw_packets {
724 for fb in &self.feedback_rx {
725 raw_packets.push_back(Box::new(RawPacket::RtcpRx(fb.clone())));
726 }
727 }
728 
729 for fb in RtcpFb::from_rtcp(self.feedback_rx.drain(..)) {
730 if let RtcpFb::AppSpecificFeedback(asf) = fb {
731 self.pending_app_feedback = Some(AppSpecificFeedback {
732 sender_ssrc: asf.sender_ssrc,
733 media_ssrc: asf.media_ssrc,
734 payload: asf.payload,
735 });
736 continue;
737 }
738 
739 if let RtcpFb::Twcc(twcc) = fb {
740 trace!("Handle TWCC: {:?}", twcc);
741 let maybe_records = self.twcc_tx_register.apply_report(twcc, now);
742 
743 if let (Some(maybe_records), Some(bwe)) = (maybe_records, &mut self.bwe) {
744 bwe.update(maybe_records, now);
745 }
746 need_configure_pacer = true;
747 
748 // The funky thing about TWCC reports is that they are never stapled
749 // together with other RTCP packet. If they were though, we want to
750 // handle more packets.
751 continue;
752 }
753 
754 if let RtcpFb::Rrtr((rrtr, ssrc)) = fb {
755 // DLRR responder (RFC 3611 §4.5): store the remote's RRTR so we can
756 // reply with a DLRR, letting the remote compute RTT.
757 let lrr = (rrtr.ntp_time.as_ntp_64() >> 16) as u32;
758 let entry = LastRrtr {
759 lrr,
760 received_at: now,
761 };
762 
763 // Update in place if this SSRC is already pending (preserves queue order).
764 if let Some(existing) = self.pending_rrtrs.iter_mut().find(|(s, _)| *s == ssrc) {
765 existing.1 = entry;
766 } else {
767 // Cap at 300 entries; drop the oldest if full.
768 if self.pending_rrtrs.len() >= MAX_PENDING_RRTRS {
769 self.pending_rrtrs.pop_front();
770 }
771 self.pending_rrtrs.push_back((ssrc, entry));
772 }
773 continue;
774 }
775 
776 if fb.is_for_rx() {
777 let Some(stream) = self.streams.stream_rx(&fb.ssrc()) else {
778 continue;
779 };
780 stream.handle_rtcp(now, fb);
781 } else {
782 let Some(stream) = self.streams.stream_tx(&fb.ssrc()) else {
783 continue;
784 };
785 stream.handle_rtcp(now, fb);
786 }
787 }
788 
789 // Not in the above if due to lifetime issues, still okay because the method
790 // doesn't do anything when BWE isn't configured.
791 if need_configure_pacer {
792 self.configure_pacer();
793 }
794 
795 Some(())
796 }
797 
798 pub fn poll_event(&mut self) -> Option<Event> {
799 #[cfg(feature = "_internal_test_exports")]
800 {
801 if let Some(probe) = self.pending_probe.take() {
802 return Some(Event::Probe(probe));
803 }
804 }
805 
806 if let Some(bitrate_estimate) = self.bwe.as_mut().and_then(|bwe| bwe.poll_estimate()) {
807 return Some(Event::EgressBitrateEstimate(BweKind::Twcc(
808 bitrate_estimate,
809 )));
810 }
811 
812 // If we're not ready to flow media, don't send any events.
813 if !self.ready_for_srtp() {
814 return None;
815 }
816 
817 if let Some(raw_packets) = &mut self.raw_packets {
818 if let Some(p) = raw_packets.pop_front() {
819 return Some(Event::RawPacket(p));
820 }
821 }
822 
823 if let Some(feedback) = self.pending_app_feedback.take() {
824 return Some(Event::AppSpecificFeedback(feedback));
825 }
826 
827 // This must be before pending_packet.take() since we need to emit the unpaused event
828 // before the first packet causing the unpause.
829 if let Some(paused) = self.streams.poll_stream_paused() {
830 if paused.paused {
831 if let Some(media) = self.medias.iter_mut().find(|m| m.mid() == paused.mid) {
832 // Drop held partial depacketizer state so pre-pause fragments can't
833 // complete into stale MediaData after the stream resumes.
834 media.reset_depayloaders_for_rid(paused.rid);
835 }
836 }
837 return Some(Event::StreamPaused(paused));
838 }
839 
840 if self.rtp_mode {
841 if let Some(packet) = self.pending_packet.take() {
842 return Some(Event::RtpPacket(packet));
843 }
844 }
845 
846 if let Some(req) = self.streams.poll_keyframe_request() {
847 return Some(Event::KeyframeRequest(req));
848 }
849 
850 if let Some(report) = self.streams.poll_sender_feedback() {
851 return Some(Event::SenderFeedback(report));
852 }
853 
854 if let Some((mid, bitrate)) = self.streams.poll_remb_request() {
855 return Some(Event::EgressBitrateEstimate(BweKind::Remb(mid, bitrate)));
856 }
857 
858 for media in &mut self.medias {
859 if media.need_open_event {
860 media.need_open_event = false;
861 
862 return Some(Event::MediaAdded(MediaAdded {
863 mid: media.mid(),
864 kind: media.kind(),
865 direction: media.direction(),
866 simulcast: media.simulcast().map(|s| s.clone().into()),
867 }));
868 }
869 
870 if media.need_changed_event {
871 media.need_changed_event = false;
872 return Some(Event::MediaChanged(MediaChanged {
873 mid: media.mid(),
874 direction: media.direction(),
875 }));
876 }
877 }
878 
879 None
880 }
881 
882 pub fn poll_event_fallible(&mut self) -> Result<Option<Event>, RtcError> {
883 // Not relevant in rtp_mode, where the packets are picked up by poll_event().
884 if self.rtp_mode {
885 return Ok(None);
886 }
887 
888 for media in &mut self.medias {
889 if let Some(e) = media.poll_sample(&self.codec_config)? {
890 return Ok(Some(Event::MediaData(e)));
891 }
892 }
893 
894 Ok(None)
895 }
896 
897 fn ready_for_srtp(&self) -> bool {
898 self.srtp_rx.is_some() && self.srtp_tx.is_some()
899 }
900 
901 pub fn poll_datagram(&mut self, now: Instant) -> Option<net::DatagramSend> {
902 // Time must have progressed forward from start value.
903 if now == already_happened() {
904 return None;
905 }
906 
907 let x = if self.rtp_closing {
908 self.poll_feedback()
909 } else {
910 None.or_else(|| self.poll_feedback())
911 .or_else(|| self.poll_packet(now))
912 };
913 
914 if let Some(x) = &x {
915 // In RTP mode we trust the API user feeds the RTP packet sizes they
916 // need for the MTU they are targeting. This warning is only for when
917 // str0m does the RTP packetization.
918 if !self.rtp_mode && x.len() > self.mtu_warn() {
919 warn!("RTP above MTU {}: {}", self.mtu_warn(), x.len());
920 }
921 }
922 
923 x
924 }
925 
926 /// To be called in lieu of [`Self::poll_datagram`] when the owner is not in a position to transmit any
927 /// generated feedback, and thus such feedback should be dropped.
928 pub fn clear_feedback(&mut self) {
929 self.feedback_rx.clear();
930 self.feedback_tx.clear();
931 }
932 
933 pub fn close_rtp(&mut self) {
934 if self.rtp_closing {
935 return;
936 }
937 
938 self.rtp_closing = true;
939 self.clear_feedback();
940 self.streams.reset_send_buffers();
941 
942 for reports in ReportList::lists_from_iter(self.streams.local_sender_ssrcs()) {
943 self.feedback_tx
944 .push_back(Rtcp::Goodbye(Goodbye { reports }));
945 }
946 }
947 
948 pub fn is_rtp_closed(&self) -> bool {
949 self.rtp_closing && self.feedback_tx.is_empty()
950 }
951 
952 fn poll_feedback(&mut self) -> Option<net::DatagramSend> {
953 if self.feedback_tx.is_empty() {
954 return None;
955 }
956 
957 // Round to nearest multiple of 4 bytes.
958 let encryptable_mtu: usize = (self.mtu() - SRTCP_OVERHEAD) & !3;
959 assert!(encryptable_mtu % 4 == 0);
960 
961 let mut data = vec![0_u8; encryptable_mtu];
962 
963 let mut raw_packets = self.raw_packets.as_mut();
964 let output = move |fb| {
965 if let Some(raw_packets) = &mut raw_packets {
966 raw_packets.push_back(Box::new(RawPacket::RtcpTx(fb)));
967 }
968 };
969 
970 let len = Rtcp::write_packet(&mut self.feedback_tx, &mut data, output);
971 
972 if len == 0 {
973 return None;
974 }
975 
976 data.truncate(len);
977 
978 let srtp = self.srtp_tx.as_mut()?;
979 let protected = srtp.protect_rtcp(&data);
980 
981 assert!(
982 protected.len() < self.mtu(),
983 "Encrypted SRTCP should be less than MTU"
984 );
985 
986 Some(protected.into())
987 }
988 
989 fn poll_packet(&mut self, now: Instant) -> Option<DatagramSend> {
990 let srtp_tx = self.srtp_tx.as_mut()?;
991 
992 // Figure out which, if any, queue to poll
993 // The cluster_id is captured by the pacer at poll time, before register_send() might clear it
994 let (midrid, cluster_id) = self.pacer.poll_queue()?;
995 let Some(media) = self.medias.iter().find(|m| m.mid() == midrid.mid()) else {
996 trace!("Pacer pointed to mid {} which has no media", midrid.mid());
997 return None;
998 };
999 
1000 let buf = &mut self.poll_packet_buf;
1001 let twcc_seq = self.twcc;
1002 
1003 let stream = self.streams.stream_tx_by_midrid(midrid)?;
1004 
1005 let params = &self.codec_config;
1006 let exts = media.remote_extmap();
1007 
1008 // TWCC might not be enabled for this m-line. Firefox do use TWCC, but not
1009 // for audio. This is indiciated via the SDP.
1010 let twcc_enabled = exts.id_of(Extension::TransportSequenceNumber).is_some();
1011 let twcc = twcc_enabled.then_some(&mut self.twcc);
1012 
1013 let receipt = stream.poll_packet(now, exts, twcc, params, buf)?;
1014 
1015 let PacketReceipt {
1016 header,
1017 seq_no,
1018 is_padding,
1019 payload_size,
1020 } = receipt;
1021 
1022 trace!(payload_size, is_padding, "Poll RTP: {:?}", header);
1023 
1024 #[cfg(feature = "_internal_dont_use_log_stats")]
1025 {
1026 let kind = if is_padding { "padding" } else { "media" };
1027 
1028 crate::log_stat!("PACKET_SENT", header.ssrc, payload_size, kind);
1029 }
1030 
1031 self.pacer.register_send(now, payload_size.into(), midrid);
1032 
1033 if let Some(raw_packets) = &mut self.raw_packets {
1034 raw_packets.push_back(Box::new(RawPacket::RtpTx(header.clone(), buf.clone())));
1035 }
1036 
1037 let protected = srtp_tx.protect_rtp(buf, &header, *seq_no);
1038 
1039 if twcc_enabled {
1040 let packet_id = if let Some(cluster) = cluster_id {
1041 TwccPacketId::with_cluster(twcc_seq, cluster)
1042 } else {
1043 TwccPacketId::new(twcc_seq)
1044 };
1045 self.twcc_tx_register
1046 .register_seq(packet_id, now, payload_size);
1047 }
1048 
1049 // Update BWE subsystem
1050 if let Some(bwe) = self.bwe.as_mut() {
1051 bwe.on_media_sent(payload_size.into(), is_padding, now);
1052 }
1053 
1054 if !is_padding && !header.ssrc.is_probe() {
1055 self.media_bytes_tx += payload_size as u64;
1056 }
1057 
1058 if !self.packet_first_sent {
1059 self.packet_first_sent = true;
1060 }
1061 
1062 Some(protected.into())
1063 }
1064 
1065 pub fn poll_timeout(&mut self) -> (Option<Instant>, Reason) {
1066 if self.rtp_closing {
1067 return if self.feedback_tx.is_empty() {
1068 (None, Reason::NotHappening)
1069 } else {
1070 (Some(already_happened()), Reason::Feedback)
1071 };
1072 }
1073 
1074 let feedback_at = self.regular_feedback_at();
1075 let nack_at = self.nack_at();
1076 let twcc_at = self.twcc_at();
1077 let pacing_at = self.pacer.poll_timeout();
1078 let packetize_at = self.medias.iter().flat_map(|m| m.poll_timeout()).next();
1079 let paused_at = self.paused_at();
1080 let send_stream_at = self.streams.send_stream();
1081 
1082 // Gives us built-in reason
1083 let bwe_at = self
1084 .bwe
1085 .as_ref()
1086 .map(|bwe| bwe.poll_timeout())
1087 // We should never see this reason.
1088 .unwrap_or((None, Reason::BweDelayControl));
1089 
1090 (feedback_at, Reason::Feedback)
1091 .soonest((nack_at, Reason::Nack))
1092 .soonest((twcc_at, Reason::Twcc))
1093 .soonest(pacing_at)
1094 .soonest((packetize_at, Reason::Packetize))
1095 .soonest((paused_at, Reason::PauseCheck))
1096 .soonest((send_stream_at, Reason::SendStream))
1097 .soonest(bwe_at)
1098 }
1099 
1100 pub fn has_mid(&self, mid: Mid) -> bool {
1101 self.medias.iter().any(|m| m.mid() == mid)
1102 }
1103 
1104 fn regular_feedback_at(&self) -> Option<Instant> {
1105 self.streams.regular_feedback_at(self.intervals)
1106 }
1107 
1108 fn paused_at(&self) -> Option<Instant> {
1109 self.streams.paused_at()
1110 }
1111 
1112 fn nack_at(&mut self) -> Option<Instant> {
1113 if !self.streams.any_nack_enabled() {
1114 return None;
1115 }
1116 
1117 Some(self.last_nack + NACK_MIN_INTERVAL)
1118 }
1119 
1120 fn twcc_at(&self) -> Option<Instant> {
1121 let is_receiving = self.streams.is_receiving();
1122 if is_receiving && self.enable_twcc_feedback && self.twcc_rx_register.has_unreported() {
1123 Some(self.last_twcc + TWCC_INTERVAL)
1124 } else {
1125 None
1126 }
1127 }
1128 
1129 pub fn enable_twcc_feedback(&mut self) {
1130 if !self.enable_twcc_feedback {
1131 debug!("Enable TWCC feedback");
1132 self.enable_twcc_feedback = true;
1133 }
1134 }
1135 
1136 pub fn visit_stats(&mut self, now: Instant, snapshot: &mut StatsSnapshot) {
1137 for stream in self.streams.streams_tx() {
1138 stream.visit_stats(snapshot, now);
1139 }
1140 
1141 for stream in self.streams.streams_rx() {
1142 stream.visit_stats(snapshot, now);
1143 }
1144 
1145 snapshot.tx = self.media_bytes_tx;
1146 snapshot.rx = self.media_bytes_rx;
1147 snapshot.bwe_tx = self.bwe.as_ref().and_then(|bwe| bwe.last_estimate());
1148 
1149 snapshot.egress_loss_fraction = self.twcc_tx_register.loss(Duration::from_secs(1), now);
1150 snapshot.rtt = self.twcc_tx_register.rtt();
1151 snapshot.ingress_loss_fraction = self.twcc_rx_register.loss();
1152 }
1153 
1154 pub fn set_bwe_desired_bitrate(&mut self, desired_bitrate: Bitrate) {
1155 if let Some(bwe) = self.bwe.as_mut() {
1156 bwe.set_desired_bitrate(desired_bitrate);
1157 self.configure_pacer();
1158 }
1159 }
1160 
1161 pub fn reset_bwe(&mut self, init_bitrate: Bitrate) {
1162 if let Some(bwe) = self.bwe.as_mut() {
1163 bwe.reset(init_bitrate);
1164 }
1165 }
1166 
1167 #[allow(dead_code)]
1168 pub fn line_count(&self) -> usize {
1169 self.medias.len() + if self.app.is_some() { 1 } else { 0 }
1170 }
1171 
1172 pub fn add_media(&mut self, media: Media) {
1173 self.medias.push(media);
1174 }
1175 
1176 pub fn medias(&self) -> &[Media] {
1177 &self.medias
1178 }
1179 
1180 pub fn remove_media(&mut self, mid: Mid) {
1181 self.medias.retain(|media| media.mid() != mid);
1182 self.streams.remove_streams_by_mid(mid);
1183 }
1184 
1185 fn configure_pacer(&mut self) {
1186 let Some(bwe) = self.bwe.as_mut() else {
1187 return;
1188 };
1189 
1190 let Some(current_estimate) = bwe.last_estimate() else {
1191 // No estimate yet, no padding
1192 return;
1193 };
1194 let is_overuse = bwe.is_overusing();
1195 
1196 let has_active_media = self.has_active_outgoing_media();
1197 
1198 // Calculate pacing and padding rates
1199 let result = self
1200 .pacer_control
1201 .calculate(has_active_media, current_estimate, is_overuse);
1202 
1203 self.pacer.set_padding_rate(result.padding_rate);
1204 self.pacer.set_pacing_rate(result.pacing_rate);
1205 }
1206 
1207 fn has_active_outgoing_media(&self) -> bool {
1208 self.medias
1209 .iter()
1210 .any(|m| m.direction().is_sending() && !m.disabled())
1211 }
1212 
1213 pub fn media_by_mid(&self, mid: Mid) -> Option<&Media> {
1214 self.medias.iter().find(|m| m.mid() == mid)
1215 }
1216 
1217 pub fn media_by_mid_mut(&mut self, mid: Mid) -> Option<&mut Media> {
1218 self.medias.iter_mut().find(|m| m.mid() == mid)
1219 }
1220 
1221 fn do_payload(&mut self) -> Result<(), RtcError> {
1222 let mtu = self.mtu();
1223 for m in &mut self.medias {
1224 m.do_payload(
1225 &mut self.streams,
1226 &self.codec_config,
1227 self.vp9_packetizer_mode,
1228 mtu,
1229 )?;
1230 }
1231 
1232 Ok(())
1233 }
1234 
1235 pub fn set_direction(&mut self, mid: Mid, direction: Direction) -> bool {
1236 let Some(media) = self.media_by_mid_mut(mid) else {
1237 return false;
1238 };
1239 let old_dir = media.direction();
1240 if old_dir == direction {
1241 return false;
1242 }
1243 
1244 media.set_direction(direction);
1245 
1246 if old_dir.is_sending() && !direction.is_sending() {
1247 self.streams.reset_buffers_tx(mid);
1248 }
1249 
1250 let max_seq_lookup = make_max_seq_lookup(&self.max_rx_seq_lookup);
1251 if old_dir.is_receiving() && !direction.is_receiving() {
1252 self.streams.reset_buffers_rx(mid, max_seq_lookup);
1253 }
1254 
1255 true
1256 }
1257 
1258 pub fn stop_media(&mut self, mid: Mid) -> bool {
1259 let Some(media) = self.media_by_mid_mut(mid) else {
1260 return false;
1261 };
1262 if media.disabled() {
1263 return false;
1264 }
1265 
1266 media.mark_stopped();
1267 self.set_direction(mid, Direction::Inactive);
1268 
1269 true
1270 }
1271 
1272 pub fn is_request_keyframe_possible(&self, kind: KeyframeRequestKind) -> bool {
1273 // TODO: It's possible to have different set of feedback enabled for different
1274 // payload types. I.e. we could have FIR enabled for H264, but not for VP8.
1275 // We might want to make this check more fine grained by testing which PT is
1276 // in "active use" right now.
1277 self.codec_config.iter().any(|r| match kind {
1278 KeyframeRequestKind::Pli => r.fb_pli,
1279 KeyframeRequestKind::Fir => r.fb_fir,
1280 })
1281 }
1282 
1283 /// Checks whether the SRTP contexts are up.
1284 pub fn is_connected(&self) -> bool {
1285 self.srtp_rx.is_some() && self.srtp_tx.is_some()
1286 }
1287}
1288 
1289pub struct PacketReceipt {
1290 pub header: RtpHeader,
1291 pub seq_no: SeqNo,
1292 pub is_padding: bool,
1293 pub payload_size: usize,
1294}
1295 
1296/// Find the PayloadParams for the given Pt, either when the Pt is the main Pt for the Codec or
1297/// when it's the RTX Pt.
1298fn main_payload_params(c: &CodecConfig, pt: Pt) -> Option<&PayloadParams> {
1299 c.iter().find(|p| p.pt == pt || p.resend == Some(pt))
1300}
1301 
1302fn make_max_seq_lookup(map: &HashMap<Ssrc, SeqNo>) -> impl Fn(Ssrc) -> Option<SeqNo> + '_ {
1303 |ssrc| map.get(&ssrc).cloned()
1304}
1305 
1306fn update_max_seq(map: &mut HashMap<Ssrc, SeqNo>, ssrc: Ssrc, seq_no: SeqNo) {
1307 let current = map.entry(ssrc).or_insert(seq_no);
1308 if seq_no > *current {
1309 *current = seq_no;
1310 }
1311}
1312 
1313#[cfg(test)]
1314mod tests {
1315 use super::*;
1316 use crate::RtcConfig;
1317 use crate::io::DATAGRAM_MTU_TARGET;
1318 
1319 #[test]
1320 fn session_mtu_matches_config() {
1321 let cfg = RtcConfig::default();
1322 let s = Session::new(&cfg);
1323 assert_eq!(s.mtu(), DATAGRAM_MTU_TARGET);
1324 
1325 let cfg = RtcConfig::default().set_mtu(900..=1280);
1326 let s = Session::new(&cfg);
1327 assert_eq!(s.mtu(), 900);
1328 }
1329}