Skip to content
File

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

rust191 lines
1use std::time::{Duration, Instant};
2 
3use crate::rtp::SeqNo;
4use crate::rtp_::{ReceptionReport, extend_u32};
5use crate::stats::{MediaEgressStats, RemoteIngressStats, StatsSnapshot};
6use crate::util::value_history::ValueHistory;
7use crate::util::{InstantExt, calculate_rtt};
8 
9use super::MidRid;
10 
11/// Holder of stats.
12#[derive(Debug)]
13pub(crate) struct StreamTxStats {
14 /// count of bytes sent, including retransmissions
15 /// <https://www.w3.org/TR/webrtc-stats/#dom-rtcsentrtpstreamstats-bytessent>
16 pub bytes: u64,
17 /// count of retransmitted bytes alone
18 bytes_resent: u64,
19 /// count of packets sent, including retransmissions
20 /// <https://www.w3.org/TR/webrtc-stats/#summary>
21 pub packets: u64,
22 /// count of retransmitted packets alone
23 packets_resent: u64,
24 /// count of FIR requests received
25 firs: u64,
26 /// count of PLI requests received
27 plis: u64,
28 /// count of NACKs received
29 nacks: u64,
30 /// round trip time
31 /// Can be null in case of missing or bad reports
32 rtt: Option<Duration>,
33 /// losses collecter from RR (known packets, lost ratio)
34 losses: Losses,
35 /// The last reception report for the stream, if any.
36 ///
37 /// The SeqNo is the extended max_seq of the reception report, extended
38 /// using the last sent sequence number.
39 last_rr: Option<(SeqNo, ReceptionReport)>,
40 
41 /// `None` if `rtx_ratio_cap` is `None`.
42 pub bytes_transmitted: Option<ValueHistory<u64>>,
43 
44 /// `None` if `rtx_ratio_cap` is `None`.
45 pub bytes_retransmitted: Option<ValueHistory<u64>>,
46}
47 
48impl StreamTxStats {
49 pub fn new(enable_stats: bool) -> Self {
50 Self {
51 bytes: 0,
52 bytes_resent: 0,
53 packets: 0,
54 packets_resent: 0,
55 firs: 0,
56 plis: 0,
57 nacks: 0,
58 rtt: None,
59 losses: Losses::new(enable_stats),
60 last_rr: None,
61 bytes_transmitted: Some(Default::default()),
62 bytes_retransmitted: Some(Default::default()),
63 }
64 }
65 
66 pub fn update_packet_counts(&mut self, bytes: u64, is_resend: bool) {
67 self.packets += 1;
68 self.bytes += bytes;
69 if is_resend {
70 self.bytes_resent += bytes;
71 self.packets_resent += 1;
72 }
73 }
74 
75 pub fn increase_nacks(&mut self) {
76 self.nacks += 1;
77 }
78 
79 pub fn increase_plis(&mut self) {
80 self.plis += 1;
81 }
82 
83 pub fn increase_firs(&mut self) {
84 self.firs += 1;
85 }
86 
87 pub fn update_with_rr(&mut self, now: Instant, last_sent_seq_no: SeqNo, r: ReceptionReport) {
88 let ntp_time = now.to_ntp_duration();
89 let rtt = calculate_rtt(ntp_time, r.last_sr_delay, r.last_sr_time);
90 self.rtt = rtt;
91 
92 // The last_sent_seq_no should be in the vicinity of the rr.max_seq.
93 let ext_seq = extend_u32(Some(*last_sent_seq_no), r.max_seq).into();
94 
95 self.last_rr = Some((ext_seq, r));
96 
97 self.losses
98 .push((*ext_seq, r.fraction_lost as f32 / u8::MAX as f32));
99 }
100 
101 pub(crate) fn fill(&mut self, snapshot: &mut StatsSnapshot, midrid: MidRid, now: Instant) {
102 if self.bytes == 0 {
103 return;
104 }
105 
106 let loss = {
107 let mut value = 0_f32;
108 let mut total_weight = 0_u64;
109 
110 // average known RR losses weighted by their number of packets
111 for it in self.losses.iterator() {
112 let [prev, next] = it else { continue };
113 let weight = next.0.saturating_sub(prev.0);
114 value += next.1 * weight as f32;
115 total_weight += weight;
116 }
117 
118 let result = value / total_weight as f32;
119 result.is_finite().then_some(result)
120 };
121 
122 self.losses.clear_all_but_last();
123 
124 snapshot.egress.insert(
125 midrid,
126 MediaEgressStats {
127 mid: midrid.mid(),
128 rid: midrid.rid(),
129 bytes: self.bytes,
130 packets: self.packets,
131 firs: self.firs,
132 plis: self.plis,
133 nacks: self.nacks,
134 rtt: self.rtt,
135 loss,
136 timestamp: now,
137 remote: self
138 .last_rr
139 .as_ref()
140 .map(|(seq_no, rr)| RemoteIngressStats {
141 jitter: rr.jitter,
142 maximum_sequence_number: *seq_no,
143 packets_lost: rr.packets_lost as u64,
144 }),
145 },
146 );
147 }
148}
149 
150/// Helper to avoid an unbounded vec if we are not enabling stats
151#[derive(Debug)]
152enum Losses {
153 Disabled,
154 Enabled(Vec<(u64, f32)>),
155}
156 
157impl Losses {
158 fn new(enabled: bool) -> Self {
159 if enabled {
160 Self::Enabled(vec![])
161 } else {
162 Self::Disabled
163 }
164 }
165 
166 fn push(&mut self, value: (u64, f32)) {
167 let Self::Enabled(losses) = self else {
168 return;
169 };
170 losses.push(value);
171 }
172 
173 fn iterator(&mut self) -> impl Iterator<Item = &[(u64, f32)]> {
174 let Self::Enabled(losses) = self else {
175 return [].windows(2);
176 };
177 
178 // just in case we received RRs out of order
179 losses.sort_by(|a, b| a.0.partial_cmp(&b.0).unwrap());
180 
181 losses.windows(2)
182 }
183 
184 fn clear_all_but_last(&mut self) {
185 let Self::Enabled(losses) = self else {
186 return;
187 };
188 losses.drain(..losses.len().saturating_sub(1));
189 }
190}