Skip to content
File

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

rust455 lines
1use std::collections::VecDeque;
2use std::ops::RangeInclusive;
3use std::time::{Duration, Instant};
4 
5use super::super::macros::{log_trendline_estimate, log_trendline_modified_trend};
6 
7use super::super::BandwidthUsage;
8use super::super::time::{TimeDelta, Timestamp};
9use super::arrival_group::InterGroupDelayDelta;
10 
11const SMOOTHING_COEF: f64 = 0.9;
12const OVER_USE_THRESHOLD_DEFAULT_MS: f64 = 12.5;
13const OVER_USE_TIME_THRESHOLD: Duration = Duration::from_millis(10);
14const MAX_ADOPT_OFFSET_MS: f64 = 15.0;
15const THRESHOLD_GAIN: f64 = 4.0;
16 
17const K_UP: f64 = 0.0087;
18const K_DOWN: f64 = 0.039;
19 
20const DELAY_COUNT_RANGE: RangeInclusive<usize> = 60..=1000;
21 
22pub struct TrendlineEstimator {
23 /// The window size in packets
24 window_size: usize,
25 
26 /// The first instant we saw, used as zero point.
27 zero_time: Option<Timestamp>,
28 
29 /// The history of observed delay variations.
30 history: VecDeque<Timing>,
31 
32 /// The total number of observed delay variations.
33 num_delay_variations: usize,
34 
35 /// Accumulated delay.
36 accumulated_delay: f64,
37 
38 /// Last smoothed delay.
39 smoothed_delay: f64,
40 
41 /// The adaptive delay threshold.
42 delay_threshold: f64,
43 
44 /// Previous trend
45 previous_trend: f64,
46 
47 /// If we are overusing, this contains data about the overuse.
48 overuse: Option<Overuse>,
49 
50 /// The last time we updated the adaptive threshold.
51 last_threshold_update: Option<Instant>,
52 
53 /// Our current hypothesis about the bandwidth usage.
54 hypothesis: BandwidthUsage,
55}
56 
57impl TrendlineEstimator {
58 pub fn new(window_size: usize) -> Self {
59 Self {
60 window_size,
61 zero_time: None,
62 history: VecDeque::default(),
63 num_delay_variations: 0,
64 accumulated_delay: 0.0,
65 smoothed_delay: 0.0,
66 delay_threshold: OVER_USE_THRESHOLD_DEFAULT_MS,
67 previous_trend: 0.0,
68 overuse: None,
69 last_threshold_update: None,
70 hypothesis: BandwidthUsage::Normal,
71 }
72 }
73 
74 pub fn add_delay_observation(&mut self, delay_variation: InterGroupDelayDelta, now: Instant) {
75 if self.history.is_empty() {
76 self.do_add_to_history(delay_variation, now);
77 return;
78 }
79 
80 self.do_add_to_history(delay_variation, now);
81 while self.history.len() > self.window_size {
82 let _ = self.history.pop_front();
83 }
84 
85 if self.history.len() == self.window_size {
86 assert!(
87 self.history
88 .iter()
89 .zip(self.history.iter().skip(1))
90 .fold(true, |acc, (a, b)| {
91 acc && a.remote_recv_time_ms <= b.remote_recv_time_ms
92 }),
93 "Out of order history {:?}",
94 self.history
95 );
96 self.update_trendline(delay_variation, now);
97 }
98 }
99 
100 pub fn hypothesis(&self) -> BandwidthUsage {
101 self.hypothesis
102 }
103 
104 fn do_add_to_history(&mut self, variation: InterGroupDelayDelta, now: Instant) {
105 let last_remote_recv_time = Timestamp::from(variation.last_remote_recv_time);
106 
107 let zero_time = *self.zero_time.get_or_insert(last_remote_recv_time);
108 
109 let delay_delta = variation.arrival_delta.as_secs_f64() * 1000.0
110 - variation.send_delta.as_secs_f64() * 1000.0;
111 
112 self.num_delay_variations += 1;
113 self.num_delay_variations = self.num_delay_variations.min(*DELAY_COUNT_RANGE.end());
114 self.accumulated_delay += delay_delta;
115 self.smoothed_delay =
116 self.smoothed_delay * SMOOTHING_COEF + (1.0 - SMOOTHING_COEF) * self.accumulated_delay;
117 
118 let remote_recv_time = last_remote_recv_time - zero_time;
119 let timing = Timing {
120 at: now,
121 remote_recv_time_ms: remote_recv_time.as_secs_f64() * 1000.0,
122 smoothed_delay_ms: self.smoothed_delay,
123 };
124 
125 let pos = self
126 .history
127 .iter()
128 .rev()
129 .position(|p| p.remote_recv_time_ms <= timing.remote_recv_time_ms)
130 .unwrap_or(self.history.len());
131 
132 // we expect pos to be 0 more often than not
133 self.history.insert(self.history.len() - pos, timing);
134 }
135 
136 fn update_trendline(&mut self, variation: InterGroupDelayDelta, now: Instant) -> Option<()> {
137 let trend = self.linear_fit().unwrap_or(self.previous_trend);
138 trace!("Computed trend {:?}", trend);
139 log_trendline_estimate!(trend);
140 
141 self.detect(trend, variation, now);
142 
143 Some(())
144 }
145 
146 fn linear_fit(&self) -> Option<f64> {
147 // Simple linear regression to compute slope.
148 assert!(self.history.len() > 2);
149 
150 let (sum_x, sum_y) = self.history.iter().fold((0.0, 0.0), |acc, t| {
151 (acc.0 + t.remote_recv_time_ms, acc.1 + t.smoothed_delay_ms)
152 });
153 
154 let avg_x = sum_x / self.history.len() as f64;
155 let avg_y = sum_y / self.history.len() as f64;
156 
157 let (numerator, denomenator) = self.history.iter().fold((0.0, 0.0), |acc, t| {
158 let x = t.remote_recv_time_ms;
159 let y = t.smoothed_delay_ms;
160 
161 let numerator = acc.0 + (x - avg_x) * (y - avg_y);
162 let denomenator = acc.1 + (x - avg_x).powi(2);
163 
164 (numerator, denomenator)
165 });
166 
167 if denomenator == 0.0 {
168 return None;
169 }
170 
171 Some(numerator / denomenator)
172 }
173 
174 fn detect(&mut self, trend: f64, variation: InterGroupDelayDelta, now: Instant) {
175 if self.num_delay_variations < 2 {
176 self.update_hypothesis(BandwidthUsage::Normal);
177 return;
178 }
179 
180 let modified_trend = self.num_delay_variations.min(*DELAY_COUNT_RANGE.start()) as f64
181 * trend
182 * THRESHOLD_GAIN;
183 
184 log_trendline_modified_trend!(modified_trend, self.delay_threshold);
185 
186 if modified_trend > self.delay_threshold {
187 let overuse = match &mut self.overuse {
188 Some(o) => {
189 o.time_overusing += variation.send_delta;
190 o
191 }
192 None => {
193 let new_overuse = Overuse {
194 count: 0,
195 // Initialize the timer. Assume that we've been
196 // over-using half of the time since the previous
197 // sample.
198 time_overusing: variation.send_delta / 2,
199 };
200 self.overuse = Some(new_overuse);
201 
202 self.overuse.as_mut().unwrap()
203 }
204 };
205 
206 overuse.count += 1;
207 trace!(
208 timeoverusing = ?overuse.time_overusing,
209 trend,
210 previous_trend = self.previous_trend,
211 "Trendline Estimator: Maybe overusing"
212 );
213 
214 if overuse.time_overusing > OVER_USE_TIME_THRESHOLD
215 && overuse.count > 1
216 && trend > self.previous_trend
217 {
218 self.overuse = None;
219 
220 self.update_hypothesis(BandwidthUsage::Overuse);
221 }
222 } else if modified_trend < -self.delay_threshold {
223 self.overuse = None;
224 self.update_hypothesis(BandwidthUsage::Underuse);
225 } else {
226 self.overuse = None;
227 self.update_hypothesis(BandwidthUsage::Normal);
228 }
229 
230 self.previous_trend = trend;
231 self.update_threshold(modified_trend, now);
232 }
233 
234 fn update_threshold(&mut self, modified_trend: f64, now: Instant) {
235 if self.last_threshold_update.is_none() {
236 self.last_threshold_update = Some(now);
237 }
238 let abs_modified_trend = modified_trend.abs();
239 
240 if abs_modified_trend > self.delay_threshold + MAX_ADOPT_OFFSET_MS {
241 // Avoid adapting the threshold to big latency spikes, caused e.g.,
242 // by a sudden capacity drop.
243 self.last_threshold_update = Some(now);
244 return;
245 }
246 
247 let k = if abs_modified_trend < self.delay_threshold {
248 K_DOWN
249 } else {
250 K_UP
251 };
252 let time_delta = now
253 .saturating_duration_since(
254 self.last_threshold_update
255 .expect("last_threshold_update must have been set"),
256 )
257 .as_millis() as f64;
258 self.delay_threshold +=
259 k * (abs_modified_trend - self.delay_threshold) * time_delta.min(100.0);
260 self.last_threshold_update = Some(now);
261 self.delay_threshold = self.delay_threshold.clamp(6.0, 600.0);
262 
263 trace!(
264 "Adaptive delay variation threshold changed to: {}",
265 self.delay_threshold
266 );
267 }
268 
269 fn update_hypothesis(&mut self, new_hypothesis: BandwidthUsage) {
270 if self.hypothesis == new_hypothesis {
271 return;
272 }
273 
274 trace!("TrendLineEstimator: Setting hypothesis to {new_hypothesis}");
275 self.hypothesis = new_hypothesis;
276 }
277}
278 
279#[derive(Debug)]
280struct Timing {
281 #[allow(unused)]
282 at: Instant,
283 remote_recv_time_ms: f64,
284 smoothed_delay_ms: f64,
285}
286 
287struct Overuse {
288 count: usize,
289 time_overusing: TimeDelta,
290}
291 
292#[cfg(test)]
293mod test {
294 use super::*;
295 use std::time::{Duration, Instant};
296 
297 #[test]
298 fn test_window_size_limit() {
299 let now = Instant::now();
300 let remote_recv_time_base = Instant::now();
301 let mut estimator = TrendlineEstimator::new(20);
302 
303 estimator.add_delay_observation(
304 delay_variation(duration_ms(1), duration_ms(1), remote_recv_time_base),
305 now,
306 );
307 
308 for _ in 0..25 {
309 estimator.add_delay_observation(
310 delay_variation(
311 duration_ms(11),
312 duration_ms(1),
313 remote_recv_time_base + duration_ms(350),
314 ),
315 now + duration_ms(500),
316 );
317 }
318 
319 assert_eq!(estimator.history.len(), 20);
320 }
321 
322 #[test]
323 fn test_overuse() {
324 let now = Instant::now();
325 let remote_recv_time_base = Instant::now();
326 let mut estimator = TrendlineEstimator::new(20);
327 
328 for g in 0..5 {
329 for i in 0..5 {
330 estimator.add_delay_observation(
331 delay_variation(
332 duration_ms(1),
333 duration_ms(1),
334 remote_recv_time_base + Duration::from_micros(5_000 * g + i * 40),
335 ),
336 now + duration_ms(g * 100),
337 );
338 }
339 }
340 
341 assert_eq!(estimator.hypothesis(), BandwidthUsage::Normal);
342 assert_eq!(estimator.history.len(), 20);
343 
344 estimator.add_delay_observation(
345 delay_variation(
346 duration_ms(17),
347 duration_ms(5),
348 remote_recv_time_base + Duration::from_micros(25_000),
349 ),
350 now + duration_ms(600),
351 );
352 assert_eq!(
353 estimator.hypothesis(),
354 BandwidthUsage::Normal,
355 "After getting an initial increasing delay the hypothesis should remain at normal"
356 );
357 
358 estimator.add_delay_observation(
359 delay_variation(
360 duration_ms(18),
361 duration_ms(5),
362 remote_recv_time_base + Duration::from_micros(25_140),
363 ),
364 now + duration_ms(600),
365 );
366 assert_eq!(
367 estimator.hypothesis(),
368 BandwidthUsage::Normal,
369 "After getting an a second increasing delay the hypothesis should remain at normal \
370 because we the time overusing threshold hasn't been reached yet"
371 );
372 
373 estimator.add_delay_observation(
374 delay_variation(
375 duration_ms(22),
376 duration_ms(8),
377 remote_recv_time_base + Duration::from_micros(25_250),
378 ),
379 now + duration_ms(600),
380 );
381 assert_eq!(
382 estimator.hypothesis(),
383 BandwidthUsage::Overuse,
384 "After getting a third increasing delay the hypothesis should move to over because \
385 we have been overusing for more than 10ms"
386 );
387 }
388 
389 fn duration_ms(ms: u64) -> Duration {
390 Duration::from_millis(ms)
391 }
392 
393 fn delay_variation(
394 recv_delta: Duration,
395 send_delta: Duration,
396 last_remote_recv_time: Instant,
397 ) -> InterGroupDelayDelta {
398 InterGroupDelayDelta {
399 send_delta: send_delta.into(),
400 arrival_delta: recv_delta.into(),
401 last_remote_recv_time,
402 }
403 }
404 
405 #[test]
406 fn test_history_insertion_regression_out_of_order() {
407 // Test for algesten/str0m#698
408 //
409 // This test reproduces the scenario where a sample with a remote receive time
410 // earlier than all existing history entries is added.
411 // Prior to the fix, this appended the smallest (negative) element to the end
412 // causing the monotonicity assert to fail when the window is full.
413 
414 let now = Instant::now();
415 let zero_time_base = Instant::now();
416 let mut estimator = TrendlineEstimator::new(8);
417 
418 // Add a few normal monotonically increasing values
419 for offset_ms in [100, 200, 300, 400] {
420 estimator.add_delay_observation(
421 delay_variation(
422 duration_ms(10),
423 duration_ms(1),
424 zero_time_base + duration_ms(offset_ms),
425 ),
426 now,
427 );
428 }
429 
430 // Add some out-of-order samples that are earlier than the zero_time
431 // (packet reordering is fun).
432 for offset_ms in [50, 100, 25, 300] {
433 estimator.add_delay_observation(
434 delay_variation(
435 duration_ms(10),
436 duration_ms(1),
437 zero_time_base - duration_ms(offset_ms),
438 ),
439 now,
440 );
441 }
442 
443 // History should remain sorted.
444 let ordered = estimator
445 .history
446 .iter()
447 .zip(estimator.history.iter().skip(1))
448 .all(|(a, b)| a.remote_recv_time_ms <= b.remote_recv_time_ms);
449 assert!(
450 ordered,
451 "history should remain sorted by remote_recv_time_ms"
452 );
453 }
454}