Skip to content
File

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

rust912 lines
1use std::collections::VecDeque;
2use std::sync::Arc;
3use std::time::{Duration, Instant};
4 
5use crate::config_mod::RtcpReportIntervals;
6use crate::media::KeyframeRequestKind;
7use crate::rtp_::MidRid;
8use crate::rtp_::{Bitrate, DlrrItem, ExtendedReport, extend_u32};
9use crate::rtp_::{Fir, FirEntry, Frequency, MediaTime, Remb};
10use crate::rtp_::{Mid, Pli, Pt, ReceiverReport};
11use crate::rtp_::{ReportBlock, ReportList, Rid, Rrtr, Rtcp};
12use crate::rtp_::{RtcpFb, RtpHeader, SenderInfo, SeqNo};
13use crate::rtp_::{SdesType, Ssrc};
14use crate::stats::{MediaIngressStats, RemoteEgressStats, StatsSnapshot};
15use crate::util::{InstantExt, SystemTimeExt};
16use crate::util::{already_happened, calculate_rtt};
17 
18use super::RtpPacket;
19use super::StreamPaused;
20use super::register::ReceiverRegister;
21 
22/// Incoming encoded stream.
23///
24/// A stream is a primary SSRC + optional RTX SSRC.
25///
26/// This is RTP level API. For frame level API see [`Rtc::writer`][crate::Rtc::writer].
27#[derive(Debug)]
28pub struct StreamRx {
29 /// Unique idenfier of the remote stream.
30 ///
31 /// If the remote changes the SSRC, we will create a new stream, not change this id.
32 ssrc: Ssrc,
33 
34 /// Identifier of a resend (RTX) stream. This can be set later, once we discover it.
35 rtx: Option<Ssrc>,
36 
37 /// Previous main SSRC. This is to ensure we never go "backwards" in terms
38 /// of changing SSRC (for FF).
39 previous_ssrc: Option<Ssrc>,
40 
41 /// The Media mid/rid this stream belongs to.
42 midrid: MidRid,
43 
44 /// Incoming CNAME in Sdes reports.
45 cname: Option<String>,
46 
47 /// Whether we explicitly want to supress NACK sending. This is normally done by not
48 /// setting an RTX, however this can be toggled off manually despite RTX being there.
49 ///
50 /// This is also set to true if the SDP negotiation disables RTX.
51 ///
52 /// Defaults to false.
53 suppress_nack: bool,
54 
55 /// Timestamp when we got some indication of remote using this stream.
56 last_used: Instant,
57 
58 /// Last seen pt and clock_rate in
59 last_clock_rate: Option<(Pt, Frequency)>,
60 
61 /// Last received sender info.
62 sender_info: Option<LastSenderInfo>,
63 
64 /// ROC to reset with on next incoming packet.
65 reset_roc: Option<u64>,
66 
67 /// Register of received packets. For NACK handling.
68 ///
69 /// Set on first ever packet.
70 register: Option<ReceiverRegister>,
71 
72 /// Register of received packets for RTX.
73 ///
74 /// Set on first ever RTXpacket.
75 register_rtx: Option<ReceiverRegister>,
76 
77 /// Last observed media time in an RTP packet.
78 last_time: Option<MediaTime>,
79 
80 /// If we have a pending keyframe request to send.
81 pending_request_keyframe: Option<KeyframeRequestKind>,
82 
83 /// If we have a pending REMB request to send.
84 pending_request_remb: Option<Bitrate>,
85 
86 /// Sequence number of the next FIR.
87 fir_seq_no: u8,
88 
89 /// Last time we produced regular feedback RR.
90 last_receiver_report: Instant,
91 
92 /// Statistics of incoming data.
93 stats: StreamRxStats,
94 
95 /// When we need to evaluate the paused state.
96 ///
97 /// now + pause_threshold
98 check_paused_at: Option<Instant>,
99 
100 /// Whether we consider this StreamRx paused.
101 ///
102 /// A stream is considered paused if it has received no packets for some (configurable) duration.
103 /// This defaults to 1.5s.
104 paused: bool,
105 
106 /// Whether we need to emit a paused event for the current paused state.
107 need_paused_event: bool,
108 
109 /// The configured threshold before considering the lack of packets as going into paused.
110 pause_threshold: Duration,
111}
112 
113/// The last sender info we recieved.
114#[derive(Debug)]
115pub(crate) struct LastSenderInfo {
116 /// When this SenderInfo was received.
117 received_at: Instant,
118 /// The sender info itself.
119 info: SenderInfo,
120 /// Whether we have emitted it yet via `poll_sender_info`
121 emitted: bool,
122}
123 
124/// Holder of stats.
125#[derive(Debug, Default)]
126pub(crate) struct StreamRxStats {
127 /// count of bytes received, including retransmissions
128 bytes: u64,
129 /// count of packets received, including retransmissions
130 packets: u64,
131 /// count of FIR requests sent
132 firs: u64,
133 /// count of PLI requests sent
134 plis: u64,
135 /// count of NACKs sent
136 nacks: u64,
137 /// interarrival jitter (RTP timestamp units) from the last RR we sent
138 jitter: u32,
139 /// round trip time from the last DLRR, if any
140 rtt: Option<Duration>,
141 /// fraction of packets lost from the last RR, if any
142 loss: Option<f32>,
143}
144 
145impl StreamRx {
146 pub(crate) fn new(ssrc: Ssrc, midrid: MidRid, suppress_nack: bool) -> Self {
147 debug!("Create StreamRx for SSRC: {}", ssrc);
148 
149 StreamRx {
150 ssrc,
151 rtx: None,
152 previous_ssrc: None,
153 midrid,
154 cname: None,
155 suppress_nack,
156 last_used: already_happened(),
157 last_clock_rate: None,
158 sender_info: None,
159 reset_roc: None,
160 register: None,
161 register_rtx: None,
162 last_time: None,
163 pending_request_keyframe: None,
164 pending_request_remb: None,
165 fir_seq_no: 0,
166 last_receiver_report: already_happened(),
167 stats: StreamRxStats::default(),
168 check_paused_at: None,
169 paused: true,
170 need_paused_event: false,
171 pause_threshold: Duration::from_millis(1500),
172 }
173 }
174 
175 /// The (primary) SSRC of this encoded stream.
176 pub fn ssrc(&self) -> Ssrc {
177 self.ssrc
178 }
179 
180 /// The resend (RTX) SSRC of this encoded stream.
181 pub fn rtx(&self) -> Option<Ssrc> {
182 self.rtx
183 }
184 
185 /// Mid for this stream.
186 ///
187 /// In SDP this corresponds to m-line and "Media".
188 pub fn mid(&self) -> Mid {
189 self.midrid.mid()
190 }
191 
192 /// Rid for this stream.
193 ///
194 /// This is used to separate streams with the same [`Mid`] when using simulcast.
195 pub fn rid(&self) -> Option<Rid> {
196 self.midrid.rid()
197 }
198 
199 /// CNAME as sent by remote peer in a Sdes.
200 ///
201 /// The value is None until we receive a first report with the value set.
202 pub fn cname(&self) -> Option<&str> {
203 self.cname.as_deref()
204 }
205 
206 /// Set threshold duration for emitting the paused event.
207 ///
208 /// This event is emitted when no packet have received for this duration.
209 pub fn set_pause_threshold(&mut self, t: Duration) {
210 self.pause_threshold = t;
211 }
212 
213 /// The last time we received a packet.
214 ///
215 /// Resets if the SSRC changes.
216 pub fn last_time(&self) -> Option<MediaTime> {
217 self.last_time
218 }
219 
220 /// Request a keyframe for an incoming encoded stream.
221 ///
222 /// * SSRC the identifier of the remote encoded stream to request a keyframe for.
223 /// * kind PLI or FIR.
224 pub fn request_keyframe(&mut self, kind: KeyframeRequestKind) {
225 self.pending_request_keyframe = Some(kind);
226 }
227 
228 /// Request max recv bitrate for an incoming encoded stream.
229 ///
230 /// * bitrate Bitrate.
231 pub fn request_remb(&mut self, bitrate: Bitrate) {
232 self.pending_request_remb = Some(bitrate);
233 }
234 
235 /// Suppress NACK sending.
236 ///
237 /// Normally NACK is disabled by not having an RTX SSRC set. In some situations it might be
238 /// desirable to manually suppress NACK sending regardless of RTX setting.
239 pub fn suppress_nack(&mut self, suppress: bool) {
240 self.suppress_nack = suppress;
241 }
242 
243 fn is_audio(&self) -> bool {
244 self.rtx.is_none() // this is maybe not correct, but it's all we got.
245 }
246 
247 pub(crate) fn receiver_report_at(&self, intervals: RtcpReportIntervals) -> Instant {
248 self.last_receiver_report + intervals.for_audio(self.is_audio())
249 }
250 
251 pub(crate) fn handle_rtcp(&mut self, now: Instant, fb: RtcpFb) {
252 use RtcpFb::*;
253 match fb {
254 SenderInfo(v) => {
255 self.set_sender_info(now, v);
256 }
257 SourceDescription(v) => {
258 for (sdes, st) in v.values {
259 if sdes == SdesType::CNAME {
260 if st.is_empty() {
261 // In simulcast, chrome doesn't send the SSRC lines, but
262 // expects us to infer that from rtp headers. It does
263 // however send the SourceDescription RTCP with an empty
264 // string CNAME. ¯\_(ツ)_/¯
265 return;
266 }
267 
268 // Here we _could_ check CNAME here matches something. But
269 // CNAMEs are a bit unfashionable.
270 self.cname = Some(st);
271 return;
272 }
273 }
274 }
275 DlrrItem(v) => {
276 self.set_dlrr_item(now, v);
277 }
278 Goodbye(_v) => {
279 // We get Goodbye at weird times, like SDP renegotiation, which makes
280 // pausing on the BYE not a good idea. Chrome also reuses the SSRC it
281 // just sent BYE on. Very not helpful.
282 }
283 _ => {}
284 }
285 }
286 
287 fn set_sender_info(&mut self, now: Instant, mut info: SenderInfo) {
288 // Extend the incoming time given our knowledge of last time.
289 let extended = {
290 let prev = self
291 .sender_info
292 .as_ref()
293 .map(|last| last.info.rtp_time.numer());
294 let r_u32 = info.rtp_time.numer() as u32;
295 extend_u32(prev, r_u32)
296 };
297 
298 // The MediaTime has a base 1 after being parsed. At this point
299 // we know whether it's audio or video and set the base accordingly.
300 let clock_rate = self
301 .last_clock_rate
302 .map(|(_, r)| r)
303 .unwrap_or(Frequency::SECONDS);
304 
305 // Clock rate is that of the last received packet.
306 info.rtp_time = MediaTime::new(extended, clock_rate);
307 
308 self.sender_info = Some(LastSenderInfo {
309 received_at: now,
310 info,
311 emitted: false,
312 });
313 }
314 
315 fn set_dlrr_item(&mut self, now: Instant, dlrr: DlrrItem) {
316 let ntp_time = now.to_ntp_duration();
317 let rtt = calculate_rtt(ntp_time, dlrr.last_rr_delay, dlrr.last_rr_time);
318 self.stats.rtt = rtt;
319 }
320 
321 pub(crate) fn paused_at(&self) -> Option<Instant> {
322 self.check_paused_at
323 }
324 
325 pub(crate) fn handle_timeout(&mut self, now: Instant) {
326 // No scheduled paused check?
327 if self.check_paused_at.is_none() {
328 return;
329 }
330 
331 // Not reached scheduled paused check?
332 if Some(now) < self.check_paused_at {
333 return;
334 }
335 
336 // Every update() schedules a paused check in the future. If we have reached that
337 // future we have implicitly also paused.
338 self.check_paused_at = None;
339 
340 self.paused = true;
341 self.need_paused_event = true;
342 }
343 
344 pub(crate) fn extend_seq(
345 &mut self,
346 header: &RtpHeader,
347 is_repair: bool,
348 max_seq_lookup: impl Fn(Ssrc) -> Option<SeqNo>,
349 ) -> SeqNo {
350 // Select reference to register to use depending on RTX or not. The RTX has a separate
351 // sequence number series to the main register.
352 let register_ref = if is_repair {
353 &mut self.register_rtx
354 } else {
355 &mut self.register
356 };
357 
358 let register =
359 register_ref.get_or_insert_with(|| ReceiverRegister::new(max_seq_lookup(header.ssrc)));
360 
361 // If the user has called `reset_seq_no`, this is the time to handle it, but only
362 // if the incoming packet is for main (not repair).
363 let mut reset_seq_no = None;
364 if !is_repair {
365 if let Some(reset_roc) = self.reset_roc.take() {
366 let s: SeqNo = (reset_roc << 16 | header.sequence_number as u64).into();
367 reset_seq_no = Some(s);
368 }
369 }
370 
371 if let Some(reset_seq_no) = reset_seq_no {
372 reset_seq_no
373 } else {
374 header.sequence_number(register.max_seq())
375 }
376 }
377 
378 pub(crate) fn is_new_packet(&self, is_repair: bool, seq_no: SeqNo) -> bool {
379 let register_ref = if is_repair {
380 self.register_rtx.as_ref()
381 } else {
382 self.register.as_ref()
383 };
384 
385 // Unwrap is OK because we always call extend_seq() for the same is_repair flag beforehand
386 register_ref.unwrap().accepts(seq_no)
387 }
388 
389 pub(crate) fn update_register(
390 &mut self,
391 now: Instant,
392 header: &RtpHeader,
393 clock_rate: Frequency,
394 is_repair: bool,
395 seq_no: SeqNo,
396 ) -> RegisterUpdateReceipt {
397 self.last_used = now;
398 
399 let was_paused = self.paused;
400 if was_paused {
401 self.paused = false;
402 self.need_paused_event = true;
403 }
404 self.check_paused_at = Some(now + self.pause_threshold);
405 
406 let register_ref = if is_repair {
407 &mut self.register_rtx
408 } else {
409 &mut self.register
410 };
411 
412 // Unwrap is OK because we always call extend_seq() for the same is_repair flag beforehand
413 let register = register_ref.as_mut().unwrap();
414 
415 let is_new_packet = register.update(seq_no, now, header.timestamp, clock_rate.get());
416 
417 // Get the previous time for comparison
418 let previous_time = self.last_time.map(|t| t.numer());
419 
420 // Calculate the extended timestamp
421 let mut time_u32 = extend_u32(previous_time, header.timestamp);
422 
423 if was_paused && Some(time_u32) < previous_time {
424 // In 32-bit RTP timestamps, adding 2^31 (MAX/2) flips to the other half of timestamp space
425 // This forces extend_u32 to produce a value in the next cycle
426 const HALF_CYCLE: u32 = 1u32 << 31;
427 let adjusted_ts = header.timestamp.wrapping_add(HALF_CYCLE);
428 
429 // Recalculate extended timestamp with adjusted value
430 let adjusted_time_u32 = extend_u32(previous_time, adjusted_ts);
431 
432 // If this adjusted timestamp moves time forward, use it
433 if adjusted_time_u32 > previous_time.unwrap() {
434 time_u32 = adjusted_time_u32;
435 } else {
436 // Fallback
437 time_u32 = header.timestamp as u64;
438 }
439 }
440 
441 let time = MediaTime::new(time_u32, clock_rate);
442 
443 if !is_repair {
444 self.last_time = Some(time);
445 }
446 
447 RegisterUpdateReceipt {
448 time,
449 is_new_packet,
450 }
451 }
452 
453 #[allow(clippy::too_many_arguments)]
454 pub(crate) fn handle_rtp(
455 &mut self,
456 now: Instant,
457 header: RtpHeader,
458 payload: Arc<[u8]>,
459 seq_no: SeqNo,
460 time: MediaTime,
461 ) -> RtpPacket {
462 trace!("Handle RTP: {:?}", header);
463 
464 let need_clock_rate = self.last_clock_rate.map(|(pt, _)| pt) != Some(header.payload_type);
465 if need_clock_rate {
466 self.last_clock_rate = Some((header.payload_type, time.frequency()));
467 
468 // If we get an SR before the first packet, we update the potential clock rate.
469 if let Some(i) = &mut self.sender_info {
470 i.info.rtp_time = MediaTime::new(i.info.rtp_time.numer(), time.frequency());
471 }
472 }
473 
474 let packet = RtpPacket {
475 seq_no,
476 time,
477 header,
478 payload,
479 vp8_patch: None,
480 nackable: false,
481 last_sender_info: self.sender_info.as_ref().map(|l| l.info),
482 timestamp: now,
483 };
484 
485 self.stats.bytes += packet.payload.len() as u64;
486 self.stats.packets += 1;
487 
488 packet
489 }
490 
491 /// Rewrites the header to be the un-RTX'd original packet, and returns
492 /// `data` with the original sequence number prefix stripped.
493 pub(crate) fn un_rtx<'a>(&self, header: &mut RtpHeader, data: &'a [u8], pt: Pt) -> &'a [u8] {
494 let mut orig_seq_no_16 = 0;
495 
496 let n = RtpHeader::read_original_sequence_number(data, &mut orig_seq_no_16);
497 
498 trace!(
499 "Repaired seq no {} -> {}",
500 header.sequence_number, orig_seq_no_16
501 );
502 
503 header.sequence_number = orig_seq_no_16;
504 
505 header.ssrc = self.ssrc;
506 header.payload_type = pt;
507 header.ext_vals.rid = header.ext_vals.rid_repair.take();
508 
509 &data[n..]
510 }
511 
512 pub(crate) fn maybe_create_keyframe_request(
513 &mut self,
514 sender_ssrc: Ssrc,
515 feedback: &mut VecDeque<Rtcp>,
516 ) {
517 let Some(kind) = self.pending_request_keyframe.take() else {
518 return;
519 };
520 
521 let ssrc = self.ssrc;
522 
523 match kind {
524 KeyframeRequestKind::Pli => {
525 self.stats.plis += 1;
526 feedback.push_back(Rtcp::Pli(Pli { sender_ssrc, ssrc }))
527 }
528 KeyframeRequestKind::Fir => {
529 self.stats.firs += 1;
530 feedback.push_back(Rtcp::Fir(Fir {
531 sender_ssrc,
532 reports: FirEntry {
533 ssrc,
534 seq_no: self.next_fir_seq_no(),
535 }
536 .into(),
537 }))
538 }
539 }
540 }
541 
542 pub(crate) fn maybe_create_remb_request(
543 &mut self,
544 sender_ssrc: Ssrc,
545 feedback: &mut VecDeque<Rtcp>,
546 ) {
547 let Some(bitrate) = self.pending_request_remb.take() else {
548 return;
549 };
550 
551 feedback.push_back(Rtcp::Remb(Remb {
552 sender_ssrc,
553 ssrc: 0.into(),
554 bitrate: bitrate.as_f64() as f32,
555 ssrcs: vec![*self.ssrc],
556 }))
557 }
558 
559 fn next_fir_seq_no(&mut self) -> u8 {
560 let x = self.fir_seq_no;
561 self.fir_seq_no = self.fir_seq_no.wrapping_add(1);
562 x
563 }
564 
565 pub(crate) fn need_rr(&self, now: Instant, intervals: RtcpReportIntervals) -> bool {
566 if self.ssrc.is_probe() {
567 return false;
568 }
569 
570 now >= self.receiver_report_at(intervals)
571 }
572 
573 pub(crate) fn create_rr_and_update(
574 &mut self,
575 now: Instant,
576 sender_ssrc: Ssrc,
577 feedback: &mut VecDeque<Rtcp>,
578 ) {
579 let mut rr = self.create_receiver_report(now);
580 rr.sender_ssrc = sender_ssrc;
581 
582 if !rr.reports.is_empty() {
583 let report = &rr.reports[rr.reports.len() - 1];
584 self.stats.update_loss(report.fraction_lost);
585 self.stats.jitter = report.jitter;
586 }
587 
588 let xr = self.create_extended_receiver_report(now, sender_ssrc);
589 
590 trace!(
591 "Created feedback RR/XR ({:?}): {:?} {:?}",
592 self.midrid, rr, xr
593 );
594 feedback.push_back(Rtcp::ReceiverReport(rr));
595 feedback.push_back(Rtcp::ExtendedReport(xr));
596 
597 self.last_receiver_report = now;
598 }
599 
600 fn create_receiver_report(&mut self, now: Instant) -> ReceiverReport {
601 let Some(mut report) = self.register.as_mut().and_then(|r| r.reception_report()) else {
602 return ReceiverReport {
603 sender_ssrc: 0.into(), // set one level up
604 reports: ReportList::new(),
605 };
606 };
607 report.ssrc = self.ssrc;
608 
609 // The middle 32 bits out of 64 in the NTP timestamp (as explained in
610 // Section 4) received as part of the most recent RTCP sender report
611 // (SR) packet from source SSRC_n. If no SR has been received yet,
612 // the field is set to zero.
613 report.last_sr_time = {
614 let t64 = self
615 .sender_info
616 .as_ref()
617 .map_or(0u64, |l| l.info.ntp_time.as_ntp_64());
618 
619 (t64 >> 16) as u32
620 };
621 
622 // The delay, expressed in units of 1/65_536 seconds, between
623 // receiving the last SR packet from source SSRC_n and sending this
624 // reception report block. If no SR packet has been received yet
625 // from SSRC_n, the DLSR field is set to zero.
626 report.last_sr_delay = if let Some(l) = self.sender_info.as_ref() {
627 let t = l.received_at;
628 let delay = now - t;
629 ((delay.as_micros() * 65_536) / 1_000_000) as u32
630 } else {
631 0
632 };
633 
634 ReceiverReport {
635 sender_ssrc: 0.into(), // set one level up
636 reports: report.into(),
637 }
638 }
639 
640 fn create_extended_receiver_report(&self, now: Instant, sender_ssrc: Ssrc) -> ExtendedReport {
641 // we only want to report our time to measure RTT,
642 // the source will answer with Dlrr feedback, allowing us to calculate RTT
643 let block = ReportBlock::Rrtr(Rrtr {
644 ntp_time: now.to_system_time(),
645 });
646 ExtendedReport {
647 ssrc: sender_ssrc,
648 blocks: vec![block],
649 }
650 }
651 
652 pub(crate) fn nack_enabled(&self) -> bool {
653 // Deliberately don't look at RTX is_some() here, since when using dynamic SSRC, we might need
654 // to send NACK before discovering the remote RTX.
655 !self.suppress_nack
656 }
657 
658 pub(crate) fn maybe_create_nack(
659 &mut self,
660 sender_ssrc: Ssrc,
661 feedback: &mut VecDeque<Rtcp>,
662 ) -> Option<()> {
663 if !self.nack_enabled() {
664 return None;
665 }
666 
667 let nacks = self.register.as_mut().and_then(|r| r.nack_report())?;
668 
669 for mut nack in nacks {
670 nack.sender_ssrc = sender_ssrc;
671 nack.ssrc = self.ssrc;
672 
673 trace!("Created feedback NACK: {:?}", nack);
674 feedback.push_back(Rtcp::Nack(nack));
675 self.stats.nacks += 1;
676 }
677 
678 Some(())
679 }
680 
681 pub(crate) fn visit_stats(&self, snapshot: &mut StatsSnapshot, now: Instant) {
682 if self.ssrc.is_probe() {
683 return;
684 }
685 
686 self.stats
687 .fill(snapshot, self.midrid, self.sender_info.as_ref(), now);
688 }
689 
690 pub(crate) fn poll_paused(&mut self) -> Option<StreamPaused> {
691 if self.ssrc.is_probe() {
692 return None;
693 }
694 
695 if !self.need_paused_event {
696 return None;
697 }
698 
699 self.need_paused_event = false;
700 
701 debug!(
702 "{} StreamRx with {:?} and SSRC: {}",
703 if self.paused { "Paused" } else { "Unpaused" },
704 self.midrid,
705 self.ssrc
706 );
707 
708 Some(StreamPaused {
709 ssrc: self.ssrc,
710 mid: self.midrid.mid(),
711 rid: self.midrid.rid(),
712 paused: self.paused,
713 })
714 }
715 
716 /// Poll the most recent sender info and when it was received
717 pub(crate) fn poll_sender_info(&mut self) -> Option<(SenderInfo, Instant)> {
718 let i = self.sender_info.as_mut()?;
719 if i.emitted {
720 return None;
721 }
722 
723 i.emitted = true;
724 
725 Some((i.info, i.received_at))
726 }
727 
728 pub(crate) fn reset_buffers(&mut self, max_seq_lookup: impl Fn(Ssrc) -> Option<SeqNo>) {
729 if let Some(r) = &mut self.register {
730 r.clear(max_seq_lookup(self.ssrc));
731 }
732 
733 if let Some(r) = &mut self.register_rtx {
734 r.clear(self.rtx.and_then(max_seq_lookup));
735 }
736 self.pending_request_keyframe = None;
737 }
738 
739 #[must_use]
740 pub(crate) fn change_ssrc(&mut self, ssrc: Ssrc) -> bool {
741 // Avoid flapping
742 if ssrc == self.ssrc || Some(ssrc) == self.previous_ssrc {
743 return false;
744 }
745 
746 debug!(
747 "Change main SSRC: {} -> {} {:?}",
748 self.ssrc, ssrc, self.midrid
749 );
750 
751 // Remember which was the previous in case a stray packet turns up
752 // so do we don't go "backwards".
753 self.previous_ssrc = Some(self.ssrc);
754 self.ssrc = ssrc;
755 
756 // Reset all SSRC-specific state
757 self.register = None;
758 self.last_time = None;
759 self.last_clock_rate = None;
760 self.sender_info = None;
761 self.last_receiver_report = already_happened();
762 self.fir_seq_no = 0;
763 self.pending_request_keyframe = None;
764 self.pending_request_remb = None;
765 self.reset_roc = None;
766 
767 // Note: We don't reset the RTX register here, as the RTX SSRC is managed separately
768 // via maybe_reset_rtx() and is not directly tied to the main SSRC change.
769 
770 true
771 }
772 
773 pub(crate) fn maybe_reset_rtx(&mut self, rtx: Ssrc) {
774 if let Some(current) = self.rtx {
775 if current == rtx {
776 return;
777 }
778 
779 debug!(
780 "Change RTX SSRC {} -> {} for main SSRC: {} {:?}",
781 current, rtx, self.ssrc, self.midrid
782 );
783 } else {
784 debug!("SSRC {} associated with RTX: {}", self.ssrc, rtx);
785 }
786 
787 self.rtx = Some(rtx);
788 self.register_rtx = None;
789 }
790 
791 /// Reset the current rollover counter (ROC).
792 ///
793 /// This is used in scenarios where we use a single sequence number across all
794 /// receivers of the same stream (as opposed to a sequence number unique per peer).
795 ///
796 /// [RFC3711](https://datatracker.ietf.org/doc/html/rfc3711#section-3.3.1):
797 ///
798 /// > Receivers joining an on-going session MUST be given the
799 /// > current ROC value using out-of-band signaling such as key-management
800 /// > signaling. Furthermore, the receiver SHALL initialize s_l to the RTP
801 /// > sequence number (SEQ) of the first observed SRTP packet (unless the
802 /// > initial value is provided by out of band signaling such as key
803 /// > management).
804 pub fn reset_roc(&mut self, roc: u64) {
805 self.register = None;
806 self.register_rtx = None;
807 self.reset_roc = Some(roc);
808 }
809 
810 pub(crate) fn is_midrid(&self, midrid: MidRid) -> bool {
811 midrid.special_equals(&self.midrid)
812 }
813}
814 
815impl StreamRxStats {
816 fn update_loss(&mut self, fraction_lost: u8) {
817 self.loss = Some(fraction_lost as f32 / u8::MAX as f32)
818 }
819 
820 pub(crate) fn fill(
821 &self,
822 snapshot: &mut StatsSnapshot,
823 midrid: MidRid,
824 sender_info: Option<&LastSenderInfo>,
825 now: Instant,
826 ) {
827 if self.bytes == 0 {
828 return;
829 }
830 
831 let stats = MediaIngressStats {
832 mid: midrid.mid(),
833 rid: midrid.rid(),
834 bytes: self.bytes,
835 packets: self.packets,
836 firs: self.firs,
837 plis: self.plis,
838 nacks: self.nacks,
839 jitter: self.jitter,
840 rtt: self.rtt,
841 loss: self.loss,
842 timestamp: now,
843 remote: sender_info.map(|l| RemoteEgressStats {
844 bytes: l.info.sender_octet_count as u64,
845 packets: l.info.sender_packet_count as u64,
846 }),
847 };
848 
849 // Several SSRCs can back a given (mid, rid) tuple. For example, Firefox creates new SSRCs
850 // when a Transceiver transitions from send -> inactive -> send. In order to continue
851 // correctly reporting stats for this (mid, rid) pair we need to merge the stats across all
852 // the SSRCs that have been used.
853 snapshot
854 .ingress
855 .entry(midrid)
856 .and_modify(|s| s.merge_by_mid_rid(&stats))
857 .or_insert(stats);
858 }
859}
860 
861#[derive(Debug, Clone, Copy)]
862pub(crate) struct RegisterUpdateReceipt {
863 pub time: MediaTime,
864 pub is_new_packet: bool,
865}
866 
867#[cfg(test)]
868mod tests {
869 use super::*;
870 
871 #[test]
872 fn paused_timestamp_repair_moves_time_forward() {
873 let now = already_happened();
874 let mut stream = StreamRx::new(7.into(), MidRid("mid".into(), None), false);
875 let previous_time = 1;
876 stream.last_time = Some(MediaTime::new(previous_time, Frequency::NINETY_KHZ));
877 stream.paused = true;
878 
879 let header = RtpHeader {
880 payload_type: Pt::new_with_value(96),
881 sequence_number: 42,
882 timestamp: 0,
883 ssrc: 7.into(),
884 ..Default::default()
885 };
886 
887 let seq_no = stream.extend_seq(&header, false, |_| None);
888 let receipt = stream.update_register(now, &header, Frequency::NINETY_KHZ, false, seq_no);
889 
890 assert!(
891 receipt.time.numer() > previous_time,
892 "expected repaired media time to move forward"
893 );
894 assert!(!stream.paused);
895 }
896 
897 #[test]
898 fn receiver_report_uses_supplied_intervals() {
899 let mut stream = StreamRx::new(7.into(), MidRid("mid".into(), None), false);
900 let now = Instant::now();
901 stream.last_receiver_report = now;
902 
903 assert_eq!(
904 stream.receiver_report_at(RtcpReportIntervals {
905 audio: Duration::from_millis(750),
906 video: Duration::from_millis(500),
907 }),
908 now + Duration::from_millis(750)
909 );
910 }
911}