Skip to content
File

Blob: firmware/vendor/str0m/src/bwe/delay/control.rs

rust227 lines
1use std::collections::VecDeque;
2use std::time::{Duration, Instant};
3 
4use super::super::macros::{log_bitrate_estimate, log_delay_variation};
5use super::super::{AckedPacket, BandwidthUsage};
6use super::arrival_group::ArrivalGroupAccumulator;
7use super::rate_control::RateControl;
8use super::trendline::TrendlineEstimator;
9use crate::rtp_::Bitrate;
10use crate::util::{MovingAverage, already_happened};
11 
12const MAX_RTT_HISTORY_WINDOW: usize = 32;
13const UPDATE_INTERVAL: Duration = Duration::from_millis(25);
14/// The maximum time we keep updating our estimate without receiving a TWCC report.
15const MAX_TWCC_GAP: Duration = Duration::from_millis(500);
16/// RFC 6298: Exponentially Weighted Moving Average smoothing factor for RTT (alpha = 1/8)
17const 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.
24pub 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 
42impl 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}