Skip to content
File

Blob: firmware/vendor/str0m/src/streams/mod.rs

rust726 lines
1use std::collections::{HashMap, VecDeque};
2use std::fmt::{self};
3use std::sync::Arc;
4use std::time::Duration;
5use std::time::Instant;
6 
7use crate::config_mod::RtcpReportIntervals;
8use crate::format::CodecConfig;
9use crate::format::PayloadParams;
10use crate::media::{KeyframeRequest, Media, SenderFeedback};
11use crate::packet::Vp8Patch;
12use crate::rtp_::MidRid;
13use crate::rtp_::Ssrc;
14use crate::rtp_::{Bitrate, Pt};
15use crate::rtp_::{MediaTime, SenderInfo};
16use crate::rtp_::{Mid, Rid, SeqNo};
17use crate::rtp_::{Rtcp, RtpHeader};
18use crate::util::already_happened;
19 
20pub use self::receive::StreamRx;
21pub use self::send::{RtpWrite, StreamTx, StreamTxQueueInfo};
22 
23mod receive;
24pub(crate) mod register;
25pub(crate) mod register_nack;
26mod rtx_cache;
27pub(crate) mod rtx_cache_buf;
28mod send;
29mod send_queue;
30mod send_stats;
31 
32pub(crate) use send::{DEFAULT_RTX_CACHE_DURATION, DEFAULT_RTX_RATIO_CAP};
33 
34pub(crate) struct StreamTimeoutConfig<'a> {
35 pub(crate) codecs: &'a CodecConfig,
36 pub(crate) intervals: RtcpReportIntervals,
37}
38 
39/// Packet of RTP data.
40///
41/// As emitted by [`Event::RtpPacket`][crate::Event::RtpPacket] when using rtp mode.
42#[derive(PartialEq, Eq)]
43pub struct RtpPacket {
44 /// Extended sequence number to avoid having to deal with ROC.
45 pub seq_no: SeqNo,
46 
47 /// Extended RTP time in the clock frequency of the codec. To avoid dealing with ROC.
48 ///
49 /// For a newly scheduled outgoing packet, the clock_rate is not correctly set until
50 /// we do the poll_output().
51 pub time: MediaTime,
52 
53 /// Parsed RTP header.
54 pub header: RtpHeader,
55 
56 /// RTP payload. This contains no header.
57 pub payload: Arc<[u8]>,
58 
59 vp8_patch: Option<Vp8Patch>,
60 
61 /// str0m server timestamp.
62 ///
63 /// This timestamp has nothing to do with RTP itself. For outgoing packets, this is when
64 /// the packet was first handed over to str0m and enqueued in the outgoing send buffers.
65 /// For incoming packets it's the time we received the network packet.
66 pub timestamp: Instant,
67 
68 /// Sender information from the most recent Sender Report(SR).
69 ///
70 /// If no Sender Report(SR) has been received or this packet is being sent by str0m this is [`None`].
71 pub last_sender_info: Option<SenderInfo>,
72 
73 /// Whether this packet can be nacked.
74 ///
75 /// This is often false for audio, but might also be false for discardable frames when
76 /// using temporal encoding as in a VP8 simulcast situation.
77 pub(crate) nackable: bool,
78}
79 
80/// Event when an encoded stream is considered paused/unpaused.
81///
82/// This means the stream has not received any data for some time (default 1.5 seconds).
83#[derive(Debug)]
84pub struct StreamPaused {
85 /// The main SSRC of the encoded stream that paused.
86 pub ssrc: Ssrc,
87 
88 /// The mid the encoded stream belongs to.
89 pub mid: Mid,
90 
91 /// The rid, if the encoded stream has a rid.
92 pub rid: Option<Rid>,
93 
94 /// Whether the stream is paused or not.
95 pub paused: bool,
96}
97 
98/// 255 is out of range for a real PT, which is 7 bit.
99const BLANK_PACKET_DEFAULT_PT: Pt = Pt::new_with_value(255);
100 
101impl RtpPacket {
102 fn blank() -> RtpPacket {
103 RtpPacket {
104 seq_no: 0.into(),
105 time: MediaTime::from_90khz(0),
106 header: RtpHeader {
107 payload_type: BLANK_PACKET_DEFAULT_PT,
108 ..Default::default()
109 },
110 payload: Arc::default(), // This payload is never used. See RtpHeader::create_padding_packet
111 vp8_patch: None,
112 nackable: false,
113 last_sender_info: None,
114 timestamp: already_happened(),
115 }
116 }
117}
118 
119/// Holder of incoming/outgoing encoded streams.
120///
121/// Each encoded stream is uniquely identified by an SSRC. The concept of mid/rid sits on the Media
122/// level together with the ability to translate a mid/rid to an encoded stream.
123#[derive(Debug)]
124pub(crate) struct Streams {
125 /// All incoming encoded streams.
126 streams_rx: HashMap<Ssrc, StreamRx>,
127 
128 /// Each incoming SSRC is mapped to a Mid/Ssrc. The Ssrc in the value is for the case
129 /// where the incoming SSRC is for an RTX and we want the "main".
130 rx_lookup: HashMap<Ssrc, RxLookup>,
131 
132 /// Time we last cleaned up unused entries from source_keys_rx.
133 last_rx_lookup_cleanup: Instant,
134 
135 /// All outgoing encoded streams.
136 streams_tx: HashMap<Ssrc, StreamTx>,
137 
138 /// Local SSRC used before we got any StreamTx. This is used for RTCP if we don't
139 /// have any reasonable value to use.
140 default_ssrc_tx: Ssrc,
141 
142 /// We need to report all RR/SR for a Mid together in one RTCP. This is a dynamic
143 /// list that we don't want to allocate on every handle_timeout.
144 mids_to_report: Vec<Mid>,
145 
146 /// Whether nack reports are enabled. This is an optimization to avoid too frequent
147 /// Session::nack_at() when we don't need to send nacks.
148 any_nack_active: Option<bool>,
149 
150 /// Whether periodic statistics reports are expected to be generated. This informs us on
151 /// whether we should be holding onto data needed for those reports or not.
152 enable_stats: bool,
153 
154 /// Threshold above which an outgoing RTP packet triggers an MTU warning. Used as a
155 /// hard cap when selecting RTX cache entries for spurious padding.
156 mtu_warn: usize,
157}
158 
159/// Delay between cleaning up the RxLookup.
160const RX_LOOKUP_CLEANUP_INTERVAL: Duration = Duration::from_millis(10_000);
161 
162/// How old an RxLookup entry may be.
163const RX_LOOKUP_EXPIRY: Duration = Duration::from_millis(30_000);
164 
165#[derive(Debug)]
166struct RxLookup {
167 mid: Mid,
168 main: Ssrc,
169 last_used: Instant,
170}
171 
172impl Streams {
173 pub(crate) fn new(enable_stats: bool, mtu_warn: usize) -> Self {
174 Self {
175 streams_rx: Default::default(),
176 rx_lookup: Default::default(),
177 last_rx_lookup_cleanup: already_happened(),
178 streams_tx: Default::default(),
179 default_ssrc_tx: 0.into(), // this will be changed
180 mids_to_report: Vec::with_capacity(10),
181 any_nack_active: None,
182 enable_stats,
183 mtu_warn,
184 }
185 }
186 
187 pub(crate) fn map_dynamic_by_rid(
188 &mut self,
189 ssrc: Ssrc,
190 midrid: MidRid,
191 media: &mut Media,
192 payload: PayloadParams,
193 is_main: bool,
194 ) {
195 // This is the point of the function.
196 let rid = midrid
197 .rid()
198 .expect("map_dynamic_by_rid to be called with Rid");
199 
200 // Check if the mid/rid combo is not expected
201 if !media.rids_rx().contains(rid) {
202 trace!("Mid does not expect rid: {} {}", midrid.mid(), rid);
203 return;
204 }
205 
206 let maybe_stream = self.stream_rx_by_midrid(midrid, false);
207 
208 let (ssrc_main, rtx) = if is_main {
209 let maybe_rtx = maybe_stream.and_then(|s| s.rtx());
210 (ssrc, maybe_rtx)
211 } else {
212 // This can bail if the main SSRC has not been discovered yet.
213 let Some(stream) = maybe_stream else {
214 return;
215 };
216 (stream.ssrc(), Some(ssrc))
217 };
218 
219 self.map_dynamic_finish(midrid, ssrc_main, rtx, media, payload);
220 }
221 
222 pub(crate) fn map_dynamic_by_pt(
223 &mut self,
224 ssrc: Ssrc,
225 midrid: MidRid,
226 media: &mut Media,
227 payload: PayloadParams,
228 is_main: bool,
229 ) {
230 if media.rids_rx().is_specific() {
231 trace!(
232 "Media expects rid and RTP packet has only mid: {:?}",
233 media.rids_rx()
234 );
235 return;
236 }
237 
238 let maybe_stream = self.stream_rx_by_midrid(midrid, false);
239 
240 let (ssrc_main, rtx) = if is_main {
241 let maybe_rtx = maybe_stream.and_then(|s| s.rtx());
242 (ssrc, maybe_rtx)
243 } else {
244 // This can bail if the main SSRC has not been discovered yet.
245 let Some(stream) = maybe_stream else {
246 return;
247 };
248 // The main is the SSRC in the stream already. The incoming is RTX.
249 let ssrc_main = stream.ssrc();
250 (ssrc_main, Some(ssrc))
251 };
252 
253 self.map_dynamic_finish(midrid, ssrc_main, rtx, media, payload);
254 }
255 
256 #[allow(clippy::too_many_arguments)]
257 fn map_dynamic_finish(
258 &mut self,
259 midrid: MidRid,
260 ssrc_main: Ssrc,
261 rtx: Option<Ssrc>,
262 media: &mut Media,
263 payload: PayloadParams,
264 ) {
265 let maybe_stream = self.stream_rx_by_midrid(midrid, false);
266 
267 if let Some(stream) = maybe_stream {
268 let ssrc_from = stream.ssrc();
269 let rtx_from = stream.rtx();
270 
271 // Handle changes in SSRC.
272 if ssrc_from != ssrc_main {
273 // We got a change in main SSRC for this stream.
274 let did_change = self.change_stream_rx_ssrc(ssrc_from, ssrc_main);
275 
276 // When the SSRCs changes the sequence number typically also does, the
277 // depayloader (if in use) relies on sequence numbers and will not handle a
278 // large jump correctly, reset it.
279 if did_change {
280 media.reset_depayloader(payload.pt(), midrid.rid());
281 }
282 }
283 
284 // Handle changes in RTX
285 if let (Some(rtx_from), Some(rtx_to)) = (rtx_from, rtx) {
286 if rtx_from != rtx_to {
287 self.change_stream_rx_rtx(rtx_from, rtx_to);
288 }
289 }
290 }
291 
292 // If we don't have an RTX PT configured, we don't want NACK.
293 let suppress_nack = payload.resend.is_none();
294 
295 // If stream already exists, this might only "fill in" the RTX.
296 self.expect_stream_rx(ssrc_main, rtx, midrid, suppress_nack);
297 }
298 
299 pub fn expect_stream_rx(
300 &mut self,
301 ssrc: Ssrc,
302 rtx: Option<Ssrc>,
303 midrid: MidRid,
304 suppress_nack: bool,
305 ) -> &mut StreamRx {
306 // New stream might have enabled nacks.
307 self.any_nack_active = None;
308 
309 let stream = self
310 .streams_rx
311 .entry(ssrc)
312 .or_insert_with(|| StreamRx::new(ssrc, midrid, suppress_nack));
313 
314 if let Some(rtx) = rtx {
315 stream.maybe_reset_rtx(rtx);
316 }
317 
318 stream
319 }
320 
321 pub fn remove_stream_rx(&mut self, ssrc: Ssrc) -> bool {
322 let stream = self.streams_rx.remove(&ssrc);
323 let existed = stream.is_some();
324 
325 self.rx_lookup.retain(|k, l| *k != ssrc && l.main != ssrc);
326 
327 existed
328 }
329 
330 pub fn declare_stream_tx(
331 &mut self,
332 ssrc: Ssrc,
333 rtx: Option<Ssrc>,
334 midrid: MidRid,
335 ) -> &mut StreamTx {
336 self.streams_tx
337 .entry(ssrc)
338 .or_insert_with(|| StreamTx::new(ssrc, rtx, midrid, self.enable_stats, self.mtu_warn))
339 }
340 
341 pub fn remove_stream_tx(&mut self, ssrc: Ssrc) -> bool {
342 self.streams_tx.remove(&ssrc).is_some()
343 }
344 
345 pub(crate) fn local_sender_ssrcs(&self) -> Vec<Ssrc> {
346 self.streams_tx
347 .values()
348 .flat_map(|s| std::iter::once(s.ssrc()).chain(s.rtx()))
349 .collect()
350 }
351 
352 pub(crate) fn reset_send_buffers(&mut self) {
353 for stream in self.streams_tx.values_mut() {
354 stream.reset_buffers();
355 }
356 }
357 
358 pub fn stream_rx(&mut self, ssrc: &Ssrc) -> Option<&mut StreamRx> {
359 self.streams_rx.get_mut(ssrc)
360 }
361 
362 pub fn stream_tx(&mut self, ssrc: &Ssrc) -> Option<&mut StreamTx> {
363 self.streams_tx.get_mut(ssrc)
364 }
365 
366 /// Lookup the "main" SSRC and mid for a given SSRC(main or RTX).
367 pub(crate) fn mid_ssrc_rx_by_ssrc_or_rtx(
368 &mut self,
369 now: Instant,
370 ssrc: Ssrc,
371 ) -> Option<(Mid, Ssrc)> {
372 // A direct hit on SSRC is to prefer. The idea is that mid/rid are only sent
373 // for the initial x seconds and then we start using SSRC only instead.
374 if let Some(r) = self.rx_lookup.get_mut(&ssrc) {
375 r.last_used = now;
376 return Some((r.mid, r.main));
377 }
378 
379 let maybe_stream = self.stream_rx_by_ssrc_or_rtx(ssrc);
380 if let Some(stream) = maybe_stream {
381 let mid = stream.mid();
382 let ssrc_main = stream.ssrc();
383 
384 self.rx_lookup.insert(
385 ssrc,
386 RxLookup {
387 mid,
388 main: ssrc_main,
389 last_used: now,
390 },
391 );
392 
393 return Some((mid, ssrc_main));
394 }
395 
396 None
397 }
398 
399 pub(crate) fn regular_feedback_at(&self, i: RtcpReportIntervals) -> Option<Instant> {
400 let r = self.streams_rx.values().map(|s| s.receiver_report_at(i));
401 let s = self.streams_tx.values().map(|s| s.sender_report_at(i));
402 r.chain(s).min()
403 }
404 
405 pub(crate) fn paused_at(&self) -> Option<Instant> {
406 self.streams_rx.values().find_map(|s| s.paused_at())
407 }
408 
409 pub(crate) fn send_stream(&self) -> Option<Instant> {
410 if self.streams_tx.values().any(|s| s.need_timeout()) {
411 Some(already_happened())
412 } else {
413 None
414 }
415 }
416 
417 pub(crate) fn is_receiving(&self) -> bool {
418 !self.streams_rx.is_empty()
419 }
420 
421 pub(crate) fn handle_timeout(
422 &mut self,
423 now: Instant,
424 sender_ssrc: Ssrc,
425 do_nack: bool,
426 medias: &[Media],
427 config: StreamTimeoutConfig<'_>,
428 feedback: &mut VecDeque<Rtcp>,
429 ) {
430 let StreamTimeoutConfig { codecs, intervals } = config;
431 
432 self.mids_to_report.clear(); // Clear for checking StreamRx.
433 for stream in self.streams_rx.values() {
434 if stream.need_rr(now, intervals) {
435 self.mids_to_report.push(stream.mid());
436 }
437 }
438 
439 for stream in self.streams_rx.values_mut() {
440 stream.maybe_create_keyframe_request(sender_ssrc, feedback);
441 stream.maybe_create_remb_request(sender_ssrc, feedback);
442 
443 // All StreamRx belonging to the same Mid are reported together.
444 if self.mids_to_report.contains(&stream.mid()) {
445 stream.create_rr_and_update(now, sender_ssrc, feedback);
446 }
447 
448 if do_nack {
449 stream.maybe_create_nack(sender_ssrc, feedback);
450 }
451 
452 stream.handle_timeout(now);
453 }
454 
455 self.mids_to_report.clear(); // start over for StreamTx.
456 for stream in self.streams_tx.values() {
457 if stream.need_sr(now, intervals) {
458 self.mids_to_report.push(stream.mid());
459 }
460 }
461 
462 for stream in self.streams_tx.values_mut() {
463 let mid = stream.mid();
464 
465 // All StreamTx belonging to the same Mid are reported together.
466 if self.mids_to_report.contains(&mid) {
467 stream.create_sr_and_update(now, feedback);
468 }
469 
470 // Finding the first (main) PT that also has RTX for the Media is expensive,
471 // this closure is run only when needed.
472 // The unwrap is okay because we cannot have StreamTx with a Mid without the corresponding Media.
473 let get_media = move || (medias.iter().find(|m| m.mid() == mid).unwrap(), codecs);
474 
475 stream.handle_timeout(now, get_media);
476 }
477 
478 if now > self.rx_lookup_at() {
479 self.rx_lookup
480 .retain(|_, l| now - l.last_used <= RX_LOOKUP_EXPIRY);
481 self.last_rx_lookup_cleanup = now;
482 }
483 }
484 
485 pub(crate) fn poll_keyframe_request(&mut self) -> Option<KeyframeRequest> {
486 self.streams_tx.values_mut().find_map(|s| {
487 let kind = s.poll_keyframe_request()?;
488 Some(KeyframeRequest {
489 mid: s.mid(),
490 rid: s.rid(),
491 kind,
492 })
493 })
494 }
495 
496 pub(crate) fn poll_sender_feedback(&mut self) -> Option<SenderFeedback> {
497 self.streams_rx.values_mut().find_map(|s| {
498 let (sender_info, at) = s.poll_sender_info()?;
499 
500 Some(SenderFeedback {
501 mid: s.mid(),
502 rid: s.rid(),
503 received_at: at,
504 sender_info,
505 })
506 })
507 }
508 
509 pub(crate) fn poll_remb_request(&mut self) -> Option<(Mid, Bitrate)> {
510 self.streams_tx
511 .values_mut()
512 .find_map(|s| s.poll_remb_request().map(|b| (s.mid(), b)))
513 }
514 
515 pub(crate) fn poll_stream_paused(&mut self) -> Option<StreamPaused> {
516 self.streams_rx.values_mut().find_map(|s| s.poll_paused())
517 }
518 
519 pub(crate) fn has_stream_rx(&self, ssrc: Ssrc) -> bool {
520 self.streams_rx.contains_key(&ssrc)
521 }
522 
523 pub(crate) fn has_stream_tx(&self, ssrc: Ssrc) -> bool {
524 self.streams_tx.contains_key(&ssrc)
525 }
526 
527 pub(crate) fn streams_rx(&mut self) -> impl Iterator<Item = &mut StreamRx> {
528 self.streams_rx.values_mut()
529 }
530 
531 pub(crate) fn streams_tx(&mut self) -> impl Iterator<Item = &mut StreamTx> {
532 self.streams_tx.values_mut()
533 }
534 
535 pub(crate) fn ssrcs_tx(&self, mid: Mid) -> Vec<(Ssrc, Option<Ssrc>)> {
536 self.streams_tx
537 .values()
538 .filter(|s| s.mid() == mid)
539 .map(|s| (s.ssrc(), s.rtx()))
540 .collect()
541 }
542 
543 pub(crate) fn new_ssrc(&self) -> Ssrc {
544 loop {
545 // This new guarantees we never get SSRC 0 (reserved for BWE probes)
546 let ssrc = Ssrc::new();
547 
548 let has_ssrc = self.has_stream_rx(ssrc) || self.has_stream_tx(ssrc);
549 
550 if has_ssrc {
551 continue;
552 }
553 
554 // Need to check RTX as well.
555 let has_rtx_rx = self.streams_rx.values().any(|s| s.rtx() == Some(ssrc));
556 if has_rtx_rx {
557 continue;
558 }
559 
560 let has_rtx_tx = self.streams_tx.values().any(|s| s.rtx() == Some(ssrc));
561 if has_rtx_tx {
562 continue;
563 }
564 
565 // Not used
566 break ssrc;
567 }
568 }
569 
570 pub fn new_ssrc_pair(&self) -> (Ssrc, Ssrc) {
571 let ssrc = self.new_ssrc();
572 
573 let rtx = loop {
574 let proposed = self.new_ssrc();
575 // Avoid clashing with just allocated main SSRC.
576 if proposed != ssrc {
577 break proposed;
578 }
579 };
580 
581 (ssrc, rtx)
582 }
583 
584 pub(crate) fn first_ssrc_remote(&self) -> Ssrc {
585 *self.streams_rx.keys().next().unwrap_or(&0.into())
586 }
587 
588 pub(crate) fn first_ssrc_local(&mut self) -> Ssrc {
589 if let Some(ssrc) = self.streams_tx.keys().next() {
590 // If there is some local Tx SSRC, use that.
591 *ssrc
592 } else {
593 // Fallback in case we don't have any Tx SSRC.
594 if *self.default_ssrc_tx == 0 {
595 // Do not use 0, allocate one that is not in session already.
596 self.default_ssrc_tx = self.new_ssrc();
597 }
598 self.default_ssrc_tx
599 }
600 }
601 
602 pub(crate) fn stream_tx_by_midrid(&mut self, midrid: MidRid) -> Option<&mut StreamTx> {
603 self.streams_tx.values_mut().find(|s| s.is_midrid(midrid))
604 }
605 
606 pub(crate) fn stream_rx_by_midrid(
607 &mut self,
608 midrid: MidRid,
609 reset_cached_nack_flag: bool,
610 ) -> Option<&mut StreamRx> {
611 if reset_cached_nack_flag {
612 // Invalidate nack_active since it's possible to manipulate the
613 // nack setting on the returned StreamRx.
614 self.any_nack_active = None;
615 }
616 
617 self.streams_rx.values_mut().find(|s| s.is_midrid(midrid))
618 }
619 
620 pub(crate) fn remove_streams_by_mid(&mut self, mid: Mid) {
621 self.streams_tx.retain(|_, s| s.mid() != mid);
622 self.streams_rx.retain(|_, s| s.mid() != mid);
623 self.rx_lookup.retain(|_, v| v.mid != mid);
624 }
625 
626 /// An iterator over all the tx streams for a given mid.
627 pub(crate) fn streams_tx_by_mid(&mut self, mid: Mid) -> impl Iterator<Item = &mut StreamTx> {
628 self.streams_tx.values_mut().filter(move |s| s.mid() == mid)
629 }
630 
631 /// An iterator over all the rx streams for a given mid.
632 pub(crate) fn streams_rx_by_mid(&mut self, mid: Mid) -> impl Iterator<Item = &mut StreamRx> {
633 self.streams_rx.values_mut().filter(move |s| s.mid() == mid)
634 }
635 
636 pub(crate) fn reset_buffers_tx(&mut self, mid: Mid) {
637 for s in self.streams_tx_by_mid(mid) {
638 s.reset_buffers();
639 }
640 }
641 
642 pub(crate) fn reset_buffers_rx(
643 &mut self,
644 mid: Mid,
645 max_seq_lookup: impl Fn(Ssrc) -> Option<SeqNo>,
646 ) {
647 for s in self.streams_rx_by_mid(mid) {
648 s.reset_buffers(&max_seq_lookup);
649 }
650 }
651 
652 pub(crate) fn change_stream_rx_ssrc(&mut self, ssrc_from: Ssrc, ssrc_to: Ssrc) -> bool {
653 // This unwrap is OK, because we can't call change_stream_rx_ssrc without first
654 // knowing there is such a StreamRx.
655 let maybe_change = self.streams_rx.get_mut(&ssrc_from).unwrap();
656 
657 // The StreamRx is allowed to not change the SSRC in case it is switching back
658 // to the previous SSRC. This is to avoid flapping in case RTP packets arrive
659 // out of order.
660 let did_change = maybe_change.change_ssrc(ssrc_to);
661 
662 if did_change {
663 // Unwrap is OK, see above.
664 let to_change = self.streams_rx.remove(&ssrc_from).unwrap();
665 
666 // Reinsert under new SSRC key.
667 self.streams_rx.insert(ssrc_to, to_change);
668 
669 // Remove previous mappings for the SSRC
670 self.rx_lookup
671 .retain(|k, l| *k != ssrc_from && l.main != ssrc_from);
672 }
673 
674 did_change
675 }
676 
677 fn change_stream_rx_rtx(&mut self, rtx_from: Ssrc, rtx_to: Ssrc) {
678 // Invalidate since we might need to enable nacks now.
679 self.any_nack_active = None;
680 
681 // Remove the SSRC mapping
682 self.rx_lookup.remove(&rtx_from);
683 
684 let Some(to_change) = self
685 .streams_rx
686 .values_mut()
687 .find(|s| s.rtx() == Some(rtx_from))
688 else {
689 // If there's no main stream associated with the RTX our job is done.
690 return;
691 };
692 
693 to_change.maybe_reset_rtx(rtx_to);
694 }
695 
696 fn stream_rx_by_ssrc_or_rtx(&self, ssrc: Ssrc) -> Option<&StreamRx> {
697 self.streams_rx
698 .values()
699 .find(|s| s.ssrc() == ssrc || s.rtx() == Some(ssrc))
700 }
701 
702 pub(crate) fn any_nack_enabled(&mut self) -> bool {
703 if self.any_nack_active.is_none() {
704 self.any_nack_active = Some(self.streams_rx.values().any(|s| s.nack_enabled()));
705 }
706 self.any_nack_active.unwrap()
707 }
708 
709 fn rx_lookup_at(&self) -> Instant {
710 self.last_rx_lookup_cleanup + RX_LOOKUP_CLEANUP_INTERVAL
711 }
712}
713 
714impl fmt::Debug for RtpPacket {
715 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
716 f.debug_struct("RtpPacket")
717 .field("seq_no", &self.seq_no)
718 .field("time", &self.time)
719 .field("header", &self.header)
720 .field("payload", &self.payload.len())
721 .field("nackable", &self.nackable)
722 .field("timestamp", &self.timestamp)
723 .finish()
724 }
725}