Skip to content
File

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

rust338 lines
1//! Statistics events.
2 
3use std::collections::{HashMap, VecDeque};
4use std::net::SocketAddr;
5use std::time::{Duration, Instant};
6 
7use crate::Bitrate;
8use crate::rtp::SeqNo;
9use crate::rtp_::{Mid, Rid};
10use crate::{io::Protocol, rtp_::MidRid};
11 
12pub(crate) struct Stats {
13 last_now: Option<Instant>,
14 events: VecDeque<StatsEvent>,
15 interval: Duration,
16}
17 
18pub(crate) struct StatsSnapshot {
19 pub peer_tx: u64,
20 pub peer_rx: u64,
21 pub tx: u64,
22 pub rx: u64,
23 pub egress_loss_fraction: Option<f32>,
24 pub ingress_loss_fraction: Option<f32>,
25 pub rtt: Option<Duration>,
26 pub ingress: HashMap<MidRid, MediaIngressStats>,
27 pub egress: HashMap<MidRid, MediaEgressStats>,
28 pub bwe_tx: Option<Bitrate>,
29 pub selected_candidate_pair: Option<CandidatePairStats>,
30 timestamp: Instant,
31}
32 
33impl StatsSnapshot {
34 pub(crate) fn new(timestamp: Instant) -> StatsSnapshot {
35 StatsSnapshot {
36 peer_rx: 0,
37 peer_tx: 0,
38 tx: 0,
39 rx: 0,
40 egress_loss_fraction: None,
41 ingress_loss_fraction: None,
42 ingress: HashMap::new(),
43 egress: HashMap::new(),
44 rtt: None,
45 bwe_tx: None,
46 selected_candidate_pair: None,
47 timestamp,
48 }
49 }
50}
51 
52// Output events
53 
54#[derive(Debug, Clone)]
55pub(crate) enum StatsEvent {
56 Peer(PeerStats),
57 MediaEgress(MediaEgressStats),
58 MediaIngress(MediaIngressStats),
59}
60 
61/// Peer statistics in [`Event::PeerStats`][crate::Event::PeerStats].
62///
63/// This event is generated roughly every second
64#[derive(Debug, Clone)]
65pub struct PeerStats {
66 /// Total bytes transmitted.
67 pub peer_bytes_rx: u64,
68 /// Total bytes received.
69 pub peer_bytes_tx: u64,
70 /// Total bytes transmitted, only counting media traffic (rtp payload).
71 pub bytes_rx: u64,
72 /// Total bytes received, only counting media traffic (rtp payload).
73 pub bytes_tx: u64,
74 /// Timestamp when this event was generated.
75 pub timestamp: Instant,
76 /// The last egress bandwidth estimate from the BWE subsystem, if enabled.
77 pub bwe_tx: Option<Bitrate>,
78 /// The egress loss over the last second.
79 pub egress_loss_fraction: Option<f32>,
80 /// The ingress loss since the last stats event.
81 pub ingress_loss_fraction: Option<f32>,
82 /// The most recent RTT since the last stats event.
83 pub rtt: Option<Duration>,
84 /// The selected ICE candidate pair, if any.
85 pub selected_candidate_pair: Option<CandidatePairStats>,
86}
87 
88#[derive(Debug, Clone)]
89/// ICE candidate pair statistics.
90pub struct CandidatePairStats {
91 /// The selected protocol.
92 pub protocol: Protocol,
93 /// The local candidate.
94 pub local: CandidateStats,
95 /// The remote candidate.
96 pub remote: CandidateStats,
97 /// Round-trip time of the most recent STUN binding transaction on the
98 /// nominated pair.
99 ///
100 /// Spec equivalent of [`RTCIceCandidatePairStats.currentRoundTripTime`][1].
101 ///
102 /// This is sourced from the ICE keepalive/consent traffic and is
103 /// available even on receive-only endpoints where RTP-based RTT
104 /// estimates aren't.
105 ///
106 /// [1]: https://www.w3.org/TR/webrtc-stats/#dom-rtcicecandidatepairstats-currentroundtriptime
107 pub current_round_trip_time: Option<Duration>,
108 /// Sum of all RTTs measured from successful STUN binding transactions
109 /// on the nominated pair over its lifetime. Divide by
110 /// [`responses_received`][Self::responses_received] for the average.
111 ///
112 /// Spec equivalent of [`RTCIceCandidatePairStats.totalRoundTripTime`][1].
113 ///
114 /// [1]: https://www.w3.org/TR/webrtc-stats/#dom-rtcicecandidatepairstats-totalroundtriptime
115 pub total_round_trip_time: Duration,
116 /// Total number of STUN binding responses received on the nominated
117 /// pair over its lifetime.
118 ///
119 /// Spec equivalent of [`RTCIceCandidatePairStats.responsesReceived`][1].
120 ///
121 /// [1]: https://www.w3.org/TR/webrtc-stats/#dom-rtcicecandidatepairstats-responsesreceived
122 pub responses_received: u64,
123}
124 
125#[derive(Debug, Clone)]
126/// ICE candidate statistics.
127pub struct CandidateStats {
128 /// The address of the candidate.
129 pub addr: SocketAddr,
130}
131 
132/// Outgoing media statistics in [`Event::MediaEgressStats`][crate::Event::MediaEgressStats].
133///
134/// note: when simulcast is disabled, `rid` is `None`
135#[derive(Debug, Clone)]
136pub struct MediaEgressStats {
137 /// The identifier of the media these stats are for.
138 pub mid: Mid,
139 /// The Rid identifier in case of simulcast.
140 pub rid: Option<Rid>,
141 /// Total bytes sent, including retransmissions.
142 ///
143 /// Spec equivalent to [`RTCSentRtpStreamStats.bytesSent`][1].
144 ///
145 /// [1]: https://www.w3.org/TR/webrtc-stats/#dom-rtcsentrtpstreamstats-bytessent
146 pub bytes: u64,
147 /// Total number of rtp packets sent, including retransmissions
148 ///
149 /// Spec equivalent of [`RTCSentRtpStreamStats.packetsSent`][1].
150 ///
151 /// [1]: https://www.w3.org/TR/webrtc-stats/#dom-rtcsentrtpstreamstats-packetssent
152 pub packets: u64,
153 /// Number of firs received.
154 pub firs: u64,
155 /// Number of plis received.
156 pub plis: u64,
157 /// Number of nacks received.
158 pub nacks: u64,
159 /// Round-trip-time extracted from the last RTCP receiver report.
160 pub rtt: Option<Duration>,
161 /// Fraction of packets lost averaged from the RTCP receiver reports received.
162 /// `None` if no reports have been received since the last event
163 pub loss: Option<f32>,
164 /// Timestamp when this event was generated
165 pub timestamp: Instant,
166 /// Stats provided by the remote peer via ReceiverReports
167 pub remote: Option<RemoteIngressStats>,
168}
169 
170/// Stats as reported by the remote side (via RTCP ReceiverReports).
171#[derive(Debug, Clone)]
172pub struct RemoteIngressStats {
173 /// The remotely calculated jitter.
174 pub jitter: u32,
175 /// The maximum extended sequence number received.
176 pub maximum_sequence_number: SeqNo,
177 /// The cumulative number of packets lost.
178 pub packets_lost: u64,
179}
180 
181/// Incoming media statistics in [`Event::MediaIngressStats`][crate::Event::MediaIngressStats].
182///
183/// note: when simulcast is disabled, `rid` is `None`
184#[derive(Debug, Clone)]
185pub struct MediaIngressStats {
186 /// The identifier of the media these stats are for.
187 pub mid: Mid,
188 /// The Rid identifier in case of simulcast.
189 pub rid: Option<Rid>,
190 /// Total bytes received, including retransmissions.
191 pub bytes: u64,
192 /// Total number of rtp packets received, including retransmissions.
193 pub packets: u64,
194 /// Number of firs sent.
195 pub firs: u64,
196 /// Number of plis sent.
197 pub plis: u64,
198 /// Number of nacks sent.
199 pub nacks: u64,
200 /// Interarrival jitter, in RTP timestamp units, as reported in the last RTCP
201 /// receiver report we sent for this stream.
202 pub jitter: u32,
203 /// Round-trip-time extracted from the last RTCP XR DLRR report block.
204 pub rtt: Option<Duration>,
205 /// Fraction of packets lost extracted from the last RTCP receiver report.
206 pub loss: Option<f32>,
207 /// Timestamp when this event was generated.
208 pub timestamp: Instant,
209 /// Stats provided by the remote peer via SenderReports
210 pub remote: Option<RemoteEgressStats>,
211}
212 
213impl MediaIngressStats {
214 /// Merge `other` into `self`, mutating `self`.
215 ///
216 /// **Panics** if called with stats that don't have the same `(mid, rid)` pair.
217 pub(crate) fn merge_by_mid_rid(&mut self, other: &Self) {
218 assert!(
219 self.mid == other.mid,
220 "Cannot merge MediaIngressStats for different mids"
221 );
222 assert!(
223 self.rid == other.rid,
224 "Cannot merge MediaIngressStats for different rids"
225 );
226 let (jitter, rtt, loss) = if self.timestamp > other.timestamp {
227 (self.jitter, self.rtt, self.loss)
228 } else {
229 (other.jitter, other.rtt, other.loss)
230 };
231 
232 *self = Self {
233 mid: self.mid,
234 rid: self.rid,
235 bytes: self.bytes + other.bytes,
236 packets: self.packets + other.packets,
237 firs: self.firs + other.firs,
238 plis: self.plis + other.plis,
239 nacks: self.nacks + other.nacks,
240 jitter,
241 rtt,
242 loss,
243 timestamp: self.timestamp.max(other.timestamp),
244 remote: match (&self.remote, &other.remote) {
245 (None, None) => None,
246 (Some(remote), None) => Some(remote.clone()),
247 (None, Some(other_remote)) => Some(other_remote.clone()),
248 (Some(remote), Some(other_remote)) => Some(RemoteEgressStats {
249 bytes: remote.bytes + other_remote.bytes,
250 packets: remote.packets + other_remote.packets,
251 }),
252 },
253 };
254 }
255}
256 
257/// Stats as reported by the remote side (via RTCP SenderReports).
258#[derive(Debug, Clone)]
259pub struct RemoteEgressStats {
260 /// Total bytes sent, including retransmissions.
261 pub bytes: u64,
262 /// Total number of rtp packets sent, including retransmissions.
263 pub packets: u64,
264}
265 
266impl Stats {
267 /// Create a new stats instance
268 ///
269 /// The internal state is market with the current `Instant::now()`.
270 /// This allows us to emit stats right away at the first upcoming timeout
271 pub fn new(interval: Duration) -> Stats {
272 Stats {
273 // by starting with the current time we can generate stats right on first timeout
274 last_now: None,
275 events: VecDeque::new(),
276 interval,
277 }
278 }
279 
280 /// Returns true if we want to handle the timeout
281 ///
282 /// The caller can use this to compute the snapshot only if needed, before calling \
283 /// [`Stats::do_handle_timeout`]
284 pub fn wants_timeout(&mut self, now: Instant) -> bool {
285 let Some(last_now) = self.last_now else {
286 // Learn our first ever `now`
287 self.last_now = Some(now);
288 return false;
289 };
290 
291 let min_step = last_now + self.interval;
292 now >= min_step
293 }
294 
295 /// Actually handles the timeout advancing the internal state and preparing the output
296 pub fn do_handle_timeout(&mut self, snapshot: &mut StatsSnapshot) {
297 // enqueue stats and timestamp them so they can be sent out
298 
299 let event = PeerStats {
300 peer_bytes_rx: snapshot.peer_rx,
301 peer_bytes_tx: snapshot.peer_tx,
302 bytes_rx: snapshot.rx,
303 bytes_tx: snapshot.tx,
304 timestamp: snapshot.timestamp,
305 bwe_tx: snapshot.bwe_tx,
306 egress_loss_fraction: snapshot.egress_loss_fraction,
307 ingress_loss_fraction: snapshot.ingress_loss_fraction,
308 rtt: snapshot.rtt,
309 selected_candidate_pair: snapshot.selected_candidate_pair.clone(),
310 };
311 
312 self.events.push_back(StatsEvent::Peer(event));
313 
314 for (_, event) in snapshot.ingress.drain() {
315 self.events.push_back(StatsEvent::MediaIngress(event));
316 }
317 
318 for (_, event) in snapshot.egress.drain() {
319 self.events.push_back(StatsEvent::MediaEgress(event));
320 }
321 
322 self.last_now = Some(snapshot.timestamp);
323 }
324 
325 /// Poll for the next time to call [`Stats::wants_timeout`] and [`Stats::do_handle_timeout`].
326 ///
327 /// NOTE: we only need Option<_> to conform to .soonest() (see caller)
328 pub fn poll_timeout(&self) -> Option<Instant> {
329 let last_now = self.last_now?;
330 Some(last_now + self.interval)
331 }
332 
333 /// Return any events ready for delivery
334 pub fn poll_output(&mut self) -> Option<StatsEvent> {
335 self.events.pop_front()
336 }
337}