File
Blob: firmware/vendor/str0m/src/bwe/delay/control.rs
| 1 | use std::collections::VecDeque; |
| 2 | use std::time::{Duration, Instant}; |
| 3 | |
| 4 | use super::super::macros::{log_bitrate_estimate, log_delay_variation}; |
| 5 | use super::super::{AckedPacket, BandwidthUsage}; |
| 6 | use super::arrival_group::ArrivalGroupAccumulator; |
| 7 | use super::rate_control::RateControl; |
| 8 | use super::trendline::TrendlineEstimator; |
| 9 | use crate::rtp_::Bitrate; |
| 10 | use crate::util::{MovingAverage, already_happened}; |
| 11 | |
| 12 | const MAX_RTT_HISTORY_WINDOW: usize = 32; |
| 13 | const UPDATE_INTERVAL: Duration = Duration::from_millis(25); |
| 14 | /// The maximum time we keep updating our estimate without receiving a TWCC report. |
| 15 | const MAX_TWCC_GAP: Duration = Duration::from_millis(500); |
| 16 | /// RFC 6298: Exponentially Weighted Moving Average smoothing factor for RTT (alpha = 1/8) |
| 17 | const RTT_SMOOTHING_FACTOR: f64 = 0.125; |
| 18 | |
| 19 | /// Delay controller for googcc inspired BWE. |
| 20 | /// |
| 21 | /// This controller attempts to estimate the available send bandwidth by looking at the variations |
| 22 | /// in packet arrival times for groups of packets sent together. Broadly, if the delay variation is |
| 23 | /// increasing this indicates overuse. |
| 24 | pub struct DelayController { |
| 25 | arrival_group_accumulator: ArrivalGroupAccumulator, |
| 26 | trendline_estimator: TrendlineEstimator, |
| 27 | rate_control: RateControl, |
| 28 | /// Last estimate produced, unlike [`next_estimate`] this will always have a value after the |
| 29 | /// first estimate. |
| 30 | last_estimate: Option<Bitrate>, |
| 31 | /// Smoothed RTT using EWMA (RFC 6298, alpha = 1/8). |
| 32 | smoothed_rtt: MovingAverage, |
| 33 | /// History of the max RTT derived for each TWCC report (kept for fallback). |
| 34 | max_rtt_history: VecDeque<Duration>, |
| 35 | |
| 36 | /// The next time we should poll. |
| 37 | next_timeout: Instant, |
| 38 | /// The last time we ingested a TWCC report. |
| 39 | last_twcc_report: Instant, |
| 40 | } |
| 41 | |
| 42 | impl DelayController { |
| 43 | pub fn new(initial_bitrate: Bitrate) -> Self { |
| 44 | Self { |
| 45 | arrival_group_accumulator: ArrivalGroupAccumulator::default(), |
| 46 | trendline_estimator: TrendlineEstimator::new(20), |
| 47 | rate_control: RateControl::new(initial_bitrate, Bitrate::kbps(40), Bitrate::gbps(10)), |
| 48 | last_estimate: Some(initial_bitrate), |
| 49 | smoothed_rtt: MovingAverage::new(RTT_SMOOTHING_FACTOR), |
| 50 | max_rtt_history: VecDeque::default(), |
| 51 | next_timeout: already_happened(), |
| 52 | last_twcc_report: already_happened(), |
| 53 | } |
| 54 | } |
| 55 | |
| 56 | /// Record a packet from a TWCC report. |
| 57 | pub fn update( |
| 58 | &mut self, |
| 59 | acked: &[AckedPacket], |
| 60 | acked_bitrate: Option<Bitrate>, |
| 61 | probe_bitrate: Option<Bitrate>, |
| 62 | now: Instant, |
| 63 | ) -> Option<Bitrate> { |
| 64 | let mut max_rtt = None; |
| 65 | |
| 66 | for acked_packet in acked { |
| 67 | max_rtt = max_rtt.max(Some(acked_packet.rtt())); |
| 68 | if let Some(delay_variation) = self |
| 69 | .arrival_group_accumulator |
| 70 | .accumulate_packet(acked_packet) |
| 71 | { |
| 72 | log_delay_variation!(delay_variation.arrival_delta); |
| 73 | |
| 74 | // Got a new delay variation, add it to the trendline. |
| 75 | // |
| 76 | // IMPORTANT: Match WebRTC's TrendlineEstimator time base. |
| 77 | // WebRTC calls Detect/UpdateThreshold with `arrival_time_ms` (remote receive time), |
| 78 | // not the local "time we processed this feedback". Using the remote receive time |
| 79 | // avoids threshold adaptation artifacts when many deltas are processed in one |
| 80 | // feedback batch (e.g. TWCC reports). |
| 81 | // |
| 82 | // Note: We use remote timestamps for relative timing only (computing time deltas |
| 83 | // between packets). Clock skew doesn't matter since we're measuring trends in |
| 84 | // delay variations, not absolute times. |
| 85 | self.trendline_estimator |
| 86 | .add_delay_observation(delay_variation, delay_variation.last_remote_recv_time); |
| 87 | } |
| 88 | } |
| 89 | |
| 90 | if let Some(rtt) = max_rtt { |
| 91 | self.update_rtt(rtt); |
| 92 | } |
| 93 | |
| 94 | let new_hypothesis = self.trendline_estimator.hypothesis(); |
| 95 | |
| 96 | self.update_estimate( |
| 97 | new_hypothesis, |
| 98 | acked_bitrate, |
| 99 | probe_bitrate, |
| 100 | self.get_smoothed_rtt(), |
| 101 | now, |
| 102 | ); |
| 103 | self.last_twcc_report = now; |
| 104 | |
| 105 | self.last_estimate |
| 106 | } |
| 107 | |
| 108 | pub fn poll_timeout(&self) -> Instant { |
| 109 | self.next_timeout |
| 110 | } |
| 111 | |
| 112 | pub fn handle_timeout(&mut self, acked_bitrate: Option<Bitrate>, now: Instant) { |
| 113 | if !self.trendline_hypothesis_valid(now) { |
| 114 | // We haven't received a TWCC report in a while. The trendline hypothesis can |
| 115 | // no longer be considered valid. We need another TWCC report before we can update |
| 116 | // estimates. |
| 117 | let next_timeout_in = self |
| 118 | .get_smoothed_rtt() |
| 119 | .unwrap_or(MAX_TWCC_GAP) |
| 120 | .min(UPDATE_INTERVAL); |
| 121 | |
| 122 | // Set this even if we didn't update, otherwise we get stuck in a poll -> handle loop |
| 123 | // that starves the run loop. |
| 124 | self.next_timeout = now + next_timeout_in; |
| 125 | return; |
| 126 | } |
| 127 | |
| 128 | self.update_estimate( |
| 129 | self.trendline_estimator.hypothesis(), |
| 130 | acked_bitrate, |
| 131 | None, |
| 132 | self.get_smoothed_rtt(), |
| 133 | now, |
| 134 | ); |
| 135 | } |
| 136 | |
| 137 | /// Get the latest estimate. |
| 138 | pub fn last_estimate(&self) -> Option<Bitrate> { |
| 139 | self.last_estimate |
| 140 | } |
| 141 | |
| 142 | /// Whether the delay-based detector currently signals overuse. |
| 143 | /// |
| 144 | /// This is useful for gating behaviors (like probing) that would otherwise |
| 145 | /// re-excite the system while we're already congested. |
| 146 | pub fn is_overusing(&self) -> bool { |
| 147 | self.trendline_estimator.hypothesis() == BandwidthUsage::Overuse |
| 148 | } |
| 149 | |
| 150 | /// Update smoothed RTT using EWMA (RFC 6298, alpha = 1/8). |
| 151 | fn update_rtt(&mut self, rtt: Duration) { |
| 152 | // Keep history as fallback in case smoothed RTT is not yet available |
| 153 | while self.max_rtt_history.len() >= MAX_RTT_HISTORY_WINDOW { |
| 154 | self.max_rtt_history.pop_front(); |
| 155 | } |
| 156 | self.max_rtt_history.push_back(rtt); |
| 157 | |
| 158 | // Update smoothed RTT using EWMA: smoothed = (7/8) * smoothed + (1/8) * sample |
| 159 | self.smoothed_rtt.update(rtt.as_secs_f64()); |
| 160 | } |
| 161 | |
| 162 | /// Get the current smoothed RTT, with fallback to mean of history if not yet available. |
| 163 | fn get_smoothed_rtt(&self) -> Option<Duration> { |
| 164 | // Try smoothed RTT first (EWMA) |
| 165 | if let Some(avg_secs) = self.smoothed_rtt.get() { |
| 166 | return Some(Duration::from_secs_f64(avg_secs)); |
| 167 | } |
| 168 | |
| 169 | // Fallback to mean of history during initialization |
| 170 | if self.max_rtt_history.is_empty() { |
| 171 | return None; |
| 172 | } |
| 173 | |
| 174 | let sum = self |
| 175 | .max_rtt_history |
| 176 | .iter() |
| 177 | .fold(Duration::ZERO, |acc, rtt| acc + *rtt); |
| 178 | Some(sum / self.max_rtt_history.len() as u32) |
| 179 | } |
| 180 | |
| 181 | fn update_estimate( |
| 182 | &mut self, |
| 183 | hypothesis: BandwidthUsage, |
| 184 | observed_bitrate: Option<Bitrate>, |
| 185 | probe_bitrate: Option<Bitrate>, |
| 186 | mean_max_rtt: Option<Duration>, |
| 187 | now: Instant, |
| 188 | ) { |
| 189 | // WebRTC's logic from delay_based_bwe.cc MaybeUpdateEstimate(): |
| 190 | // - If we have a probe result, apply it directly and skip delay-based updates |
| 191 | // - Otherwise, apply normal delay-based rate control |
| 192 | // |
| 193 | // This prevents probe results from being immediately overridden by delay-based |
| 194 | // decreases caused by the probe itself (probes cause temporary queuing delay). |
| 195 | |
| 196 | if let Some(probe_rate) = probe_bitrate { |
| 197 | // Apply probe result directly, bypassing delay-based updates |
| 198 | self.rate_control.set_probe_result(probe_rate, now); |
| 199 | let estimated_rate = self.rate_control.estimated_bitrate(); |
| 200 | log_bitrate_estimate!(estimated_rate.as_f64()); |
| 201 | self.last_estimate = Some(estimated_rate); |
| 202 | } else if let Some(observed_bitrate) = observed_bitrate { |
| 203 | // No probe result, apply normal delay-based rate control |
| 204 | self.rate_control |
| 205 | .update(hypothesis.into(), observed_bitrate, mean_max_rtt, now); |
| 206 | let estimated_rate = self.rate_control.estimated_bitrate(); |
| 207 | |
| 208 | log_bitrate_estimate!(estimated_rate.as_f64()); |
| 209 | self.last_estimate = Some(estimated_rate); |
| 210 | } |
| 211 | |
| 212 | // Set this even if we didn't update, otherwise we get stuck in a poll -> handle loop |
| 213 | // that starves the run loop. |
| 214 | self.next_timeout = now + UPDATE_INTERVAL; |
| 215 | } |
| 216 | |
| 217 | /// Whether the current trendline hypothesis is valid i.e. not too old. |
| 218 | fn trendline_hypothesis_valid(&self, now: Instant) -> bool { |
| 219 | now.duration_since(self.last_twcc_report) |
| 220 | <= self |
| 221 | .get_smoothed_rtt() |
| 222 | .map(|rtt| rtt * 2) |
| 223 | .unwrap_or(MAX_TWCC_GAP) |
| 224 | .min(UPDATE_INTERVAL * 2) |
| 225 | } |
| 226 | } |