File
Blob: firmware/vendor/str0m/src/streams/send_stats.rs
| 1 | use std::time::{Duration, Instant}; |
| 2 | |
| 3 | use crate::rtp::SeqNo; |
| 4 | use crate::rtp_::{ReceptionReport, extend_u32}; |
| 5 | use crate::stats::{MediaEgressStats, RemoteIngressStats, StatsSnapshot}; |
| 6 | use crate::util::value_history::ValueHistory; |
| 7 | use crate::util::{InstantExt, calculate_rtt}; |
| 8 | |
| 9 | use super::MidRid; |
| 10 | |
| 11 | /// Holder of stats. |
| 12 | #[derive(Debug)] |
| 13 | pub(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 | |
| 48 | impl 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)] |
| 152 | enum Losses { |
| 153 | Disabled, |
| 154 | Enabled(Vec<(u64, f32)>), |
| 155 | } |
| 156 | |
| 157 | impl 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 | } |