Skip to content
File

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

rust1331 lines
1use std::collections::VecDeque;
2use std::sync::Arc;
3use std::time::Duration;
4use std::time::Instant;
5 
6use crate::config_mod::RtcpReportIntervals;
7use crate::format::CodecConfig;
8use crate::format::PayloadParams;
9use crate::io::DATAGRAM_MAX_PACKET_SIZE;
10use crate::io::MAX_RTP_OVERHEAD;
11use crate::media::KeyframeRequestKind;
12use crate::media::Media;
13use crate::media::MediaKind;
14use crate::pacer::QueuePriority;
15use crate::pacer::QueueSnapshot;
16use crate::pacer::QueueState;
17use crate::packet::Vp8Patch;
18use crate::rtp_::MidRid;
19use crate::rtp_::{Bitrate, Descriptions, Extension, ExtensionMap, ExtensionValues, Frequency};
20use crate::rtp_::{MAX_BLANK_PADDING_PAYLOAD_SIZE, Sdes, SdesType};
21use crate::rtp_::{MediaTime, Mid, NackEntry, ReportList, Rtcp, RtpHeader};
22use crate::rtp_::{Pt, Rid, RtcpFb, SenderInfo, SenderReport, Ssrc};
23use crate::rtp_::{SRTP_BLOCK_SIZE, SeqNo};
24use crate::session::PacketReceipt;
25use crate::stats::StatsSnapshot;
26use crate::util::value_history::ValueHistory;
27use crate::util::{InstantExt, already_happened, not_happening};
28 
29use super::RtpPacket;
30use super::rtx_cache::RtxCache;
31use super::send_queue::SendQueue;
32use super::send_stats::StreamTxStats;
33 
34/// The smallest size of padding for which we attempt to use a spurious resend. For padding
35/// requests smaller than this we use blank packets instead.
36const MIN_SPURIOUS_PADDING_SIZE: usize = 50;
37 
38pub const DEFAULT_RTX_CACHE_DURATION: Duration = Duration::from_secs(3);
39 
40pub const DEFAULT_RTX_RATIO_CAP: Option<f32> = Some(0.15f32);
41 
42/// A recently sampled view of an outgoing stream's send queue.
43///
44/// str0m refreshes this information as part of its regular pacer and bandwidth-estimation
45/// processing. Reading it does not cause a new value to be computed.
46pub trait StreamTxQueueInfo {
47 /// When this queue information was sampled.
48 ///
49 /// Compare this value between reads to determine whether str0m has computed a new sample.
50 fn created_at(&self) -> Instant;
51 
52 /// Total number of bytes in the queue.
53 fn byte_size(&self) -> usize;
54 
55 /// Total number of packets in the queue.
56 fn packet_count(&self) -> usize;
57 
58 /// When the first packet still in the queue was enqueued.
59 fn first_unsent(&self) -> Option<Instant>;
60}
61 
62impl StreamTxQueueInfo for QueueSnapshot {
63 fn created_at(&self) -> Instant {
64 self.created_at
65 }
66 
67 fn byte_size(&self) -> usize {
68 self.byte_size
69 }
70 
71 fn packet_count(&self) -> usize {
72 self.packet_count as usize
73 }
74 
75 fn first_unsent(&self) -> Option<Instant> {
76 self.first_unsent
77 }
78}
79 
80/// Outgoing encoded stream.
81///
82/// A stream is a primary SSRC + optional RTX SSRC.
83///
84/// This is RTP level API. For frame level API see [`Rtc::writer`][crate::Rtc::writer].
85#[derive(Debug)]
86pub struct StreamTx {
87 /// Unique identifier of the remote encoded stream.
88 ssrc: Ssrc,
89 
90 /// Identifier of a resend (RTX) stream. If we are doing resends.
91 rtx: Option<Ssrc>,
92 
93 /// The Media mid and rid this stream belongs to.
94 midrid: MidRid,
95 
96 /// Set on first handle_timeout.
97 kind: Option<MediaKind>,
98 
99 /// Set on first handle_timeout.
100 cname: Option<String>,
101 
102 /// The last main payload clock rate that was sent.
103 clock_rate: Option<Frequency>,
104 
105 /// If we are doing seq_no ourselves (when writing frame mode).
106 seq_no: SeqNo,
107 
108 /// If we are using RTX, this is the seq no counter.
109 seq_no_rtx: SeqNo,
110 
111 /// The last seq_no that we sent, either by increasing seq_no ourselves (media API), or by
112 /// direct RTP mode writing.
113 last_sent_seq_no: SeqNo,
114 
115 /// When we last sent something for this encoded stream, packet or RTCP.
116 last_used: Instant,
117 
118 /// Last written media + wallclock time.
119 rtp_and_wallclock: Option<(u32, Instant)>,
120 
121 /// Queue of packets to send.
122 ///
123 /// The packets here do not have correct sequence numbers, header extension values etc.
124 /// They must be updated when we are about to send.
125 send_queue: SendQueue,
126 
127 /// The queue information sampled during the latest pacer update.
128 queue_info: Option<QueueSnapshot>,
129 
130 /// Whether this sender is to be unpaced in BWE situations.
131 ///
132 /// Audio defaults to not being paced.
133 unpaced: Option<bool>,
134 
135 /// Scheduled resends due to NACK or spurious padding.
136 resends: VecDeque<Resend>,
137 
138 /// Requested padding, that has not been turned into packets yet.
139 padding: usize,
140 
141 /// Dummy packet for resends. Used between poll_packet and poll_packet_padding
142 blank_packet: RtpPacket,
143 
144 /// Cache of sent packets to be able to answer to NACKs as well as
145 /// sending spurious resends as padding.
146 rtx_cache: RtxCache,
147 
148 /// Determines retransmitted bytes ratio value to clear queued resends.
149 rtx_ratio_cap: Option<f32>,
150 
151 /// Last time we produced a SR.
152 last_sender_report: Instant,
153 
154 /// If we have a pending incoming keyframe request.
155 pending_request_keyframe: Option<KeyframeRequestKind>,
156 
157 /// If we have a pending incoming remb request.
158 pending_request_remb: Option<Bitrate>,
159 
160 /// Statistics of outgoing data.
161 ///
162 /// Stats are use to calculate the rtx ratio also when statistics events are disabled.
163 stats: StreamTxStats,
164 
165 // downsampled rtx ratio (value, last calculation)
166 rtx_ratio: (f32, Instant),
167 
168 // The _main_ PT to use for padding. This is main PT, since the poll_packet() loop
169 // figures out the param.resend() RTX PT using main.
170 pt_for_padding: Option<Pt>,
171 
172 /// Whether a receiver report has been received for this SSRC, thus acknowledging
173 /// that the receiver has bound the Mid/Rid tuple to the SSRC and no longer
174 /// needs to be sent on every packet
175 remote_acked_ssrc: bool,
176 
177 // Same as `remote_acked_ssrc`, but for the RTX SSRC
178 remote_acked_rtx_ssrc: bool,
179 
180 /// MTU warn threshold; used to cap spurious-padding RTX cache lookups.
181 mtu_warn: usize,
182}
183 
184/// RTP packet data to enqueue on a direct RTP send stream.
185///
186/// The payload is the RTP payload only, without the RTP header.
187///
188/// Optional RTP fields default to the common unmarked non-nackable packet with
189/// no header extension values, no CSRC entries and no VP8 rewrite.
190#[derive(Debug)]
191pub struct RtpWrite {
192 pt: Pt,
193 seq_no: SeqNo,
194 time: u32,
195 wallclock: Instant,
196 marker: bool,
197 ext_vals: ExtensionValues,
198 nackable: bool,
199 payload: Arc<[u8]>,
200 csrc_count: usize,
201 csrc: [u32; 15],
202 vp8_patch: Option<Vp8Patch>,
203}
204 
205impl RtpWrite {
206 /// Create a direct RTP write command.
207 ///
208 /// The `payload` argument is expected to be only the RTP payload, not the
209 /// RTP packet header.
210 ///
211 /// * `pt` Payload type. Declared in the [`Media`][crate::media::Media] this
212 /// encoded stream belongs to.
213 /// * `seq_no` Sequence number to use for this packet.
214 /// * `time` Time in whatever the clock rate is for the media in question
215 /// (normally 90_000 for video and 48_000 for audio).
216 /// * `wallclock` Real world time that corresponds to the media time in the
217 /// RTP packet. For an SFU, this can be hard to know because
218 /// RTP packets typically only contain the media time (RTP
219 /// time). In the simplest SFU setup, the wallclock could
220 /// simply be the arrival time of the incoming RTP data. For
221 /// better synchronization the SFU probably needs to weigh in
222 /// clock drifts and data provided via the statistics, receiver
223 /// reports etc.
224 /// * `payload` RTP packet payload, without header.
225 ///
226 /// Optional fields default to an unmarked packet with no header extension
227 /// values, no CSRC entries, no VP8 patch and no NACK support.
228 /// [`RtpWrite::marker`] for video frame boundaries,
229 /// [`RtpWrite::ext_vals`] for RTP header extension values,
230 /// [`RtpWrite::nackable`] for packets that may answer incoming NACKs,
231 /// [`RtpWrite::csrc`] for contributing source identifiers,
232 /// [`RtpWrite::vp8_patch`] to apply a prevalidated VP8 payload descriptor patch during RTP serialization.
233 pub fn new(
234 pt: Pt,
235 seq_no: SeqNo,
236 time: u32,
237 wallclock: Instant,
238 payload: impl Into<Arc<[u8]>>,
239 ) -> Self {
240 Self {
241 pt,
242 seq_no,
243 time,
244 wallclock,
245 marker: false,
246 ext_vals: ExtensionValues::default(),
247 nackable: false,
248 payload: payload.into(),
249 csrc_count: 0,
250 csrc: [0; 15],
251 vp8_patch: None,
252 }
253 }
254 
255 /// Set the RTP marker bit.
256 pub fn marker(mut self, marker: bool) -> Self {
257 self.marker = marker;
258 self
259 }
260 
261 /// Set RTP header extension values.
262 pub fn ext_vals(mut self, ext_vals: ExtensionValues) -> Self {
263 self.ext_vals = ext_vals;
264 self
265 }
266 
267 /// Set whether this RTP packet should respond to incoming NACKs.
268 pub fn nackable(mut self, nackable: bool) -> Self {
269 self.nackable = nackable;
270 self
271 }
272 
273 /// Set contributing source identifiers.
274 ///
275 /// # Panics
276 ///
277 /// Panics if more than 15 entries are supplied.
278 pub fn csrc(mut self, csrc: &[u32]) -> Self {
279 assert!(csrc.len() <= 15, "CSRC count must be <= 15");
280 
281 self.csrc = [0; 15];
282 self.csrc[..csrc.len()].copy_from_slice(csrc);
283 self.csrc_count = csrc.len();
284 self
285 }
286 
287 /// Apply a prevalidated VP8 payload descriptor patch during RTP serialization.
288 ///
289 /// Build the patch with [`Vp8Descriptor::patch`][crate::rtp::Vp8Descriptor::patch]
290 /// when the caller already parsed the VP8 payload descriptor before
291 /// constructing this write command.
292 pub fn vp8_patch(mut self, patch: Vp8Patch) -> Self {
293 self.vp8_patch = Some(patch);
294 self
295 }
296}
297 
298impl StreamTx {
299 pub(crate) fn new(
300 ssrc: Ssrc,
301 rtx: Option<Ssrc>,
302 midrid: MidRid,
303 enable_stats: bool,
304 mtu_warn: usize,
305 ) -> Self {
306 debug!("Create StreamTx for SSRC: {}", ssrc);
307 
308 StreamTx {
309 ssrc,
310 rtx,
311 midrid,
312 kind: None,
313 cname: None,
314 clock_rate: None,
315 seq_no: SeqNo::default(),
316 seq_no_rtx: SeqNo::default(),
317 last_sent_seq_no: SeqNo::default(),
318 last_used: already_happened(),
319 rtp_and_wallclock: None,
320 send_queue: SendQueue::new(),
321 queue_info: None,
322 unpaced: None,
323 resends: VecDeque::new(),
324 padding: 0,
325 blank_packet: RtpPacket::blank(),
326 rtx_cache: RtxCache::new(2000, DEFAULT_RTX_CACHE_DURATION),
327 rtx_ratio_cap: DEFAULT_RTX_RATIO_CAP,
328 last_sender_report: already_happened(),
329 pending_request_keyframe: None,
330 pending_request_remb: None,
331 stats: StreamTxStats::new(enable_stats),
332 rtx_ratio: (0.0, already_happened()),
333 pt_for_padding: None,
334 remote_acked_ssrc: false,
335 remote_acked_rtx_ssrc: false,
336 mtu_warn,
337 }
338 }
339 
340 /// The (primary) SSRC of this encoded stream.
341 pub fn ssrc(&self) -> Ssrc {
342 self.ssrc
343 }
344 
345 /// The resend (RTX) SSRC of this encoded stream.
346 pub fn rtx(&self) -> Option<Ssrc> {
347 self.rtx
348 }
349 
350 /// Mid for this stream.
351 ///
352 /// In SDP this corresponds to m-line and "Media".
353 pub fn mid(&self) -> Mid {
354 self.midrid.mid()
355 }
356 
357 /// Rid for this stream.
358 ///
359 /// This is used to separate streams with the same [`Mid`] when using simulcast.
360 pub fn rid(&self) -> Option<Rid> {
361 self.midrid.rid()
362 }
363 
364 /// Information sampled during the latest pacer update for this stream's send queue.
365 ///
366 /// This is `None` until the pacer has observed the stream for the first time.
367 /// Calling this method does not refresh the information; use
368 /// [`StreamTxQueueInfo::created_at()`] to determine whether a new sample is available.
369 pub fn queue_info(&self) -> Option<&dyn StreamTxQueueInfo> {
370 self.queue_info
371 .as_ref()
372 .map(|info| info as &dyn StreamTxQueueInfo)
373 }
374 
375 /// Configure the RTX (resend) cache.
376 ///
377 /// This determines how old incoming NACKs we can reply to.
378 ///
379 /// `rtx_ratio_cap` determines when to clear queued resends because of too many resends,
380 /// i.e. if `tx_sum / (rtx_sum + tx_sum) > rtx_ratio_cap`. `None` disables this functionality
381 /// so all queued resends will be sent.
382 ///
383 /// The default is 1024 packets over 3 seconds and RTX cache drop ratio of 0.15.
384 pub fn set_rtx_cache(
385 &mut self,
386 max_packets: usize,
387 max_age: Duration,
388 rtx_ratio_cap: Option<f32>,
389 ) {
390 // Dump old cache to avoid having to deal with resizing logic inside the cache impl.
391 self.rtx_cache = RtxCache::new(max_packets, max_age);
392 if rtx_ratio_cap.is_some() {
393 self.stats
394 .bytes_transmitted
395 .get_or_insert_with(ValueHistory::default);
396 self.stats
397 .bytes_retransmitted
398 .get_or_insert_with(ValueHistory::default);
399 } else {
400 self.stats.bytes_transmitted = None;
401 self.stats.bytes_retransmitted = None;
402 }
403 self.rtx_ratio_cap = rtx_ratio_cap;
404 }
405 
406 /// Set whether this stream is unpaced or not.
407 ///
408 /// This is only relevant when BWE (Bandwidth Estimation) is enabled. By default, audio is unpaced
409 /// thus not held to a steady send rate by the Pacer.
410 ///
411 /// This overrides the default behavior.
412 pub fn set_unpaced(&mut self, unpaced: bool) {
413 self.unpaced = Some(unpaced);
414 }
415 
416 /// Write RTP packet to a send stream.
417 ///
418 /// The write command carries the RTP payload and optional RTP metadata.
419 pub fn write_rtp(&mut self, rtp: RtpWrite) {
420 let RtpWrite {
421 pt,
422 seq_no,
423 time,
424 wallclock,
425 marker,
426 ext_vals,
427 nackable,
428 payload,
429 csrc_count,
430 csrc,
431 vp8_patch,
432 } = rtp;
433 
434 let first_call = self.rtp_and_wallclock.is_none();
435 
436 if first_call && seq_no.roc() > 0 {
437 // TODO: make it possible to supress this.
438 warn!(
439 "First SeqNo has non-zero ROC ({}), which needs out-of-band signalling \
440 to remote peer",
441 seq_no.roc()
442 );
443 }
444 
445 // This 1 in clock frequency will be fixed in poll_output.
446 let media_time = MediaTime::from_secs(time as u64);
447 self.rtp_and_wallclock = Some((time, wallclock));
448 
449 let header = RtpHeader {
450 csrc_count,
451 sequence_number: *seq_no as u16,
452 marker,
453 payload_type: pt,
454 timestamp: time,
455 ssrc: self.ssrc,
456 csrc,
457 ext_vals,
458 ..Default::default()
459 };
460 
461 let packet = RtpPacket {
462 seq_no,
463 time: media_time,
464 header,
465 payload,
466 vp8_patch,
467 nackable,
468 // The overall idea for str0m is to only drive time forward from handle_input. If we
469 // used a "now" argument to write_rtp(), we effectively get a second point that also need
470 // to move time forward _for all of Rtc_ โ€“ that's too complicated.
471 //
472 // Instead we set a future timestamp here. When time moves forward in the "regular way",
473 // in handle_timeout() we delegate to self.send_queue.handle_timeout() to mark the enqueued
474 // timestamp of all packets that are about to be sent.
475 timestamp: not_happening(),
476 
477 // This is only relevant for incoming RTP packets.
478 last_sender_info: None,
479 };
480 
481 self.send_queue.push(packet);
482 }
483 
484 fn padding_enabled(&self) -> bool {
485 self.rtx.is_some() && self.pt_for_padding.is_some()
486 }
487 
488 pub(crate) fn poll_packet(
489 &mut self,
490 now: Instant,
491 exts: &ExtensionMap,
492 twcc: Option<&mut u64>,
493 params: &[PayloadParams],
494 buf: &mut Vec<u8>,
495 ) -> Option<PacketReceipt> {
496 let mid = self.midrid.mid();
497 let rid = self.midrid.rid();
498 let ssrc_rtx = self.rtx;
499 let remote_acked_ssrc = self.remote_acked_ssrc;
500 let remote_acked_rtx_ssrc = self.remote_acked_rtx_ssrc;
501 
502 let (next, is_padding) = if let Some(next) = self.poll_packet_resend(now) {
503 (next, false)
504 } else if let Some(next) = self.poll_packet_regular(now) {
505 (next, false)
506 } else {
507 let next = self.poll_packet_padding(now)?;
508 (next, true)
509 };
510 
511 let pop_send_queue = next.kind == NextPacketKind::Regular;
512 
513 // Need the header for the receipt and modifications
514 // TODO: Can we remove this?
515 let header_ref = &mut next.pkt.header;
516 
517 // <https://webrtc.googlesource.com/src/+/refs/heads/main/modules/rtp_rtcp/source/rtp_sender.cc#537>
518 // BUNDLE requires that the receiver "bind" the received SSRC to the values
519 // in the MID and/or (R)RID header extensions if present. Therefore, the
520 // sender can reduce overhead by omitting these header extensions once it
521 // knows that the receiver has "bound" the SSRC.
522 // <snip>
523 // The algorithm here is fairly simple: Always attach a MID and/or RID (if
524 // configured) to the outgoing packets until an RTCP receiver report comes
525 // back for this SSRC. That feedback indicates the receiver must have
526 // received a packet with the SSRC and header extension(s), so the sender
527 // then stops attaching the MID and RID.
528 if !remote_acked_ssrc {
529 header_ref.ext_vals.mid = Some(mid);
530 header_ref.ext_vals.rid = rid;
531 }
532 
533 let pt_main = header_ref.payload_type;
534 
535 // The pt in next.pkt is the "main" pt.
536 let Some(param) = params.iter().find(|p| p.pt() == pt_main) else {
537 // PT does not exist in the connected media.
538 warn!("Media is missing PT ({}) used in RTP packet", pt_main);
539 
540 // Get rid of this packet we can't send.
541 if pop_send_queue {
542 self.send_queue.pop(now);
543 }
544 
545 return None;
546 };
547 
548 let mut set_pt_for_padding = None;
549 let mut set_cr = None;
550 
551 let mut header = match next.kind {
552 NextPacketKind::Regular => {
553 let rtx_possible = param.resend().is_some();
554 
555 if rtx_possible {
556 // Remember PT We want to set these directly on `self` here, but can't
557 // because we already have a mutable borrow. We are using pt_main
558 // since the above loop figuring out param needs to be correct also
559 // for the NextPacketKind::Blank case.
560 set_pt_for_padding = Some(pt_main);
561 } else {
562 // If the PT we're sending on doesn't have a corresponding RTX PT,
563 // the packet is de-facto not nackable.
564 //
565 // This blocks incoming NACK requests and thus ensures there are no
566 // entries in self.retries without a RTX PT.
567 next.pkt.nackable = false;
568 }
569 
570 let clock_rate = param.spec().rtp_clock_rate();
571 set_cr = Some(clock_rate);
572 
573 // Modify the cached packet time. This is so write_rtp can use u32 media time without
574 // worrying about lengthening or the clock rate.
575 let time = MediaTime::new(next.pkt.time.numer(), clock_rate);
576 next.pkt.time = time;
577 
578 // Modify the original (and also cached) header value.
579 header_ref.ext_vals.rid_repair = None;
580 
581 header_ref.clone()
582 }
583 NextPacketKind::Resend(_) | NextPacketKind::Blank(_) => {
584 // * For the Resend case, we will not have accepted/cached the packet unless
585 // we have a RTX PT (see logic setting next.pkt.nackable above).
586 // * For the Blank case, we will only have produced blank packets if we
587 // got a "real" PTX RT, either via set_pt_for_padding above, or via
588 // the on_first_timeout() further down.
589 // Either way, unwrapping this optional _should_ be correct.
590 let pt_rtx = param.resend().expect("PT for resend or blank");
591 
592 // Clone header to not change the original (cached) header.
593 let mut header = header_ref.clone();
594 
595 // Update clone of header (to not change the cached value).
596 header.payload_type = pt_rtx;
597 header.ssrc = ssrc_rtx.expect("Should have RTX SSRC for resends");
598 header.sequence_number = *next.seq_no as u16;
599 
600 header.ext_vals.rid = None;
601 header.ext_vals.rid_repair = rid;
602 
603 if !remote_acked_rtx_ssrc {
604 header.ext_vals.mid = Some(mid);
605 }
606 
607 header
608 }
609 };
610 
611 // These need to match `Extension::is_supported()` so we are sending what we are
612 // declaring we support.
613 
614 // Absolute Send Time might not be enabled for this m-line.
615 if exts.id_of(Extension::AbsoluteSendTime).is_some() {
616 header.ext_vals.abs_send_time = Some(now);
617 }
618 
619 // TWCC might not be enabled for this m-line.
620 if let Some(twcc) = twcc {
621 header.ext_vals.transport_cc = Some(*twcc as u16);
622 *twcc += 1;
623 }
624 
625 buf.resize(DATAGRAM_MAX_PACKET_SIZE, 0);
626 
627 let header_len = header.write_to(buf, exts);
628 assert!(header_len % 4 == 0, "RTP header must be multiple of 4");
629 header.header_len = header_len;
630 
631 let mut body_out = &mut buf[header_len..];
632 
633 // For resends, the original seq_no is inserted before the payload.
634 let mut original_seq_len = 0;
635 if let NextPacketKind::Resend(orig_seq_no) = next.kind {
636 original_seq_len = RtpHeader::write_original_sequence_number(body_out, orig_seq_no);
637 body_out = &mut body_out[original_seq_len..];
638 }
639 
640 let pkt = &next.pkt;
641 
642 let body_len = match next.kind {
643 NextPacketKind::Regular | NextPacketKind::Resend(_) => {
644 let body_len = pkt.payload.len();
645 if let Some(patch) = pkt.vp8_patch.as_ref() {
646 patch.copy_to(pkt.payload.as_ref(), &mut body_out[..body_len]);
647 } else {
648 body_out[..body_len].copy_from_slice(pkt.payload.as_ref());
649 }
650 
651 // media packets are sent unpadded: the SRTP ciphers handle
652 // arbitrary body lengths (the CTR scratch alignment is internal
653 // to protect_rtp), and RTP-level padding breaks SFUs that
654 // re-marshal forwarded packets without it while keeping the P
655 // bit (observed with LiveKit -> Chrome, str0m issue #1014)
656 body_len + original_seq_len
657 }
658 NextPacketKind::Blank(len) => {
659 let len = RtpHeader::create_padding_packet(
660 &mut buf[..],
661 header_len,
662 len,
663 SRTP_BLOCK_SIZE,
664 );
665 
666 if len == 0 {
667 return None;
668 }
669 
670 len
671 }
672 };
673 
674 buf.truncate(header_len + body_len);
675 
676 #[cfg(feature = "_internal_dont_use_log_stats")]
677 {
678 let queued_at = match next.kind {
679 NextPacketKind::Regular => Some(pkt.timestamp),
680 _ => {
681 // TODO: We don't have queued at stats for Resends or blank padding.
682 None
683 }
684 };
685 
686 if let Some(delay) = queued_at.map(|i| now.duration_since(i)) {
687 crate::log_stat!("QUEUE_DELAY", header.ssrc, delay.as_secs_f64() * 1000.0);
688 }
689 }
690 
691 let seq_no = next.seq_no;
692 if next.kind == NextPacketKind::Regular {
693 self.last_sent_seq_no = seq_no;
694 }
695 
696 self.last_used = now;
697 
698 // Padding comes in two forms, "spurious resends" of sent packets where
699 // the remote side didn't ask for a resend. The other variant are blank
700 // packets, containing nothing but zeroes. Such packets must be sent from
701 // _some_ RTX PT. A good pick is the RTX for the PT last used to send
702 // regular media data.
703 //
704 // This is set here due to borrow checker.
705 if set_pt_for_padding.is_some() && self.pt_for_padding != set_pt_for_padding {
706 self.pt_for_padding = set_pt_for_padding;
707 }
708 
709 if set_cr.is_some() && self.clock_rate != set_cr {
710 self.clock_rate = set_cr;
711 }
712 
713 if pop_send_queue {
714 // poll_packet_regular leaves the packet in the head of the send_queue
715 let pkt = self
716 .send_queue
717 .pop(now)
718 .expect("head of send_queue to be there");
719 if pkt.nackable {
720 self.rtx_cache.cache_sent_packet(pkt, now);
721 }
722 }
723 
724 Some(PacketReceipt {
725 header,
726 seq_no,
727 is_padding,
728 payload_size: body_len,
729 })
730 }
731 
732 fn rtx_ratio_downsampled(&mut self, now: Instant) -> f32 {
733 assert!(
734 self.stats.bytes_transmitted.is_some(),
735 "rtx_ratio_cap must be enabled"
736 );
737 assert!(
738 self.stats.bytes_retransmitted.is_some(),
739 "rtx_ratio_cap must be enabled"
740 );
741 
742 let (value, ts) = self.rtx_ratio;
743 if now - ts < Duration::from_millis(50) {
744 // not worth re-evaluating, return the old value
745 return value;
746 }
747 
748 // bytes stats refer to the last second by default
749 self.stats
750 .bytes_transmitted
751 .as_mut()
752 .unwrap()
753 .purge_old(now);
754 self.stats
755 .bytes_retransmitted
756 .as_mut()
757 .unwrap()
758 .purge_old(now);
759 
760 let bytes_transmitted = self.stats.bytes_transmitted.as_mut().unwrap().sum();
761 let bytes_retransmitted = self.stats.bytes_retransmitted.as_mut().unwrap().sum();
762 let ratio = bytes_retransmitted as f32 / (bytes_retransmitted + bytes_transmitted) as f32;
763 let ratio = if ratio.is_finite() { ratio } else { 0_f32 };
764 self.rtx_ratio = (ratio, now);
765 ratio
766 }
767 
768 fn poll_packet_resend(&mut self, now: Instant) -> Option<NextPacket<'_>> {
769 if let Some(ratio_cap) = self.rtx_ratio_cap {
770 let ratio = self.rtx_ratio_downsampled(now);
771 
772 // If we hit the cap, stop doing resends by clearing those we have queued.
773 if ratio > ratio_cap {
774 self.resends.clear();
775 return None;
776 }
777 }
778 
779 let seq_no = loop {
780 let resend = self.resends.pop_front()?;
781 
782 let pkt = self.rtx_cache.get_cached_packet_by_seq_no(resend.seq_no);
783 
784 // The seq_no could simply be too old to exist in the buffer, in which
785 // case we will not do a resend.
786 let Some(pkt) = pkt else {
787 continue;
788 };
789 
790 // Cached packets must be nackable. This is ensured before adding the
791 // entry to the self.rtx_cache.
792 assert!(pkt.nackable);
793 
794 break pkt.seq_no;
795 };
796 
797 // Borrow checker gymnastics.
798 let pkt = self.rtx_cache.get_cached_packet_by_seq_no(seq_no).unwrap();
799 
800 let len = pkt.payload.len() as u64;
801 self.stats.update_packet_counts(len, true);
802 if let Some(h) = &mut self.stats.bytes_retransmitted {
803 h.push(now, len);
804 }
805 
806 let seq_no = self.seq_no_rtx.inc();
807 
808 let orig_seq_no = pkt.seq_no;
809 
810 Some(NextPacket {
811 kind: NextPacketKind::Resend(orig_seq_no),
812 seq_no,
813 pkt,
814 })
815 }
816 
817 fn poll_packet_regular(&mut self, now: Instant) -> Option<NextPacket<'_>> {
818 // exit via ? here is ok since that means there is nothing to send.
819 // The packet remains in the head of the send queue until we
820 // finish poll_packet, at which point we move it to the cache.
821 let pkt = self.send_queue.peek()?;
822 
823 pkt.timestamp = now;
824 
825 let len = pkt.payload.len() as u64;
826 self.stats.update_packet_counts(len, false);
827 if let Some(h) = &mut self.stats.bytes_transmitted {
828 h.push(now, len)
829 }
830 
831 let seq_no = pkt.seq_no;
832 
833 Some(NextPacket {
834 kind: NextPacketKind::Regular,
835 seq_no,
836 pkt,
837 })
838 }
839 
840 fn poll_packet_padding(&mut self, _now: Instant) -> Option<NextPacket> {
841 if !self.padding_enabled() {
842 self.padding = 0;
843 return None;
844 }
845 
846 if self.padding == 0 {
847 return None;
848 }
849 
850 #[allow(clippy::unnecessary_operation)]
851 'outer: {
852 if self.padding > MIN_SPURIOUS_PADDING_SIZE {
853 // Find a historic packet that is smaller than this max size. The max size
854 // is a headroom since we can accept slightly larger padding than asked for.
855 let max_size = (self.padding * 2).min(self.mtu_warn - MAX_RTP_OVERHEAD);
856 
857 let Some(pkt) = self.rtx_cache.get_cached_packet_smaller_than(max_size) else {
858 // Couldn't find spurious packet, try a blank packet instead.
859 break 'outer;
860 };
861 
862 let orig_seq_no = pkt.seq_no;
863 let seq_no = self.seq_no_rtx.inc();
864 
865 self.padding = self.padding.saturating_sub(pkt.payload.len());
866 
867 return Some(NextPacket {
868 kind: NextPacketKind::Resend(orig_seq_no),
869 seq_no,
870 pkt,
871 });
872 }
873 };
874 
875 let seq_no = self.seq_no_rtx.inc();
876 
877 let pkt = &mut self.blank_packet;
878 pkt.seq_no = seq_no;
879 // Unwrap here is correct because self.padding_enabled() above checks the we got the PT set.
880 pkt.header.payload_type = self.pt_for_padding.unwrap();
881 
882 let len = self
883 .padding
884 .clamp(SRTP_BLOCK_SIZE, MAX_BLANK_PADDING_PAYLOAD_SIZE);
885 assert!(len <= 255); // should fit in a byte
886 
887 self.padding = self.padding.saturating_sub(len);
888 
889 Some(NextPacket {
890 kind: NextPacketKind::Blank(len as u8),
891 seq_no,
892 pkt,
893 })
894 }
895 
896 fn is_audio(&self) -> bool {
897 self.kind.is_some_and(|kind| kind.is_audio())
898 }
899 
900 pub(crate) fn sender_report_at(&self, intervals: RtcpReportIntervals) -> Instant {
901 if self.kind.is_none() {
902 // First handle_timeout sets the kind. No sender report until then.
903 return not_happening();
904 }
905 self.last_sender_report + intervals.for_audio(self.is_audio())
906 }
907 
908 pub(crate) fn poll_keyframe_request(&mut self) -> Option<KeyframeRequestKind> {
909 self.pending_request_keyframe.take()
910 }
911 
912 pub(crate) fn poll_remb_request(&mut self) -> Option<Bitrate> {
913 self.pending_request_remb.take()
914 }
915 
916 pub(crate) fn handle_rtcp(&mut self, now: Instant, fb: RtcpFb) {
917 use RtcpFb::*;
918 match fb {
919 ReceptionReport(r) => {
920 if let Some(rtx_ssrc) = self.rtx {
921 if rtx_ssrc == r.ssrc {
922 // Receiver has bound MidRid to RTX SSRC
923 self.remote_acked_rtx_ssrc = true;
924 }
925 } else if r.ssrc == self.ssrc {
926 // Receiver has bound MidRid to SSRC
927 self.remote_acked_ssrc = true;
928 }
929 
930 self.stats.update_with_rr(now, self.last_sent_seq_no, r)
931 }
932 Nack(_, list) => {
933 self.stats.increase_nacks();
934 let entries = list.into_iter();
935 self.handle_nack(entries, now);
936 }
937 Pli(_) => {
938 self.stats.increase_plis();
939 self.pending_request_keyframe = Some(KeyframeRequestKind::Pli);
940 }
941 Fir(_) => {
942 self.stats.increase_firs();
943 self.pending_request_keyframe = Some(KeyframeRequestKind::Fir);
944 }
945 Remb(r) => {
946 self.pending_request_remb = Some(Bitrate::from(r.bitrate as f64));
947 }
948 Twcc(_) => unreachable!("TWCC should be handled on session level"),
949 _ => {}
950 }
951 }
952 
953 pub(crate) fn handle_nack(
954 &mut self,
955 entries: impl Iterator<Item = NackEntry>,
956 now: Instant,
957 ) -> Option<()> {
958 // Turning NackEntry into SeqNo we need to know a SeqNo "close by" to lengthen the 16 bit
959 // sequence number into the 64 bit we have in SeqNo.
960 let seq_no = self.rtx_cache.last_cached_seq_no()?;
961 let iter = entries.flat_map(|n| n.into_iter(seq_no));
962 
963 // Schedule all resends. They will be handled on next poll_packet
964 for seq_no in iter {
965 let Some(packet) = self.rtx_cache.get_cached_packet_by_seq_no(seq_no) else {
966 // Packet was not available in RTX cache, it has probably expired.
967 continue;
968 };
969 
970 let resend = Resend {
971 seq_no,
972 queued_at: now,
973 payload_size: packet.payload.len(),
974 };
975 self.resends.push_back(resend);
976 }
977 
978 Some(())
979 }
980 
981 pub(crate) fn need_sr(&self, now: Instant, intervals: RtcpReportIntervals) -> bool {
982 now >= self.sender_report_at(intervals)
983 }
984 
985 pub(crate) fn create_sr_and_update(&mut self, now: Instant, feedback: &mut VecDeque<Rtcp>) {
986 let sr = self.create_sender_report(now);
987 
988 trace!("Created feedback SR: {:?}", sr);
989 feedback.push_back(Rtcp::SenderReport(sr));
990 
991 if let Some(ds) = self.create_sdes() {
992 feedback.push_back(Rtcp::SourceDescription(ds));
993 }
994 
995 // Update timestamp to move time when next is created.
996 self.last_sender_report = now;
997 }
998 
999 fn create_sender_report(&self, now: Instant) -> SenderReport {
1000 SenderReport {
1001 sender_info: self.sender_info(now),
1002 reports: ReportList::new(),
1003 }
1004 }
1005 
1006 fn create_sdes(&self) -> Option<Descriptions> {
1007 // CNAME is set on first handle_timeout. No SDES before that.
1008 let cname = self.cname.as_ref()?;
1009 let mut s = Sdes {
1010 ssrc: self.ssrc,
1011 values: ReportList::new(),
1012 };
1013 s.values.push((SdesType::CNAME, cname.to_string()));
1014 
1015 let mut d = Descriptions {
1016 reports: Box::new(ReportList::new()),
1017 };
1018 d.reports.push(s);
1019 
1020 Some(d)
1021 }
1022 
1023 fn sender_info(&self, now: Instant) -> SenderInfo {
1024 let rtp_time = self.current_rtp_time(now).unwrap_or(MediaTime::ZERO);
1025 
1026 SenderInfo {
1027 ssrc: self.ssrc,
1028 ntp_time: now.to_system_time(),
1029 rtp_time,
1030 sender_packet_count: self.stats.packets as u32,
1031 sender_octet_count: self.stats.bytes as u32,
1032 }
1033 }
1034 
1035 fn current_rtp_time(&self, now: Instant) -> Option<MediaTime> {
1036 // This is the RTP time and the wallclock from the last written media.
1037 // We use that as an offset to current time (now), to calculate the
1038 // current RTP time.
1039 let (t_u32, w) = self.rtp_and_wallclock?;
1040 
1041 let clock_rate = self.clock_rate?;
1042 let t = MediaTime::new(t_u32 as u64, clock_rate);
1043 
1044 // Wallclock needs to be in the past.
1045 if w > now {
1046 let delta = w - now;
1047 debug!("write_rtp wallclock is in the future: {:?}", delta);
1048 return None;
1049 }
1050 let offset = now - w;
1051 
1052 // This might be in the wrong base.
1053 let rtp_time = t + offset.into();
1054 
1055 Some(rtp_time.rebase(clock_rate))
1056 }
1057 
1058 pub(crate) fn next_seq_no(&mut self) -> SeqNo {
1059 self.seq_no.inc()
1060 }
1061 
1062 pub(crate) fn last_packet(&self) -> Option<&[u8]> {
1063 if self.send_queue.is_empty() {
1064 self.rtx_cache.last_packet()
1065 } else {
1066 self.send_queue.last().map(|q| q.payload.as_ref())
1067 }
1068 }
1069 
1070 pub(crate) fn visit_stats(&mut self, snapshot: &mut StatsSnapshot, now: Instant) {
1071 self.stats.fill(snapshot, self.midrid, now);
1072 }
1073 
1074 pub(crate) fn queue_state(&mut self, now: Instant) -> QueueState {
1075 // The unpaced flag is set to a default value on first handle_timeout. The
1076 // default is to not pace audio. We unwrap default to "true" here to not
1077 // apply any pacing until we know what kind of content we are sending.
1078 let unpaced = self.unpaced.unwrap_or(true);
1079 
1080 // It's only possible to use this sender for padding if RTX is enabled and
1081 // we know a PT to use for it.
1082 let use_for_padding = self.padding_enabled();
1083 
1084 let mut snapshot = self.send_queue.snapshot(now);
1085 
1086 if let Some(snapshot_resend) = self.queue_state_resend(now) {
1087 snapshot.merge(&snapshot_resend);
1088 }
1089 
1090 if let Some(snapshot_padding) = self.queue_state_padding(now) {
1091 snapshot.merge(&snapshot_padding);
1092 }
1093 
1094 self.queue_info = Some(snapshot);
1095 
1096 QueueState {
1097 midrid: self.midrid,
1098 unpaced,
1099 use_for_padding,
1100 snapshot,
1101 }
1102 }
1103 
1104 fn queue_state_resend(&self, now: Instant) -> Option<QueueSnapshot> {
1105 if self.resends.is_empty() {
1106 return None;
1107 }
1108 
1109 // Outstanding resends
1110 let mut snapshot = self
1111 .resends
1112 .iter()
1113 .fold(QueueSnapshot::default(), |mut snapshot, r| {
1114 snapshot.total_queue_time_origin += now.duration_since(r.queued_at);
1115 snapshot.byte_size += r.payload_size;
1116 snapshot.packet_count += 1;
1117 snapshot.first_unsent = snapshot
1118 .first_unsent
1119 .map(|i| i.min(r.queued_at))
1120 .or(Some(r.queued_at));
1121 
1122 snapshot
1123 });
1124 snapshot.created_at = now;
1125 snapshot.update_priority(QueuePriority::Media);
1126 
1127 Some(snapshot)
1128 }
1129 
1130 fn queue_state_padding(&self, now: Instant) -> Option<QueueSnapshot> {
1131 if self.padding == 0 {
1132 return None;
1133 }
1134 
1135 // TODO: Be more scientific about these factors.
1136 const AVERAGE_PADDING_PACKET_SIZE: usize = 800;
1137 const FAKE_PADDING_DURATION_MILLIS: usize = 5;
1138 
1139 let fake_packets = self.padding.div_ceil(AVERAGE_PADDING_PACKET_SIZE);
1140 let fake_millis = fake_packets * FAKE_PADDING_DURATION_MILLIS;
1141 let fake_duration = Duration::from_millis(fake_millis as u64);
1142 
1143 Some(QueueSnapshot {
1144 created_at: now,
1145 byte_size: self.padding,
1146 packet_count: fake_packets as u32,
1147 total_queue_time_origin: fake_duration,
1148 priority: QueuePriority::Padding,
1149 ..Default::default()
1150 })
1151 }
1152 
1153 pub(crate) fn generate_padding(&mut self, padding: usize) {
1154 if !self.padding_enabled() {
1155 return;
1156 }
1157 self.padding += padding;
1158 }
1159 
1160 pub(crate) fn need_timeout(&self) -> bool {
1161 self.send_queue.need_timeout()
1162 }
1163 
1164 pub(crate) fn handle_timeout<'a>(
1165 &mut self,
1166 now: Instant,
1167 get_media: impl FnOnce() -> (&'a Media, &'a CodecConfig),
1168 ) {
1169 // If kind is None, this is the first time we ever get a handle_timeout.
1170 if self.kind.is_none() {
1171 let (media, config) = get_media();
1172 self.on_first_timeout(media, config);
1173 }
1174 
1175 self.send_queue.handle_timeout(now);
1176 }
1177 
1178 fn on_first_timeout(&mut self, media: &Media, config: &CodecConfig) {
1179 // Always set on first timeout.
1180 self.kind = Some(media.kind());
1181 self.cname = Some(media.cname().to_string());
1182 
1183 // Set on first timeout, if not set already by configuration.
1184 if self.unpaced.is_none() {
1185 // Default audio to be unpaced.
1186 self.unpaced = Some(media.kind().is_audio());
1187 }
1188 
1189 // To allow for sending padding on a newly created StreamTx, before any regular
1190 // packet has been sent, we need any main PT that has associated RTX. This is
1191 // later be overwritten when we send the first regular packet.
1192 if self.rtx.is_some() && self.pt_for_padding.is_none() {
1193 if let Some(pt) = media.first_pt_with_rtx(config) {
1194 trace!(
1195 "StreamTx {:?} PT {} before first regular packet",
1196 self.midrid, pt
1197 );
1198 self.pt_for_padding = Some(pt);
1199 
1200 // Setting the pt_for_rtx should enable RTX.
1201 assert!(self.padding_enabled());
1202 }
1203 }
1204 }
1205 
1206 pub(crate) fn reset_buffers(&mut self) {
1207 self.send_queue.clear();
1208 self.queue_info = None;
1209 self.rtx_cache.clear();
1210 self.resends.clear();
1211 self.padding = 0;
1212 }
1213 
1214 /// Reset this stream to use a new SSRC and optionally a new RTX SSRC.
1215 ///
1216 /// This updates the SSRCs and resets all relevant internal fields.
1217 pub(crate) fn reset_ssrc(&mut self, new_ssrc: Ssrc, new_rtx: Option<Ssrc>) {
1218 // Update the SSRC and RTX
1219 self.ssrc = new_ssrc;
1220 self.rtx = new_rtx;
1221 
1222 // Reset sequence numbers
1223 self.seq_no = SeqNo::default();
1224 self.seq_no_rtx = SeqNo::default();
1225 
1226 // Reset timing related fields
1227 self.last_used = already_happened();
1228 self.rtp_and_wallclock = None;
1229 self.last_sender_report = already_happened();
1230 
1231 // Reset blank packet's SSRC
1232 self.blank_packet.header.ssrc = new_ssrc;
1233 
1234 // Clear any pending requests
1235 self.pending_request_keyframe = None;
1236 self.pending_request_remb = None;
1237 
1238 // Reset all statistics - preserve whether stats tracking is enabled
1239 let stats_enabled = self.stats.bytes_transmitted.is_some();
1240 self.stats = StreamTxStats::new(stats_enabled);
1241 self.rtx_ratio = (0.0, already_happened());
1242 self.remote_acked_ssrc = false;
1243 
1244 // Clear all buffers
1245 self.reset_buffers();
1246 }
1247 
1248 pub(crate) fn is_midrid(&self, midrid: MidRid) -> bool {
1249 midrid.special_equals(&self.midrid)
1250 }
1251}
1252 
1253struct NextPacket<'a> {
1254 kind: NextPacketKind,
1255 seq_no: SeqNo,
1256 pkt: &'a mut RtpPacket,
1257}
1258 
1259#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1260enum NextPacketKind {
1261 Regular,
1262 Resend(SeqNo),
1263 Blank(u8),
1264}
1265 
1266#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1267struct Resend {
1268 seq_no: SeqNo,
1269 queued_at: Instant,
1270 payload_size: usize,
1271}
1272 
1273#[cfg(test)]
1274mod test {
1275 use super::*;
1276 
1277 #[test]
1278 fn queue_info_is_cached_on_queue_state_update() {
1279 let mut stream = StreamTx::new(42.into(), None, MidRid(Mid::from("0"), None), false, 1200);
1280 
1281 assert!(stream.queue_info().is_none());
1282 
1283 let now = Instant::now();
1284 let first_unsent = now - Duration::from_millis(10);
1285 stream.resends.push_back(Resend {
1286 seq_no: 7.into(),
1287 queued_at: first_unsent,
1288 payload_size: 123,
1289 });
1290 stream.queue_state(now);
1291 
1292 let info = stream.queue_info().expect("queue info after pacer update");
1293 assert_eq!(info.created_at(), now);
1294 assert_eq!(info.byte_size(), 123);
1295 assert_eq!(info.packet_count(), 1);
1296 assert_eq!(info.first_unsent(), Some(first_unsent));
1297 
1298 stream.reset_buffers();
1299 assert!(stream.queue_info().is_none());
1300 }
1301 
1302 #[test]
1303 fn sub_average_padding_is_visible_to_the_pacer() {
1304 let mut stream = StreamTx::new(42.into(), None, MidRid(Mid::from("0"), None), false, 1200);
1305 stream.padding = 240;
1306 
1307 let snapshot = stream
1308 .queue_state_padding(Instant::now())
1309 .expect("nonzero padding must produce a queue snapshot");
1310 
1311 assert_eq!(snapshot.byte_size, 240);
1312 assert_eq!(snapshot.packet_count, 1);
1313 }
1314 
1315 #[test]
1316 fn sender_report_uses_supplied_intervals() {
1317 let mut stream = StreamTx::new(42.into(), None, MidRid(Mid::from("0"), None), false, 1200);
1318 let now = Instant::now();
1319 stream.kind = Some(MediaKind::Audio);
1320 stream.last_sender_report = now;
1321 
1322 assert_eq!(
1323 stream.sender_report_at(RtcpReportIntervals {
1324 audio: Duration::from_millis(750),
1325 video: Duration::from_millis(500),
1326 }),
1327 now + Duration::from_millis(750)
1328 );
1329 }
1330}