File
Blob: firmware/vendor/str0m/src/stats.rs
| 1 | //! Statistics events. |
| 2 | |
| 3 | use std::collections::{HashMap, VecDeque}; |
| 4 | use std::net::SocketAddr; |
| 5 | use std::time::{Duration, Instant}; |
| 6 | |
| 7 | use crate::Bitrate; |
| 8 | use crate::rtp::SeqNo; |
| 9 | use crate::rtp_::{Mid, Rid}; |
| 10 | use crate::{io::Protocol, rtp_::MidRid}; |
| 11 | |
| 12 | pub(crate) struct Stats { |
| 13 | last_now: Option<Instant>, |
| 14 | events: VecDeque<StatsEvent>, |
| 15 | interval: Duration, |
| 16 | } |
| 17 | |
| 18 | pub(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 | |
| 33 | impl 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)] |
| 55 | pub(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)] |
| 65 | pub 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. |
| 90 | pub 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. |
| 127 | pub 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)] |
| 136 | pub 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)] |
| 172 | pub 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)] |
| 185 | pub 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 | |
| 213 | impl 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)] |
| 259 | pub 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 | |
| 266 | impl 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 | } |