File
Blob: firmware/vendor/str0m/src/bwe/probe/control.rs
| 1 | //! Bandwidth probing controller - decides when and how to probe network capacity. |
| 2 | //! |
| 3 | //! This module implements WebRTC's `ProbeController` state machine for discovering available |
| 4 | //! bandwidth through intentional bursts of packets at rates higher than current estimates. |
| 5 | |
| 6 | use std::collections::VecDeque; |
| 7 | use std::time::{Duration, Instant}; |
| 8 | |
| 9 | use super::{ProbeClusterConfig, ProbeKind}; |
| 10 | use crate::rtp_::{Bitrate, TwccClusterId}; |
| 11 | use crate::util::{already_happened, not_happening}; |
| 12 | |
| 13 | // Port notes: |
| 14 | // This module ports WebRTC's `ProbeController` behavior from: |
| 15 | // `webrtc/modules/congestion_controller/goog_cc/probe_controller.cc` |
| 16 | // |
| 17 | // Key integration difference: WebRTC returns vectors of probe clusters, while str0m |
| 18 | // returns a single `ProbeClusterConfig` per `handle_timeout()` call. Configs are queued |
| 19 | // internally and `poll_timeout()` returns `already_happened()` until the queue is drained. |
| 20 | |
| 21 | /// WebRTC: `kMaxWaitingTimeForProbingResult`. |
| 22 | const MAX_WAITING_TIME_FOR_PROBING_RESULT: Duration = Duration::from_secs(1); |
| 23 | |
| 24 | /// WebRTC: `kBitrateDropThreshold`, `kBitrateDropTimeout`, `kProbeFractionAfterDrop`, |
| 25 | /// `kProbeUncertainty`, `kAlrEndedTimeout`, `kMinTimeBetweenAlrProbes`. |
| 26 | const BITRATE_DROP_THRESHOLD: f64 = 0.66; |
| 27 | const BITRATE_DROP_TIMEOUT: Duration = Duration::from_secs(5); |
| 28 | const PROBE_FRACTION_AFTER_DROP: f64 = 0.85; |
| 29 | const PROBE_UNCERTAINTY: f64 = 0.05; |
| 30 | const ALR_ENDED_TIMEOUT: Duration = Duration::from_secs(3); |
| 31 | const MIN_TIME_BETWEEN_ALR_PROBES: Duration = Duration::from_secs(5); |
| 32 | |
| 33 | /// WebRTC: inline `* 2` in probe_controller.cc InitiateProbing(). |
| 34 | /// Allows probing up to 2x max_bitrate to account for bursty streams. |
| 35 | const MAX_PROBE_BITRATE_FACTOR: f64 = 2.0; |
| 36 | |
| 37 | /// Minimum time between stagnant periodic probes to avoid excessive probing when at capacity. |
| 38 | const MIN_TIME_BETWEEN_STAGNANT_PROBES: Duration = Duration::from_secs(15); |
| 39 | |
| 40 | /// Threshold for considering an estimate change significant (5%). |
| 41 | const ESTIMATE_CHANGE_THRESHOLD: f64 = 0.05; |
| 42 | |
| 43 | /// Probe rate scale for stagnation probes (2× current estimate). |
| 44 | const STAGNANT_PROBE_SCALE: f64 = 2.0; |
| 45 | |
| 46 | /// WebRTC's `BandwidthLimitedCause` (subset used by probing gating). |
| 47 | #[derive(Debug, Clone, Copy, PartialEq, Eq)] |
| 48 | pub enum BandwidthLimitedCause { |
| 49 | LossLimitedBweIncreasing, |
| 50 | LossLimitedBwe, |
| 51 | DelayBasedLimited, |
| 52 | DelayBasedLimitedDelayIncreased, |
| 53 | } |
| 54 | |
| 55 | pub struct ProbeControl { |
| 56 | config: Config, |
| 57 | next_timeout: Instant, |
| 58 | enabled: bool, |
| 59 | |
| 60 | desired_bitrate: Option<Bitrate>, |
| 61 | prev_desired: Option<Bitrate>, |
| 62 | |
| 63 | last_estimate: Option<Bitrate>, |
| 64 | last_estimate_change: Option<Instant>, |
| 65 | last_cause: BandwidthLimitedCause, |
| 66 | |
| 67 | prev_estimate: Option<Bitrate>, |
| 68 | |
| 69 | alr_start: Option<Instant>, |
| 70 | alr_stop: Option<Instant>, |
| 71 | |
| 72 | last_probe: Option<LastProbe>, |
| 73 | |
| 74 | large_drop: Option<LargeDrop>, |
| 75 | |
| 76 | last_stagnant: Option<Instant>, |
| 77 | |
| 78 | next_cluster_id: TwccClusterId, |
| 79 | pending: VecDeque<ProbeClusterConfig>, |
| 80 | |
| 81 | scheduled_exponential: Option<Instant>, |
| 82 | scheduled_periodic_alr: Option<Instant>, |
| 83 | scheduled_stagnant: Option<Instant>, |
| 84 | } |
| 85 | |
| 86 | #[derive(Debug, Clone, Copy, PartialEq)] |
| 87 | struct LastProbe { |
| 88 | when: Instant, |
| 89 | kind: ProbeKind, |
| 90 | further: Bitrate, |
| 91 | was_estimate: Option<Bitrate>, |
| 92 | } |
| 93 | |
| 94 | struct LargeDrop { |
| 95 | when: Instant, |
| 96 | bitrate_before: Bitrate, |
| 97 | } |
| 98 | |
| 99 | impl Default for ProbeControl { |
| 100 | fn default() -> Self { |
| 101 | Self { |
| 102 | config: Config::default(), |
| 103 | enabled: false, |
| 104 | next_timeout: not_happening(), |
| 105 | desired_bitrate: None, |
| 106 | prev_desired: None, |
| 107 | last_estimate: None, |
| 108 | last_estimate_change: None, |
| 109 | last_cause: BandwidthLimitedCause::DelayBasedLimited, |
| 110 | prev_estimate: None, |
| 111 | alr_start: None, |
| 112 | alr_stop: None, |
| 113 | next_cluster_id: 0.into(), |
| 114 | last_probe: None, |
| 115 | large_drop: None, |
| 116 | last_stagnant: None, |
| 117 | pending: VecDeque::new(), |
| 118 | scheduled_exponential: None, |
| 119 | scheduled_periodic_alr: None, |
| 120 | scheduled_stagnant: None, |
| 121 | } |
| 122 | } |
| 123 | } |
| 124 | |
| 125 | impl ProbeControl { |
| 126 | pub fn new() -> Self { |
| 127 | Self::default() |
| 128 | } |
| 129 | |
| 130 | pub fn enable(&mut self, v: bool) { |
| 131 | if !self.enabled && v { |
| 132 | self.enabled = true; |
| 133 | self.request_immediate(); |
| 134 | } else if self.enabled && !v { |
| 135 | self.enabled = false; |
| 136 | self.pending.clear(); |
| 137 | self.last_estimate = None; |
| 138 | self.desired_bitrate = None; |
| 139 | self.last_estimate_change = None; |
| 140 | self.last_stagnant = None; |
| 141 | self.last_probe = None; |
| 142 | self.prev_estimate = None; |
| 143 | self.scheduled_exponential = None; |
| 144 | self.scheduled_periodic_alr = None; |
| 145 | self.scheduled_stagnant = None; |
| 146 | self.next_timeout = not_happening(); |
| 147 | } |
| 148 | } |
| 149 | |
| 150 | pub fn set_desired_bitrate(&mut self, v: Bitrate) { |
| 151 | // Don't accept Bitrate::ZERO as first ever value. |
| 152 | if self.desired_bitrate.is_none() && v.is_zero() { |
| 153 | return; |
| 154 | } |
| 155 | self.desired_bitrate = Some(v); |
| 156 | self.request_immediate(); |
| 157 | } |
| 158 | |
| 159 | pub fn set_estimated_bitrate(&mut self, v: Bitrate, cause: BandwidthLimitedCause) { |
| 160 | // Don't accept Bitrate::ZERO as first ever value. |
| 161 | if self.last_estimate.is_none() && v.is_zero() { |
| 162 | return; |
| 163 | } |
| 164 | |
| 165 | // Check if estimate changed significantly (>5%) or cause changed. |
| 166 | let dominated_by_last = self.last_estimate.is_some_and(|last| { |
| 167 | let upper = last * (1.0 + ESTIMATE_CHANGE_THRESHOLD); |
| 168 | let lower = last * (1.0 - ESTIMATE_CHANGE_THRESHOLD); |
| 169 | v <= upper && v >= lower |
| 170 | }); |
| 171 | |
| 172 | if dominated_by_last && self.last_cause == cause { |
| 173 | return; |
| 174 | } |
| 175 | |
| 176 | self.last_estimate = Some(v); |
| 177 | self.last_cause = cause; |
| 178 | self.request_immediate(); |
| 179 | } |
| 180 | |
| 181 | pub fn set_alr_start_time(&mut self, t: Instant) { |
| 182 | if self.alr_start.is_some() { |
| 183 | return; |
| 184 | } |
| 185 | self.alr_start = Some(t); |
| 186 | self.alr_stop = None; |
| 187 | self.request_immediate(); |
| 188 | } |
| 189 | |
| 190 | pub fn set_alr_stop_time(&mut self, t: Instant) { |
| 191 | if self.alr_start.is_none() || self.alr_stop.is_some() { |
| 192 | return; |
| 193 | } |
| 194 | self.alr_start = None; |
| 195 | self.alr_stop = Some(t); |
| 196 | self.request_immediate(); |
| 197 | } |
| 198 | |
| 199 | fn request_immediate(&mut self) { |
| 200 | self.next_timeout = already_happened(); |
| 201 | self.scheduled_exponential = None; |
| 202 | self.scheduled_periodic_alr = None; |
| 203 | self.scheduled_stagnant = None; |
| 204 | } |
| 205 | |
| 206 | pub fn poll_timeout(&self) -> Instant { |
| 207 | self.next_timeout |
| 208 | } |
| 209 | |
| 210 | pub fn handle_timeout(&mut self, now: Instant) -> Option<ProbeClusterConfig> { |
| 211 | // Spurious call before timeout is due - ignore. |
| 212 | if now < self.next_timeout { |
| 213 | return None; |
| 214 | } |
| 215 | |
| 216 | // Timeout fired - reset to not_happening until we compute the next one. |
| 217 | self.next_timeout = not_happening(); |
| 218 | |
| 219 | // Probing is disabled until first packet sent and padding queue exists. |
| 220 | if !self.enabled { |
| 221 | return None; |
| 222 | } |
| 223 | |
| 224 | // We need to have both desired AND last_estimate set to |
| 225 | // start considering probing. |
| 226 | let desired = self.desired_bitrate?; |
| 227 | let estimate = self.last_estimate?; |
| 228 | |
| 229 | // Return pending probes first. |
| 230 | if let Some(config) = self.pending.pop_front() { |
| 231 | // Schedule another. |
| 232 | self.request_immediate(); |
| 233 | return Some(config); |
| 234 | } |
| 235 | |
| 236 | // Can't probe in certain bandwidth-limited states. |
| 237 | if !self.can_probe(estimate) { |
| 238 | return None; |
| 239 | } |
| 240 | |
| 241 | // Try each probe type in order - only one fires per timeout. |
| 242 | let _ = self.maybe_initial(now, desired, estimate) |
| 243 | || self.maybe_exponential(now, desired, estimate) |
| 244 | || self.maybe_increase_alr(now, desired, estimate) |
| 245 | || self.maybe_large_drop(now, desired, estimate) |
| 246 | || self.maybe_periodic_alr(now, desired) |
| 247 | || self.maybe_stagnant(now, desired, estimate); |
| 248 | |
| 249 | self.update_estimate_change(now, estimate); |
| 250 | |
| 251 | // Update prev_estimate for next cycle (used by large drop and stagnation detection). |
| 252 | self.prev_estimate = Some(estimate); |
| 253 | |
| 254 | // Update timeout based on current state. |
| 255 | self.next_timeout = self.compute_next_timeout(now); |
| 256 | |
| 257 | if !self.pending.is_empty() { |
| 258 | self.request_immediate(); |
| 259 | } |
| 260 | |
| 261 | self.pending.pop_front() |
| 262 | } |
| 263 | |
| 264 | fn update_estimate_change(&mut self, now: Instant, estimate: Bitrate) { |
| 265 | // Track when estimate last changed significantly (>5%). |
| 266 | if let Some(prev) = self.prev_estimate { |
| 267 | if estimate != prev { |
| 268 | self.last_estimate_change = Some(now); |
| 269 | } |
| 270 | } |
| 271 | |
| 272 | // Initialize baseline if not set yet. |
| 273 | if self.last_estimate_change.is_none() { |
| 274 | self.last_estimate_change = Some(now); |
| 275 | } |
| 276 | } |
| 277 | |
| 278 | fn maybe_initial(&mut self, now: Instant, desired: Bitrate, estimate: Bitrate) -> bool { |
| 279 | // Initial probes only fire once at startup. |
| 280 | if self.last_probe.is_some() { |
| 281 | return false; |
| 282 | } |
| 283 | |
| 284 | // Queue 3× and 6× of estimate. |
| 285 | let p1 = estimate * self.config.first_exponential_probe_scale; |
| 286 | let p2 = estimate * self.config.second_exponential_probe_scale; |
| 287 | |
| 288 | self.queue_probe(p1, ProbeKind::Initial, desired, now); |
| 289 | self.queue_probe(p2, ProbeKind::Initial, desired, now); |
| 290 | true |
| 291 | } |
| 292 | |
| 293 | fn maybe_exponential(&mut self, now: Instant, desired: Bitrate, estimate: Bitrate) -> bool { |
| 294 | // Wait for pending probes to be dispatched first. |
| 295 | if !self.pending.is_empty() { |
| 296 | return false; |
| 297 | } |
| 298 | |
| 299 | // Need a previous probe to continue from. |
| 300 | let Some(last) = self.last_probe else { |
| 301 | return false; |
| 302 | }; |
| 303 | |
| 304 | // Estimate must exceed 70% of last probe rate to trigger further probing. |
| 305 | if estimate < last.further { |
| 306 | return false; |
| 307 | } |
| 308 | |
| 309 | let is_same = Some(estimate) == last.was_estimate; |
| 310 | let time_since = self.time_since_last_probe(now); |
| 311 | |
| 312 | // Don't re-probe at the same estimate; wait for new result or timeout. |
| 313 | if is_same && time_since < MAX_WAITING_TIME_FOR_PROBING_RESULT { |
| 314 | return false; |
| 315 | } |
| 316 | |
| 317 | let scale = self.last_cause.probe_scale(&self.config); |
| 318 | let target = estimate * scale; |
| 319 | |
| 320 | // Already probed at max rate; no point probing again. |
| 321 | let max = desired * MAX_PROBE_BITRATE_FACTOR; |
| 322 | if target >= max && last.further >= max * self.config.further_probe_threshold { |
| 323 | return false; |
| 324 | } |
| 325 | |
| 326 | self.queue_probe(target, ProbeKind::Exponential, desired, now); |
| 327 | |
| 328 | true |
| 329 | } |
| 330 | |
| 331 | fn maybe_increase_alr(&mut self, now: Instant, desired: Bitrate, estimate: Bitrate) -> bool { |
| 332 | // Don't interfere with initial probing phase. |
| 333 | if self.is_during_initial(now) { |
| 334 | return false; |
| 335 | } |
| 336 | |
| 337 | // Allocation probes only fire in ALR (application-limited region). |
| 338 | if !self.in_alr() { |
| 339 | return false; |
| 340 | } |
| 341 | |
| 342 | let prev = self.prev_desired; |
| 343 | self.prev_desired = Some(desired); |
| 344 | |
| 345 | // Need a previous desired value to compare against. |
| 346 | let Some(prev) = prev else { |
| 347 | return false; |
| 348 | }; |
| 349 | |
| 350 | // Only probe if desired increased |
| 351 | if desired <= prev { |
| 352 | return false; |
| 353 | } |
| 354 | |
| 355 | // No point probing if we already have enough bandwidth. |
| 356 | if desired <= estimate { |
| 357 | return false; |
| 358 | } |
| 359 | |
| 360 | // Allocation probes at 1× and 2× of desired, capped by 2× estimate |
| 361 | let current_bwe_limit = estimate * self.config.allocation_probe_limit_by_current_scale; |
| 362 | |
| 363 | let p1 = (desired * self.config.first_allocation_probe_scale).min(current_bwe_limit); |
| 364 | self.queue_probe(p1, ProbeKind::IncreaseAlr, desired, now); |
| 365 | |
| 366 | let p2 = desired * self.config.second_allocation_probe_scale; |
| 367 | if p2 <= current_bwe_limit && p2 > p1 { |
| 368 | self.queue_probe(p2, ProbeKind::IncreaseAlr, desired, now); |
| 369 | } |
| 370 | |
| 371 | true |
| 372 | } |
| 373 | |
| 374 | fn maybe_periodic_alr(&mut self, now: Instant, desired: Bitrate) -> bool { |
| 375 | // Don't interfere with initial probing phase. |
| 376 | if self.is_during_initial(now) { |
| 377 | return false; |
| 378 | } |
| 379 | |
| 380 | // Periodic probes only fire in ALR (application-limited region). |
| 381 | if !self.in_alr() { |
| 382 | return false; |
| 383 | } |
| 384 | |
| 385 | // Respect minimum interval between ALR probes. |
| 386 | if self.time_since_last_probe(now) < MIN_TIME_BETWEEN_ALR_PROBES { |
| 387 | return false; |
| 388 | } |
| 389 | |
| 390 | // Periodic ALR probe at 2× desired (capped by queue_probe to 2× desired anyway). |
| 391 | // Using desired rather than estimate allows discovering higher capacity when |
| 392 | // the app wants more bandwidth than currently estimated. |
| 393 | let target = desired * self.config.further_exponential_probe_scale; |
| 394 | self.queue_probe(target, ProbeKind::PeriodicAlr, desired, now); |
| 395 | true |
| 396 | } |
| 397 | |
| 398 | /// Probe when estimate has stagnated (no change for 15+ seconds) despite unmet demand. |
| 399 | /// |
| 400 | /// ## Why This Exists (str0m Addition) |
| 401 | /// |
| 402 | /// This probe type addresses a deadlock scenario in the BWE system where AIMD recovery |
| 403 | /// cannot make progress after network capacity is restored: |
| 404 | /// |
| 405 | /// **The Deadlock:** |
| 406 | /// 1. Network degrades from 5 Mbps → 1 Mbps, estimate drops to ~900 kbps |
| 407 | /// 2. Application reduces send rate to ~500 kbps (below estimate) |
| 408 | /// 3. Network recovers to 5 Mbps |
| 409 | /// 4. AIMD tries to increase but is capped at 1.5× observed throughput: |
| 410 | /// 500 kbps × 1.5 = 750 kbps maximum |
| 411 | /// 5. Sending at 500 kbps = 71% of estimate, which is above ALR threshold (65%) |
| 412 | /// 6. ALR never triggers → no periodic probing |
| 413 | /// 7. Large-drop probe requires ALR or recent ALR exit (see `maybe_large_drop`) |
| 414 | /// 8. System is stuck: estimate ~700 kbps on a 5 Mbps network |
| 415 | /// |
| 416 | /// **AIMD's 1.5× Cap (line 191 in rate_control.rs):** |
| 417 | /// `observed_bitrate * 1.5 + Bitrate::kbps(10)` |
| 418 | /// This prevents runaway growth beyond actual sending rate. It's conservative but |
| 419 | /// necessary - without it, the estimate could grow unbounded even when we're barely |
| 420 | /// sending anything. |
| 421 | /// |
| 422 | /// **ALR Detection Threshold:** |
| 423 | /// ALR triggers when sending < 65% of estimate consistently for 500ms with budget |
| 424 | /// accumulation > 80%. At 60-70% send rate, you're in the deadlock zone: too high |
| 425 | /// to trigger ALR, too low for AIMD to help much. |
| 426 | /// |
| 427 | /// **Loss Controller's 1.5× Cap:** |
| 428 | /// The loss controller also applies a 1.5× cap during recovery (line 303 in |
| 429 | /// loss_controller.rs), compounding the AIMD limitation. |
| 430 | /// |
| 431 | /// ## How This Differs from WebRTC |
| 432 | /// |
| 433 | /// WebRTC does not have stagnation-based probing. They rely on: |
| 434 | /// 1. Large-drop recovery probe (requires ALR or recent ALR exit) |
| 435 | /// 2. Rapid recovery field trial (`WebRTC-BweRapidRecoveryExperiment`) which removes |
| 436 | /// the ALR requirement from large-drop probes |
| 437 | /// |
| 438 | /// str0m adds stagnant probing as a complementary mechanism that: |
| 439 | /// - Catches deadlock regardless of whether a drop was detected |
| 440 | /// - Provides periodic escape from any stagnation scenario, not just post-drop |
| 441 | /// - Uses a conservative 15-second wait to avoid probing at convergence |
| 442 | /// - Rate-limited to once per 30 seconds to prevent oscillation |
| 443 | fn maybe_stagnant(&mut self, now: Instant, desired: Bitrate, estimate: Bitrate) -> bool { |
| 444 | // Don't interfere with initial probing phase. |
| 445 | if self.is_during_initial(now) { |
| 446 | return false; |
| 447 | } |
| 448 | |
| 449 | // Don't probe in ALR (periodic ALR handles that). |
| 450 | if self.in_alr() { |
| 451 | return false; |
| 452 | } |
| 453 | |
| 454 | let Some(last_change) = self.last_estimate_change else { |
| 455 | return false; |
| 456 | }; |
| 457 | |
| 458 | if now.saturating_duration_since(last_change) < MIN_TIME_BETWEEN_STAGNANT_PROBES { |
| 459 | return false; |
| 460 | } |
| 461 | |
| 462 | // Only if there's unmet demand. |
| 463 | if desired <= estimate { |
| 464 | return false; |
| 465 | } |
| 466 | |
| 467 | // Rate limit: at least 30 seconds between stagnation probes. |
| 468 | if let Some(last_probe) = self.last_stagnant { |
| 469 | if now.saturating_duration_since(last_probe) < MIN_TIME_BETWEEN_STAGNANT_PROBES { |
| 470 | return false; |
| 471 | } |
| 472 | } |
| 473 | |
| 474 | // Probe at 2× estimate (conservative, won't overwhelm if at capacity). |
| 475 | let probe_rate = estimate * STAGNANT_PROBE_SCALE; |
| 476 | self.queue_probe(probe_rate, ProbeKind::Stagnant, desired, now); |
| 477 | self.last_stagnant = Some(now); |
| 478 | |
| 479 | true |
| 480 | } |
| 481 | |
| 482 | fn maybe_large_drop(&mut self, now: Instant, desired: Bitrate, estimate: Bitrate) -> bool { |
| 483 | // Don't interfere with initial probing phase. |
| 484 | if self.is_during_initial(now) { |
| 485 | return false; |
| 486 | } |
| 487 | |
| 488 | // Detect large drops: estimate fell below 66% of previous. |
| 489 | if self.large_drop.is_none() { |
| 490 | if let Some(prev) = self.prev_estimate { |
| 491 | if estimate < prev * BITRATE_DROP_THRESHOLD { |
| 492 | self.large_drop = Some(LargeDrop { |
| 493 | when: now, |
| 494 | bitrate_before: prev, |
| 495 | }); |
| 496 | } |
| 497 | } |
| 498 | } |
| 499 | |
| 500 | // No large drop detected. |
| 501 | let Some(drop) = &self.large_drop else { |
| 502 | return false; |
| 503 | }; |
| 504 | |
| 505 | // Drop expires after 5 seconds. |
| 506 | if now.saturating_duration_since(drop.when) > BITRATE_DROP_TIMEOUT { |
| 507 | self.large_drop = None; |
| 508 | return false; |
| 509 | } |
| 510 | |
| 511 | // Large-drop probing requires ALR context (in ALR or recently exited). |
| 512 | if !self.in_alr() && !self.alr_ended_recently(now) { |
| 513 | return false; |
| 514 | } |
| 515 | |
| 516 | // Respect minimum interval between ALR probes. |
| 517 | if self.time_since_last_probe(now) < MIN_TIME_BETWEEN_ALR_PROBES { |
| 518 | return false; |
| 519 | } |
| 520 | |
| 521 | // Probe at 85% of pre-drop bitrate. |
| 522 | let target = drop.bitrate_before * PROBE_FRACTION_AFTER_DROP; |
| 523 | self.queue_probe(target, ProbeKind::LargeDrop, desired, now); |
| 524 | |
| 525 | self.large_drop = None; |
| 526 | true |
| 527 | } |
| 528 | |
| 529 | fn queue_probe(&mut self, bitrate: Bitrate, kind: ProbeKind, desired: Bitrate, now: Instant) { |
| 530 | // Cap at 2× desired bitrate. |
| 531 | let max = desired * MAX_PROBE_BITRATE_FACTOR; |
| 532 | let bitrate = bitrate.min(max); |
| 533 | |
| 534 | // No probe at too small values. |
| 535 | if bitrate < Bitrate::kbps(5) { |
| 536 | return; |
| 537 | } |
| 538 | |
| 539 | let cluster_id = self.next_cluster_id.inc(); |
| 540 | |
| 541 | let config = ProbeClusterConfig::new(cluster_id, bitrate, kind) |
| 542 | .with_min_packet_count(self.config.min_probe_packets_sent) |
| 543 | .with_duration(self.config.min_probe_duration) |
| 544 | .with_min_probe_delta(self.config.min_probe_delta); |
| 545 | |
| 546 | // Threshold for further exponential probing (probe_bitrate * 0.7). |
| 547 | let probe_further = bitrate * self.config.further_probe_threshold; |
| 548 | |
| 549 | self.pending.push_back(config); |
| 550 | self.last_probe = Some(LastProbe { |
| 551 | when: now, |
| 552 | kind, |
| 553 | further: probe_further, |
| 554 | was_estimate: self.last_estimate, |
| 555 | }); |
| 556 | } |
| 557 | |
| 558 | fn compute_next_timeout(&mut self, now: Instant) -> Instant { |
| 559 | // Exponential probing: wait for probe result before re-probing at same estimate. |
| 560 | // This handles the case where we sent a probe but haven't received updated estimate yet. |
| 561 | if let Some(last) = &self.last_probe { |
| 562 | if matches!(last.kind, ProbeKind::Initial | ProbeKind::Exponential) { |
| 563 | if self.scheduled_exponential.is_none() { |
| 564 | self.scheduled_exponential = Some(now + MAX_WAITING_TIME_FOR_PROBING_RESULT); |
| 565 | } |
| 566 | return self.scheduled_exponential.unwrap(); |
| 567 | } |
| 568 | } |
| 569 | |
| 570 | // ALR periodic probing |
| 571 | if self.in_alr() { |
| 572 | if self.scheduled_periodic_alr.is_none() { |
| 573 | self.scheduled_periodic_alr = Some(now + MIN_TIME_BETWEEN_ALR_PROBES); |
| 574 | } |
| 575 | return self.scheduled_periodic_alr.unwrap(); |
| 576 | } |
| 577 | |
| 578 | // Stagnant probing (only when not in ALR) |
| 579 | if !self.in_alr() { |
| 580 | if self.scheduled_stagnant.is_none() { |
| 581 | self.scheduled_stagnant = Some(now + MIN_TIME_BETWEEN_STAGNANT_PROBES); |
| 582 | } |
| 583 | return self.scheduled_stagnant.unwrap(); |
| 584 | } |
| 585 | |
| 586 | not_happening() |
| 587 | } |
| 588 | |
| 589 | fn can_probe(&self, estimate: Bitrate) -> bool { |
| 590 | // Infinite estimate indicates no valid measurement yet. |
| 591 | if estimate == Bitrate::INFINITY { |
| 592 | return false; |
| 593 | } |
| 594 | |
| 595 | // Only probe when delay-limited or loss-limited-but-increasing. |
| 596 | // Don't probe during active congestion (loss-limited, delay-increased). |
| 597 | matches!( |
| 598 | self.last_cause, |
| 599 | BandwidthLimitedCause::LossLimitedBweIncreasing |
| 600 | | BandwidthLimitedCause::DelayBasedLimited |
| 601 | ) |
| 602 | } |
| 603 | |
| 604 | fn in_alr(&self) -> bool { |
| 605 | self.alr_start.is_some() && self.alr_stop.is_none() |
| 606 | } |
| 607 | |
| 608 | fn alr_ended_recently(&self, now: Instant) -> bool { |
| 609 | self.alr_stop |
| 610 | .map(|stop| now.saturating_duration_since(stop) < ALR_ENDED_TIMEOUT) |
| 611 | .unwrap_or(false) |
| 612 | } |
| 613 | |
| 614 | fn is_during_initial(&self, now: Instant) -> bool { |
| 615 | let is_initial = matches!( |
| 616 | self.last_probe.map(|p| p.kind), |
| 617 | Some(ProbeKind::Initial) | Some(ProbeKind::Exponential) |
| 618 | ); |
| 619 | is_initial && self.time_since_last_probe(now) <= MAX_WAITING_TIME_FOR_PROBING_RESULT |
| 620 | } |
| 621 | |
| 622 | fn last_when(&self) -> Option<Instant> { |
| 623 | self.last_probe.map(|p| p.when) |
| 624 | } |
| 625 | |
| 626 | fn time_since_last_probe(&self, now: Instant) -> Duration { |
| 627 | self.last_when() |
| 628 | .map(|t| now.saturating_duration_since(t)) |
| 629 | .unwrap_or(Duration::MAX) |
| 630 | } |
| 631 | } |
| 632 | |
| 633 | /// Configuration using WebRTC default constants (no field-trial plumbing). |
| 634 | #[derive(Debug, Clone, Copy)] |
| 635 | struct Config { |
| 636 | // Initial/exponential probing |
| 637 | first_exponential_probe_scale: f64, // p1 = 3.0 |
| 638 | second_exponential_probe_scale: f64, // p2 = 6.0 |
| 639 | further_exponential_probe_scale: f64, // step_size = 2.0 |
| 640 | further_probe_threshold: f64, // 0.7 |
| 641 | |
| 642 | // Allocation probing |
| 643 | first_allocation_probe_scale: f64, // 1.0 |
| 644 | second_allocation_probe_scale: f64, // 2.0 |
| 645 | allocation_probe_limit_by_current_scale: f64, // 2.0 |
| 646 | |
| 647 | // Probe cluster config defaults |
| 648 | min_probe_packets_sent: usize, // 5 |
| 649 | min_probe_duration: Duration, // 15ms |
| 650 | min_probe_delta: Duration, // 2ms |
| 651 | |
| 652 | // Gating / limits |
| 653 | loss_limited_probe_scale: f64, // 1.5 |
| 654 | } |
| 655 | |
| 656 | impl Default for Config { |
| 657 | fn default() -> Self { |
| 658 | Self { |
| 659 | first_exponential_probe_scale: 3.0, |
| 660 | second_exponential_probe_scale: 6.0, |
| 661 | further_exponential_probe_scale: 2.0, |
| 662 | further_probe_threshold: 0.7, |
| 663 | |
| 664 | first_allocation_probe_scale: 1.0, |
| 665 | second_allocation_probe_scale: 2.0, |
| 666 | allocation_probe_limit_by_current_scale: 2.0, |
| 667 | |
| 668 | min_probe_packets_sent: 5, |
| 669 | min_probe_duration: Duration::from_millis(15), |
| 670 | min_probe_delta: Duration::from_millis(2), |
| 671 | |
| 672 | loss_limited_probe_scale: 1.5, |
| 673 | } |
| 674 | } |
| 675 | } |
| 676 | |
| 677 | impl BandwidthLimitedCause { |
| 678 | /// Probe scale factor for exponential probing. |
| 679 | /// |
| 680 | /// When loss-limited but increasing, use a more conservative 1.575× (1.5 * 1.05). |
| 681 | /// Otherwise use the standard 2× scale. |
| 682 | fn probe_scale(&self, config: &Config) -> f64 { |
| 683 | match self { |
| 684 | BandwidthLimitedCause::LossLimitedBweIncreasing => { |
| 685 | config.loss_limited_probe_scale * (1.0 + PROBE_UNCERTAINTY) |
| 686 | } |
| 687 | _ => config.further_exponential_probe_scale, |
| 688 | } |
| 689 | } |
| 690 | } |
| 691 | |
| 692 | #[cfg(test)] |
| 693 | mod test { |
| 694 | use super::*; |
| 695 | |
| 696 | #[test] |
| 697 | fn initial_exponential_probes_are_queued_and_emitted_one_per_tick() { |
| 698 | let mut pc = ProbeControl::new(); |
| 699 | pc.enable(true); |
| 700 | let now = Instant::now(); |
| 701 | |
| 702 | pc.set_desired_bitrate(Bitrate::mbps(50)); |
| 703 | pc.set_estimated_bitrate(Bitrate::kbps(300), BandwidthLimitedCause::DelayBasedLimited); |
| 704 | |
| 705 | // First handle_timeout triggers initial probing and returns first probe. |
| 706 | let p1 = pc.handle_timeout(now).unwrap(); |
| 707 | |
| 708 | // poll_timeout returns already_happened while there are pending probes. |
| 709 | assert_eq!(pc.poll_timeout(), already_happened()); |
| 710 | |
| 711 | // Second handle_timeout returns the second queued probe. |
| 712 | let p2 = pc.handle_timeout(now).unwrap(); |
| 713 | |
| 714 | assert_eq!(p1.target_bitrate(), Bitrate::kbps(900)); |
| 715 | assert_eq!(p2.target_bitrate(), Bitrate::kbps(1800)); |
| 716 | assert_eq!(p1.min_packet_count(), 5); |
| 717 | assert_eq!(p1.min_probe_delta(), Duration::from_millis(2)); |
| 718 | assert!(!p1.is_alr_probe()); |
| 719 | |
| 720 | // Queue drained - no more probes. |
| 721 | assert!(pc.handle_timeout(now).is_none()); |
| 722 | } |
| 723 | |
| 724 | #[test] |
| 725 | fn further_probe_is_triggered_when_probe_result_is_high_enough() { |
| 726 | let mut pc = ProbeControl::new(); |
| 727 | pc.enable(true); |
| 728 | let now = Instant::now(); |
| 729 | |
| 730 | pc.enable(true); |
| 731 | pc.set_desired_bitrate(Bitrate::mbps(50)); |
| 732 | pc.set_estimated_bitrate(Bitrate::mbps(1), BandwidthLimitedCause::DelayBasedLimited); |
| 733 | |
| 734 | // Drain initial two probes. |
| 735 | let _ = pc.handle_timeout(now).unwrap(); |
| 736 | let _ = pc.handle_timeout(now).unwrap(); |
| 737 | |
| 738 | // WebRTC rule: if measured bitrate > min_bitrate_to_probe_further, probe at 2x measured. |
| 739 | // min_bitrate_to_probe_further is 0.7 * last_probe_rate (6x start) = 4.2 Mbps. |
| 740 | pc.set_estimated_bitrate(Bitrate::mbps(5), BandwidthLimitedCause::DelayBasedLimited); |
| 741 | |
| 742 | let p = pc.handle_timeout(now + Duration::from_millis(10)).unwrap(); |
| 743 | assert_eq!(p.target_bitrate(), Bitrate::mbps(10)); |
| 744 | } |
| 745 | |
| 746 | #[test] |
| 747 | fn allocation_probe_is_triggered_in_alr_when_allocation_increases() { |
| 748 | let mut pc = ProbeControl::new(); |
| 749 | pc.enable(true); |
| 750 | let now = Instant::now(); |
| 751 | |
| 752 | pc.set_desired_bitrate(Bitrate::mbps(1)); |
| 753 | pc.set_estimated_bitrate(Bitrate::mbps(1), BandwidthLimitedCause::DelayBasedLimited); |
| 754 | |
| 755 | // Drain initial probes. |
| 756 | let _ = pc.handle_timeout(now).unwrap(); |
| 757 | let _ = pc.handle_timeout(now).unwrap(); |
| 758 | |
| 759 | // Time out waiting for probing result -> probing complete. |
| 760 | assert!(pc.handle_timeout(now + Duration::from_secs(2)).is_none()); |
| 761 | |
| 762 | // Enter ALR |
| 763 | pc.set_alr_start_time(now + Duration::from_secs(2)); |
| 764 | |
| 765 | // No probe yet - desired hasn't increased |
| 766 | assert!(pc.handle_timeout(now + Duration::from_secs(2)).is_none()); |
| 767 | |
| 768 | // Increase desired bitrate while in ALR (desired > prev AND desired > estimate) |
| 769 | pc.set_desired_bitrate(Bitrate::mbps(4)); |
| 770 | |
| 771 | // Should trigger allocation probe: p1 = 4 Mbps * 1.0 = 4 Mbps, capped by 2× estimate = 2 Mbps |
| 772 | let p = pc.handle_timeout(now + Duration::from_secs(2)).unwrap(); |
| 773 | assert_eq!(p.target_bitrate(), Bitrate::mbps(2)); |
| 774 | } |
| 775 | |
| 776 | #[test] |
| 777 | fn handles_bitrate_infinity_without_panic() { |
| 778 | let mut pc = ProbeControl::new(); |
| 779 | pc.enable(true); |
| 780 | let now = Instant::now(); |
| 781 | |
| 782 | pc.set_desired_bitrate(Bitrate::mbps(50)); |
| 783 | |
| 784 | // Should not panic with Infinity |
| 785 | pc.set_estimated_bitrate(Bitrate::INFINITY, BandwidthLimitedCause::DelayBasedLimited); |
| 786 | |
| 787 | // Verify behavior is reasonable (no probing with infinite estimate) |
| 788 | assert!(pc.handle_timeout(now).is_none()); |
| 789 | } |
| 790 | |
| 791 | #[test] |
| 792 | fn handles_clock_skew_gracefully() { |
| 793 | let mut pc = ProbeControl::new(); |
| 794 | pc.enable(true); |
| 795 | let now = Instant::now(); |
| 796 | |
| 797 | pc.set_desired_bitrate(Bitrate::mbps(50)); |
| 798 | pc.set_estimated_bitrate(Bitrate::kbps(300), BandwidthLimitedCause::DelayBasedLimited); |
| 799 | |
| 800 | // Drain initial probes |
| 801 | let _ = pc.handle_timeout(now); |
| 802 | let _ = pc.handle_timeout(now); |
| 803 | |
| 804 | // Simulate time going backwards (clock skew) |
| 805 | let earlier = now - Duration::from_secs(5); |
| 806 | |
| 807 | // Should handle gracefully with saturating_duration_since |
| 808 | let _ = pc.handle_timeout(earlier); |
| 809 | |
| 810 | // Should still be able to continue normally |
| 811 | let _ = pc.handle_timeout(now + Duration::from_secs(1)); |
| 812 | } |
| 813 | |
| 814 | #[test] |
| 815 | fn handles_max_bitrate_zero() { |
| 816 | let mut pc = ProbeControl::new(); |
| 817 | pc.enable(true); |
| 818 | let now = Instant::now(); |
| 819 | |
| 820 | // Set max_bitrate to zero - this is rejected as first value to avoid |
| 821 | // creating probes with zero cap. |
| 822 | pc.set_desired_bitrate(Bitrate::ZERO); |
| 823 | pc.set_estimated_bitrate(Bitrate::kbps(300), BandwidthLimitedCause::DelayBasedLimited); |
| 824 | |
| 825 | // No probes should be created since desired was rejected. |
| 826 | let p1 = pc.handle_timeout(now); |
| 827 | assert!(p1.is_none(), "Should not create probes with zero desired"); |
| 828 | } |
| 829 | |
| 830 | #[test] |
| 831 | fn allocation_probe_fires_when_desired_increases_in_alr() { |
| 832 | let mut pc = ProbeControl::new(); |
| 833 | pc.enable(true); |
| 834 | let now = Instant::now(); |
| 835 | |
| 836 | pc.set_desired_bitrate(Bitrate::kbps(500)); |
| 837 | pc.set_estimated_bitrate(Bitrate::kbps(500), BandwidthLimitedCause::DelayBasedLimited); |
| 838 | |
| 839 | // Drain initial probes |
| 840 | let _ = pc.handle_timeout(now); |
| 841 | let _ = pc.handle_timeout(now); |
| 842 | |
| 843 | // Timeout to reach probing complete |
| 844 | assert!(pc.handle_timeout(now + Duration::from_secs(2)).is_none()); |
| 845 | |
| 846 | // Enter ALR |
| 847 | pc.set_alr_start_time(now + Duration::from_secs(3)); |
| 848 | |
| 849 | // No probe on ALR entry alone |
| 850 | assert!(pc.handle_timeout(now + Duration::from_secs(3)).is_none()); |
| 851 | |
| 852 | // Increase desired while in ALR (desired > prev AND desired > estimate) |
| 853 | pc.set_desired_bitrate(Bitrate::mbps(4)); |
| 854 | |
| 855 | // Should trigger allocation probe |
| 856 | let probe = pc.handle_timeout(now + Duration::from_secs(3)); |
| 857 | assert!( |
| 858 | probe.is_some(), |
| 859 | "Allocation probe should trigger when desired increases in ALR" |
| 860 | ); |
| 861 | } |
| 862 | |
| 863 | #[test] |
| 864 | fn large_drop_probing_after_alr_ended() { |
| 865 | let mut pc = ProbeControl::new(); |
| 866 | pc.enable(true); |
| 867 | let now = Instant::now(); |
| 868 | |
| 869 | pc.set_desired_bitrate(Bitrate::mbps(5)); |
| 870 | pc.set_estimated_bitrate(Bitrate::mbps(5), BandwidthLimitedCause::DelayBasedLimited); |
| 871 | |
| 872 | // Drain initial probes |
| 873 | let _ = pc.handle_timeout(now); |
| 874 | let _ = pc.handle_timeout(now); |
| 875 | |
| 876 | // Timeout to probing complete |
| 877 | assert!(pc.handle_timeout(now + Duration::from_secs(2)).is_none()); |
| 878 | |
| 879 | // Enter and exit ALR (large-drop works when ALR ended recently) |
| 880 | pc.set_alr_start_time(now + Duration::from_secs(2)); |
| 881 | pc.set_alr_stop_time(now + Duration::from_secs(3)); |
| 882 | |
| 883 | // Simulate large drop (to 60% of original = 3 Mbps, below 66% threshold) |
| 884 | pc.set_estimated_bitrate(Bitrate::mbps(3), BandwidthLimitedCause::DelayBasedLimited); |
| 885 | |
| 886 | // Check at now+5s (within 3s of ALR ending, so alr_ended_recently is true) |
| 887 | let later = now + Duration::from_secs(5); |
| 888 | |
| 889 | // Should trigger large-drop recovery probe at 85% of pre-drop rate (4.25 Mbps) |
| 890 | let p = pc.handle_timeout(later); |
| 891 | assert!(p.is_some(), "Large-drop recovery should schedule probe"); |
| 892 | if let Some(probe) = p { |
| 893 | // 85% of 5 Mbps = 4.25 Mbps |
| 894 | assert!(probe.target_bitrate() >= Bitrate::mbps(4)); |
| 895 | assert!(probe.target_bitrate() <= Bitrate::mbps(5)); |
| 896 | } |
| 897 | } |
| 898 | |
| 899 | #[test] |
| 900 | fn allocation_probe_requires_desired_increase_in_alr() { |
| 901 | let mut pc = ProbeControl::new(); |
| 902 | pc.enable(true); |
| 903 | let now = Instant::now(); |
| 904 | |
| 905 | pc.set_desired_bitrate(Bitrate::mbps(5)); |
| 906 | pc.set_estimated_bitrate(Bitrate::mbps(1), BandwidthLimitedCause::DelayBasedLimited); |
| 907 | |
| 908 | // Drain initial probes |
| 909 | let _ = pc.handle_timeout(now); |
| 910 | let _ = pc.handle_timeout(now); |
| 911 | |
| 912 | // Timeout to probing complete |
| 913 | assert!(pc.handle_timeout(now + Duration::from_secs(2)).is_none()); |
| 914 | |
| 915 | // Enter ALR with estimate < max_bitrate |
| 916 | pc.set_alr_start_time(now + Duration::from_secs(2)); |
| 917 | |
| 918 | // No allocation probe on ALR entry - need desired to increase |
| 919 | let probe = pc.handle_timeout(now + Duration::from_secs(2)); |
| 920 | assert!( |
| 921 | probe.is_none(), |
| 922 | "Should NOT trigger allocation probe on ALR entry alone" |
| 923 | ); |
| 924 | |
| 925 | // Increase desired while in ALR |
| 926 | pc.set_desired_bitrate(Bitrate::mbps(10)); |
| 927 | |
| 928 | // Now should trigger allocation probe |
| 929 | let probe = pc.handle_timeout(now + Duration::from_secs(2)); |
| 930 | assert!( |
| 931 | probe.is_some(), |
| 932 | "Should trigger allocation probe when desired increases in ALR" |
| 933 | ); |
| 934 | } |
| 935 | |
| 936 | #[test] |
| 937 | fn periodic_alr_probing() { |
| 938 | let mut pc = ProbeControl::new(); |
| 939 | pc.enable(true); |
| 940 | let now = Instant::now(); |
| 941 | |
| 942 | pc.set_desired_bitrate(Bitrate::mbps(5)); |
| 943 | pc.set_estimated_bitrate(Bitrate::mbps(1), BandwidthLimitedCause::DelayBasedLimited); |
| 944 | |
| 945 | // Drain initial probes |
| 946 | let _ = pc.handle_timeout(now); |
| 947 | let _ = pc.handle_timeout(now); |
| 948 | |
| 949 | // Timeout to probing complete |
| 950 | assert!(pc.handle_timeout(now + Duration::from_secs(2)).is_none()); |
| 951 | |
| 952 | // Enter ALR |
| 953 | pc.set_alr_start_time(now + Duration::from_secs(2)); |
| 954 | |
| 955 | // No immediate probe on ALR entry |
| 956 | assert!(pc.handle_timeout(now + Duration::from_secs(2)).is_none()); |
| 957 | |
| 958 | // Wait 5 seconds for periodic probe (2s to complete initial + 5s = 7s) |
| 959 | let probe = pc.handle_timeout(now + Duration::from_secs(7)); |
| 960 | assert!( |
| 961 | probe.is_some(), |
| 962 | "Should trigger periodic ALR probe after 5 seconds in ALR" |
| 963 | ); |
| 964 | assert!(probe.unwrap().is_alr_probe()); |
| 965 | } |
| 966 | |
| 967 | #[test] |
| 968 | fn periodic_alr_probing_continues_even_when_estimate_reaches_max() { |
| 969 | let mut pc = ProbeControl::new(); |
| 970 | pc.enable(true); |
| 971 | let now = Instant::now(); |
| 972 | |
| 973 | pc.set_desired_bitrate(Bitrate::mbps(2)); |
| 974 | pc.set_estimated_bitrate(Bitrate::mbps(1), BandwidthLimitedCause::DelayBasedLimited); |
| 975 | |
| 976 | // Drain initial probes |
| 977 | let _ = pc.handle_timeout(now); |
| 978 | let _ = pc.handle_timeout(now); |
| 979 | |
| 980 | // Timeout to probing complete |
| 981 | assert!(pc.handle_timeout(now + Duration::from_secs(2)).is_none()); |
| 982 | |
| 983 | // Enter ALR |
| 984 | pc.set_alr_start_time(now + Duration::from_secs(2)); |
| 985 | |
| 986 | // No immediate probe on ALR entry |
| 987 | assert!(pc.handle_timeout(now + Duration::from_secs(2)).is_none()); |
| 988 | |
| 989 | // Now increase estimate to match max_bitrate |
| 990 | pc.set_estimated_bitrate(Bitrate::mbps(2), BandwidthLimitedCause::DelayBasedLimited); |
| 991 | |
| 992 | // Wait 5 seconds - should still trigger periodic probe in ALR |
| 993 | // even though estimate >= max_bitrate, to maintain confidence in the estimate |
| 994 | let probe = pc.handle_timeout(now + Duration::from_secs(7)); |
| 995 | assert!( |
| 996 | probe.is_some(), |
| 997 | "Should continue periodic probing in ALR even when estimate >= max_bitrate" |
| 998 | ); |
| 999 | } |
| 1000 | } |