File
Blob: firmware/vendor/str0m/src/bwe/probe/estimator.rs
| 1 | use std::collections::VecDeque; |
| 2 | use std::fmt; |
| 3 | use std::time::{Duration, Instant}; |
| 4 | |
| 5 | use super::super::macros::log_probe_bitrate_estimate; |
| 6 | use super::ProbeClusterConfig; |
| 7 | use crate::rtp_::{Bitrate, DataSize, TwccClusterId, TwccSendRecord}; |
| 8 | use crate::util::not_happening; |
| 9 | |
| 10 | /// Minimum ratio of packets we need to receive for a valid probe (80%). |
| 11 | const MIN_RECEIVED_PROBES_RATIO: f64 = 0.80; |
| 12 | |
| 13 | /// Minimum ratio of bytes we need to receive for a valid probe (80%). |
| 14 | const MIN_RECEIVED_BYTES_RATIO: f64 = 0.80; |
| 15 | |
| 16 | /// Minimum packet count for a valid probe cluster (WebRTC's kMinClusterSize). |
| 17 | const MIN_CLUSTER_SIZE: usize = 4; |
| 18 | |
| 19 | /// The maximum valid duration between first and last probe packet on send/receive side. |
| 20 | /// Matches WebRTC's `kMaxProbeInterval` in `probe_bitrate_estimator.cc`. |
| 21 | const MAX_PROBE_INTERVAL: Duration = Duration::from_secs(1); |
| 22 | |
| 23 | /// The maximum |receive rate| / |send rate| ratio for a valid estimate. |
| 24 | /// Matches WebRTC's `kMaxValidRatio`. |
| 25 | const MAX_VALID_RATIO: f64 = 2.0; |
| 26 | |
| 27 | /// Minimum |receive rate| / |send rate| ratio to consider the link unsaturated. |
| 28 | /// Matches WebRTC's `kMinRatioForUnsaturatedLink`. |
| 29 | const MIN_RATIO_FOR_UNSATURATED_LINK: f64 = 0.9; |
| 30 | |
| 31 | /// Target utilization when we believe we've found the true capacity. |
| 32 | /// Matches WebRTC's `kTargetUtilizationFraction`. |
| 33 | const TARGET_UTILIZATION_FRACTION: f64 = 0.95; |
| 34 | |
| 35 | /// Maximum number of active probes before we trigger cleanup. |
| 36 | const MAX_ACTIVE_PROBES: usize = 20; |
| 37 | |
| 38 | /// Probes older than this are considered stale and will be removed when hitting the cap. |
| 39 | const STALE_PROBE_THRESHOLD: Duration = Duration::from_secs(5); |
| 40 | |
| 41 | /// Analyzes probe cluster results from TWCC feedback. |
| 42 | /// |
| 43 | /// This component takes packets tagged with a `TwccClusterId` and calculates the |
| 44 | /// achieved bitrate for each probe cluster. |
| 45 | /// |
| 46 | /// **Important:** This follows WebRTC's `ProbeBitrateEstimator` semantics: |
| 47 | /// only packets with a known remote receive timestamp are included in the probe |
| 48 | /// result. Probe packets reported as lost (no remote receive timestamp) are ignored. |
| 49 | #[derive(Debug)] |
| 50 | pub struct ProbeEstimator { |
| 51 | /// Active probe states (VecDeque for efficient front removal). |
| 52 | states: VecDeque<ProbeEstimatorState>, |
| 53 | |
| 54 | /// Clusters that were updated in the last call to `update`. |
| 55 | did_update: VecDeque<TwccClusterId>, |
| 56 | } |
| 57 | |
| 58 | #[derive(Debug)] |
| 59 | struct ProbeEstimatorState { |
| 60 | /// Configuration of the active probe (targets for validation). |
| 61 | config: ProbeClusterConfig, |
| 62 | |
| 63 | /// When this probe was created. |
| 64 | created_at: Instant, |
| 65 | |
| 66 | /// When to erase this cluster's state (cluster history expiry). |
| 67 | finalize_at: Instant, |
| 68 | |
| 69 | /// First (earliest) send time among packets included in this probe. |
| 70 | first_send_time: Option<Instant>, |
| 71 | /// Last (latest) send time among packets included in this probe. |
| 72 | last_send_time: Option<Instant>, |
| 73 | /// Size of the packet with the last send time (excluded from send-rate calculation). |
| 74 | size_last_send: DataSize, |
| 75 | |
| 76 | /// First (earliest) receive time among packets included in this probe. |
| 77 | first_recv_time: Option<Instant>, |
| 78 | /// Last (latest) receive time among packets included in this probe. |
| 79 | last_recv_time: Option<Instant>, |
| 80 | /// Size of the packet with the first receive time (excluded from receive-rate calculation). |
| 81 | size_first_receive: DataSize, |
| 82 | |
| 83 | /// Total bytes for packets included in this probe (received packets only). |
| 84 | total_bytes: DataSize, |
| 85 | /// Number of packets included in this probe (received packets only). |
| 86 | packet_count: usize, |
| 87 | } |
| 88 | |
| 89 | impl ProbeEstimator { |
| 90 | pub fn new() -> Self { |
| 91 | Self { |
| 92 | states: VecDeque::new(), |
| 93 | did_update: VecDeque::with_capacity(10), |
| 94 | } |
| 95 | } |
| 96 | |
| 97 | /// Start analyzing a new probe cluster. |
| 98 | /// |
| 99 | /// Resets all accumulated state and begins watching for packets with the |
| 100 | /// given cluster ID. Returns `true` if the probe was started, `false` if |
| 101 | /// it was rejected due to too many active probes. |
| 102 | pub fn probe_start(&mut self, config: ProbeClusterConfig, now: Instant) -> bool { |
| 103 | // Under normal operation, we expect at most 2-4 active probes: |
| 104 | // - Initial exponential probing: 2 probes (3×, 6×) |
| 105 | // - Further probing: 1-2 additional probes |
| 106 | // - Allocation/recovery probing: 1-2 more |
| 107 | // Even with rapid probe sequences and 1-second cleanup delay, we shouldn't |
| 108 | // accumulate more than ~8 probes. If we hit 20, it indicates a bug in probe |
| 109 | // lifecycle management (missing end_probe() calls or handle_timeout() not running). |
| 110 | |
| 111 | // If we've hit the cap, try to clean up stale probes first. |
| 112 | if self.states.len() >= MAX_ACTIVE_PROBES { |
| 113 | let before_cleanup = self.states.len(); |
| 114 | self.states |
| 115 | .retain(|s| now.saturating_duration_since(s.created_at) < STALE_PROBE_THRESHOLD); |
| 116 | |
| 117 | let removed = before_cleanup - self.states.len(); |
| 118 | if removed > 0 { |
| 119 | debug!( |
| 120 | removed, |
| 121 | remaining = self.states.len(), |
| 122 | "Cleaned up stale probes" |
| 123 | ); |
| 124 | } |
| 125 | |
| 126 | // If still at cap after cleanup, reject the new probe. |
| 127 | if self.states.len() >= MAX_ACTIVE_PROBES { |
| 128 | debug!("Rejecting new probe: too many active probes with none stale"); |
| 129 | return false; |
| 130 | } |
| 131 | } |
| 132 | |
| 133 | self.states.push_back(ProbeEstimatorState::new(config, now)); |
| 134 | true |
| 135 | } |
| 136 | |
| 137 | /// Process TWCC feedback records. |
| 138 | /// |
| 139 | /// Only accumulates packets that match the active cluster ID. All other |
| 140 | /// packets are ignored. |
| 141 | pub fn update<'t>( |
| 142 | &mut self, |
| 143 | records: impl Iterator<Item = &'t TwccSendRecord>, |
| 144 | ) -> impl Iterator<Item = (ProbeClusterConfig, Bitrate)> + '_ { |
| 145 | // Keep track of which clusters were updated in this call. |
| 146 | self.did_update.clear(); |
| 147 | |
| 148 | for record in records { |
| 149 | let Some(cluster) = record.cluster() else { |
| 150 | continue; |
| 151 | }; |
| 152 | |
| 153 | // Find the state for this cluster. |
| 154 | let maybe_state = self |
| 155 | .states |
| 156 | .iter_mut() |
| 157 | .find(|s| s.config.cluster() == cluster); |
| 158 | |
| 159 | let Some(state) = maybe_state else { |
| 160 | continue; |
| 161 | }; |
| 162 | |
| 163 | let did_update = state.update(record); |
| 164 | |
| 165 | if did_update { |
| 166 | // The correct behavior is that the _last updated_ |
| 167 | // is emitted last, so that the consumer of the returned |
| 168 | // iterator gets the latest probe result last. |
| 169 | self.did_update.retain(|c| *c != cluster); |
| 170 | self.did_update.push_back(cluster); |
| 171 | } |
| 172 | } |
| 173 | |
| 174 | self.did_update |
| 175 | .iter() |
| 176 | .filter_map(|cluster| self.states.iter().find(|s| s.config.cluster() == *cluster)) |
| 177 | .filter_map(|s| s.calculate_bitrate()) |
| 178 | } |
| 179 | |
| 180 | /// Mark the probe as ended. |
| 181 | /// |
| 182 | /// The probe will continue collecting feedback during a cluster history |
| 183 | /// period after the probe is finished. This period must be shorter than |
| 184 | /// the time between probe clusters to avoid overlap. |
| 185 | pub fn end_probe(&mut self, now: Instant, cluster_id: TwccClusterId) { |
| 186 | let maybe_state = self |
| 187 | .states |
| 188 | .iter_mut() |
| 189 | .find(|s| s.config.cluster() == cluster_id); |
| 190 | |
| 191 | let Some(state) = maybe_state else { |
| 192 | return; |
| 193 | }; |
| 194 | |
| 195 | state.end_probe(now); |
| 196 | } |
| 197 | |
| 198 | pub fn poll_timeout(&self) -> Instant { |
| 199 | self.states |
| 200 | .iter() |
| 201 | .map(|s| s.finalize_at) |
| 202 | .min() |
| 203 | .unwrap_or(not_happening()) |
| 204 | } |
| 205 | |
| 206 | /// Finalize probes that are ready. |
| 207 | pub fn handle_timeout(&mut self, now: Instant) { |
| 208 | self.states.retain(|s| { |
| 209 | let do_keep = now < s.finalize_at; |
| 210 | if do_keep { |
| 211 | return true; |
| 212 | } |
| 213 | |
| 214 | let result = s.do_calculate_bitrate(); |
| 215 | if let ProbeResult::Estimate(_) = result { |
| 216 | // Already logged in calculate_bitrate() during update(). |
| 217 | } else { |
| 218 | // Log the final rejection reason for the probe. |
| 219 | trace!(%result, "Probe result"); |
| 220 | } |
| 221 | |
| 222 | false |
| 223 | }); |
| 224 | } |
| 225 | |
| 226 | /// Clear all active probes. |
| 227 | /// |
| 228 | /// This should be called when probing is no longer possible. |
| 229 | pub fn clear_probes(&mut self) { |
| 230 | self.states.clear(); |
| 231 | } |
| 232 | } |
| 233 | |
| 234 | impl ProbeEstimatorState { |
| 235 | pub fn new(config: ProbeClusterConfig, now: Instant) -> Self { |
| 236 | Self { |
| 237 | config, |
| 238 | created_at: now, |
| 239 | finalize_at: not_happening(), |
| 240 | first_send_time: None, |
| 241 | last_send_time: None, |
| 242 | size_last_send: DataSize::ZERO, |
| 243 | first_recv_time: None, |
| 244 | last_recv_time: None, |
| 245 | size_first_receive: DataSize::ZERO, |
| 246 | total_bytes: DataSize::ZERO, |
| 247 | packet_count: 0, |
| 248 | } |
| 249 | } |
| 250 | |
| 251 | fn update(&mut self, record: &TwccSendRecord) -> bool { |
| 252 | // Only packets with a known remote receive time participate in probe estimation. |
| 253 | let Some(recv_time) = record.remote_recv_time() else { |
| 254 | return false; // lost/unreceived packet -> ignore for probe result |
| 255 | }; |
| 256 | |
| 257 | let packet_size = DataSize::from(record.size()); |
| 258 | let send_time = record.local_send_time(); |
| 259 | |
| 260 | // Track min/max send time among included packets. |
| 261 | let first = self.first_send_time.get_or_insert(send_time); |
| 262 | *first = (*first).min(send_time); |
| 263 | |
| 264 | let last = self.last_send_time.get_or_insert(send_time); |
| 265 | if send_time >= *last { |
| 266 | *last = send_time; |
| 267 | self.size_last_send = packet_size; |
| 268 | } |
| 269 | |
| 270 | // Track min/max receive time among included packets. |
| 271 | let first_recv = self.first_recv_time.get_or_insert(recv_time); |
| 272 | if recv_time <= *first_recv { |
| 273 | *first_recv = recv_time; |
| 274 | self.size_first_receive = packet_size; |
| 275 | } |
| 276 | |
| 277 | let last_recv = self.last_recv_time.get_or_insert(recv_time); |
| 278 | *last_recv = (*last_recv).max(recv_time); |
| 279 | |
| 280 | self.total_bytes += packet_size; |
| 281 | self.packet_count += 1; |
| 282 | |
| 283 | true |
| 284 | } |
| 285 | |
| 286 | fn calculate_bitrate(&self) -> Option<(ProbeClusterConfig, Bitrate)> { |
| 287 | let result = self.do_calculate_bitrate(); |
| 288 | |
| 289 | let ProbeResult::Estimate(bitrate) = result else { |
| 290 | return None; |
| 291 | }; |
| 292 | |
| 293 | // Log the estimates continuously during the probe. |
| 294 | trace!(%result, "Probe result"); |
| 295 | log_probe_bitrate_estimate!(bitrate.as_f64()); |
| 296 | |
| 297 | Some((self.config, bitrate)) |
| 298 | } |
| 299 | |
| 300 | /// Calculate the estimated bitrate for this probe cluster. |
| 301 | fn do_calculate_bitrate(&self) -> ProbeResult { |
| 302 | // WebRTC requires at least kMinClusterSize (4) packets received. |
| 303 | // We may send more, but packet loss can result in fewer received packets. |
| 304 | if self.packet_count < MIN_CLUSTER_SIZE { |
| 305 | return ProbeResult::ClusterTooSmall { |
| 306 | recv: self.packet_count, |
| 307 | limit: MIN_CLUSTER_SIZE, |
| 308 | }; |
| 309 | } |
| 310 | |
| 311 | // Also check we received enough of what was sent |
| 312 | let min_packets = |
| 313 | (self.config.min_packet_count() as f64 * MIN_RECEIVED_PROBES_RATIO) as usize; |
| 314 | let min_bytes = DataSize::bytes( |
| 315 | (self.config.target_bytes().as_bytes_usize() as f64 * MIN_RECEIVED_BYTES_RATIO) as i64, |
| 316 | ); |
| 317 | |
| 318 | if self.packet_count < min_packets { |
| 319 | return ProbeResult::InsufficientPackets { |
| 320 | recv: self.packet_count, |
| 321 | limit: min_packets, |
| 322 | }; |
| 323 | } |
| 324 | if self.total_bytes < min_bytes { |
| 325 | return ProbeResult::InsufficientBytes { |
| 326 | recv: self.total_bytes, |
| 327 | limit: min_bytes, |
| 328 | }; |
| 329 | } |
| 330 | |
| 331 | // Get timing bounds |
| 332 | let Some(first_send) = self.first_send_time else { |
| 333 | return ProbeResult::MissingTimingInfo; |
| 334 | }; |
| 335 | let Some(last_send) = self.last_send_time else { |
| 336 | return ProbeResult::MissingTimingInfo; |
| 337 | }; |
| 338 | let send_interval = last_send.saturating_duration_since(first_send); |
| 339 | |
| 340 | let Some(first_recv) = self.first_recv_time else { |
| 341 | return ProbeResult::MissingTimingInfo; |
| 342 | }; |
| 343 | let Some(last_recv) = self.last_recv_time else { |
| 344 | return ProbeResult::MissingTimingInfo; |
| 345 | }; |
| 346 | let recv_interval = last_recv.saturating_duration_since(first_recv); |
| 347 | |
| 348 | // Intervals must be positive and within bounds. |
| 349 | if send_interval.is_zero() { |
| 350 | return ProbeResult::SendIntervalInvalid { |
| 351 | interval: send_interval, |
| 352 | }; |
| 353 | } |
| 354 | if send_interval > MAX_PROBE_INTERVAL { |
| 355 | return ProbeResult::SendIntervalTooLong { |
| 356 | interval: send_interval, |
| 357 | }; |
| 358 | } |
| 359 | if recv_interval.is_zero() || recv_interval > MAX_PROBE_INTERVAL { |
| 360 | return ProbeResult::RecvIntervalInvalid { |
| 361 | interval: recv_interval, |
| 362 | }; |
| 363 | } |
| 364 | |
| 365 | // WebRTC boundary exclusions: |
| 366 | // - exclude the last sent packet size when computing send rate |
| 367 | // - exclude the first received packet size when computing receive rate |
| 368 | let send_size = self.total_bytes.saturating_sub(self.size_last_send); |
| 369 | let recv_size = self.total_bytes.saturating_sub(self.size_first_receive); |
| 370 | if send_size <= DataSize::ZERO || recv_size <= DataSize::ZERO { |
| 371 | return ProbeResult::InvalidDataSize; |
| 372 | } |
| 373 | |
| 374 | let recv_rate = recv_size / recv_interval; |
| 375 | let send_rate = send_size / send_interval; |
| 376 | |
| 377 | // WebRTC validation: reject if receive/send ratio is too high. |
| 378 | let ratio = recv_rate.as_f64() / send_rate.as_f64(); |
| 379 | if ratio > MAX_VALID_RATIO { |
| 380 | return ProbeResult::InvalidSendReceiveRatio { |
| 381 | ratio, |
| 382 | limit: MAX_VALID_RATIO, |
| 383 | }; |
| 384 | } |
| 385 | |
| 386 | // Match WebRTC semantics: |
| 387 | // - estimate is the min(send_rate, recv_rate) |
| 388 | // - if recv_rate is significantly lower than send_rate, assume saturation and |
| 389 | // return a conservative fraction of recv_rate. |
| 390 | let mut estimate = send_rate.min(recv_rate); |
| 391 | if recv_rate < send_rate * MIN_RATIO_FOR_UNSATURATED_LINK { |
| 392 | estimate = recv_rate * TARGET_UTILIZATION_FRACTION; |
| 393 | } |
| 394 | |
| 395 | ProbeResult::Estimate(estimate) |
| 396 | } |
| 397 | |
| 398 | fn end_probe(&mut self, now: Instant) { |
| 399 | self.finalize_at = now + Duration::from_secs(1); |
| 400 | } |
| 401 | } |
| 402 | |
| 403 | /// Result of a probe cluster estimation. |
| 404 | #[derive(Debug, Clone, Copy, PartialEq)] |
| 405 | enum ProbeResult { |
| 406 | /// Successfully estimated bitrate |
| 407 | Estimate(Bitrate), |
| 408 | /// Not enough packets in cluster (< 4) |
| 409 | ClusterTooSmall { recv: usize, limit: usize }, |
| 410 | /// Insufficient packets received (< 80% of sent) |
| 411 | InsufficientPackets { recv: usize, limit: usize }, |
| 412 | /// Insufficient bytes received (< 80% of sent) |
| 413 | InsufficientBytes { recv: DataSize, limit: DataSize }, |
| 414 | /// Send interval too long (> 1 second) |
| 415 | SendIntervalTooLong { interval: Duration }, |
| 416 | /// Send interval invalid (zero) |
| 417 | SendIntervalInvalid { interval: Duration }, |
| 418 | /// Receive interval invalid (zero or > 1 second) |
| 419 | RecvIntervalInvalid { interval: Duration }, |
| 420 | /// Invalid receive/send ratio (recv_rate / send_rate too high) |
| 421 | InvalidSendReceiveRatio { ratio: f64, limit: f64 }, |
| 422 | /// Calculated data size is zero |
| 423 | InvalidDataSize, |
| 424 | /// Missing timing information |
| 425 | MissingTimingInfo, |
| 426 | } |
| 427 | |
| 428 | impl fmt::Display for ProbeResult { |
| 429 | fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { |
| 430 | match self { |
| 431 | ProbeResult::Estimate(bitrate) => write!(f, "estimate={}", bitrate), |
| 432 | ProbeResult::ClusterTooSmall { |
| 433 | recv: received, |
| 434 | limit: required, |
| 435 | } => { |
| 436 | write!(f, "cluster too small ({} < {})", received, required) |
| 437 | } |
| 438 | ProbeResult::InsufficientPackets { |
| 439 | recv: received, |
| 440 | limit: required, |
| 441 | } => { |
| 442 | write!(f, "insufficient packets ({} < {})", received, required) |
| 443 | } |
| 444 | ProbeResult::InsufficientBytes { |
| 445 | recv: received, |
| 446 | limit: required, |
| 447 | } => { |
| 448 | write!(f, "insufficient bytes ({} < {})", received, required) |
| 449 | } |
| 450 | ProbeResult::SendIntervalTooLong { interval } => { |
| 451 | write!(f, "send interval too long ({:?})", interval) |
| 452 | } |
| 453 | ProbeResult::SendIntervalInvalid { interval } => { |
| 454 | write!(f, "send interval invalid ({:?})", interval) |
| 455 | } |
| 456 | ProbeResult::RecvIntervalInvalid { interval } => { |
| 457 | write!(f, "recv interval invalid ({:?})", interval) |
| 458 | } |
| 459 | ProbeResult::InvalidSendReceiveRatio { ratio, limit } => { |
| 460 | write!(f, "invalid receive/send ratio ({ratio:.3} > {limit:.3})") |
| 461 | } |
| 462 | ProbeResult::InvalidDataSize => write!(f, "invalid data size"), |
| 463 | ProbeResult::MissingTimingInfo => write!(f, "missing timing info"), |
| 464 | } |
| 465 | } |
| 466 | } |
| 467 | |
| 468 | #[cfg(test)] |
| 469 | mod test { |
| 470 | use super::*; |
| 471 | use crate::bwe_::probe::ProbeKind; |
| 472 | use crate::rtp_::{TwccPacketId, TwccSeq}; |
| 473 | |
| 474 | #[test] |
| 475 | fn probe_estimator_starts_with_no_active_probe() { |
| 476 | let estimator = ProbeEstimator::new(); |
| 477 | assert_eq!(estimator.poll_timeout(), not_happening()); |
| 478 | } |
| 479 | |
| 480 | #[test] |
| 481 | fn probe_estimator_lifecycle() { |
| 482 | let mut estimator = ProbeEstimator::new(); |
| 483 | let now = Instant::now(); |
| 484 | |
| 485 | // Start probe |
| 486 | let config = ProbeClusterConfig::new(1.into(), Bitrate::mbps(2), ProbeKind::Initial); |
| 487 | assert!(estimator.probe_start(config, now)); |
| 488 | assert!(estimator.states.len() == 1, "Should have one active probe"); |
| 489 | assert_eq!(estimator.poll_timeout(), not_happening()); |
| 490 | |
| 491 | // End probe with 1 second cluster history retention |
| 492 | estimator.end_probe(now, config.cluster()); |
| 493 | let timeout = estimator.poll_timeout(); |
| 494 | assert!( |
| 495 | timeout > now && timeout <= now + Duration::from_secs(1), |
| 496 | "Expected timeout between now and now+1s, got: {:?}", |
| 497 | timeout.duration_since(now) |
| 498 | ); |
| 499 | |
| 500 | // Handle timeout clears expired probes |
| 501 | estimator.handle_timeout(now + Duration::from_secs(1)); |
| 502 | assert!(estimator.states.is_empty(), "All probes should be cleared"); |
| 503 | assert_eq!(estimator.poll_timeout(), not_happening()); |
| 504 | } |
| 505 | |
| 506 | #[test] |
| 507 | fn lost_probe_packets_do_not_affect_estimate() { |
| 508 | let mut estimator = ProbeEstimator::new(); |
| 509 | let cluster: TwccClusterId = 7.into(); |
| 510 | let config = ProbeClusterConfig::new(cluster, Bitrate::mbps(2), ProbeKind::Initial); |
| 511 | |
| 512 | let base = Instant::now(); |
| 513 | let received = (0..5).map(|i| { |
| 514 | let seq: TwccSeq = (1000 + i).into(); |
| 515 | let pid = TwccPacketId::with_cluster(seq, cluster); |
| 516 | // send spaced 4ms, recv spaced 4ms (same ordering) |
| 517 | crate::rtp_::TwccSendRecord::test_new( |
| 518 | pid, |
| 519 | base + Duration::from_millis(i as u64 * 4), |
| 520 | 1200, |
| 521 | base + Duration::from_millis(i as u64 * 4 + 1), |
| 522 | Some(base + Duration::from_millis(i as u64 * 4 + 2)), |
| 523 | ) |
| 524 | }); |
| 525 | |
| 526 | // Add extra lost probe packets with later send times. These should not change the result. |
| 527 | let lost = (0..20).map(|i| { |
| 528 | let seq: TwccSeq = (2000 + i).into(); |
| 529 | let pid = TwccPacketId::with_cluster(seq, cluster); |
| 530 | crate::rtp_::TwccSendRecord::test_new( |
| 531 | pid, |
| 532 | base + Duration::from_millis(100 + i as u64), |
| 533 | 1200, |
| 534 | base + Duration::from_millis(150 + i as u64), |
| 535 | None, // lost |
| 536 | ) |
| 537 | }); |
| 538 | |
| 539 | // First run: only received packets |
| 540 | estimator.probe_start(config, base); |
| 541 | let recv_vec: Vec<_> = received.collect(); |
| 542 | let results: Vec<_> = estimator.update(recv_vec.iter()).collect(); |
| 543 | let estimate_only_received = results |
| 544 | .last() |
| 545 | .map(|(_, bitrate)| *bitrate) |
| 546 | .expect("expected a probe estimate"); |
| 547 | |
| 548 | // Second run: received + lost |
| 549 | let mut estimator2 = ProbeEstimator::new(); |
| 550 | estimator2.probe_start(config, base); |
| 551 | let mut all_vec = recv_vec; |
| 552 | all_vec.extend(lost); |
| 553 | let results: Vec<_> = estimator2.update(all_vec.iter()).collect(); |
| 554 | let estimate_with_lost = results |
| 555 | .last() |
| 556 | .map(|(_, bitrate)| *bitrate) |
| 557 | .expect("expected a probe estimate"); |
| 558 | |
| 559 | assert_eq!( |
| 560 | estimate_only_received, estimate_with_lost, |
| 561 | "lost packets must not change probe estimate" |
| 562 | ); |
| 563 | } |
| 564 | |
| 565 | #[test] |
| 566 | fn invalid_receive_send_ratio_is_rejected() { |
| 567 | let mut estimator = ProbeEstimator::new(); |
| 568 | let cluster: TwccClusterId = 9.into(); |
| 569 | let config = ProbeClusterConfig::new(cluster, Bitrate::mbps(2), ProbeKind::Initial); |
| 570 | |
| 571 | let base = Instant::now(); |
| 572 | // Make send times span 200ms, but receive times span only 1ms. |
| 573 | // This yields receive_rate >> send_rate. WebRTC would reject this via ratio check. |
| 574 | let records: Vec<_> = (0..5) |
| 575 | .map(|i| { |
| 576 | let seq: TwccSeq = (3000 + i).into(); |
| 577 | let pid = TwccPacketId::with_cluster(seq, cluster); |
| 578 | crate::rtp_::TwccSendRecord::test_new( |
| 579 | pid, |
| 580 | base + Duration::from_millis(i as u64 * 50), |
| 581 | 1200, |
| 582 | base + Duration::from_millis(250 + i as u64), |
| 583 | Some(base + Duration::from_millis(300 + (i as u64 % 2))), // ~0-1ms spread |
| 584 | ) |
| 585 | }) |
| 586 | .collect(); |
| 587 | |
| 588 | estimator.probe_start(config, base); |
| 589 | let results: Vec<_> = estimator.update(records.iter()).collect(); |
| 590 | |
| 591 | assert!( |
| 592 | results.is_empty(), |
| 593 | "probe should be rejected by ratio validation, got: {:?}", |
| 594 | results |
| 595 | ); |
| 596 | } |
| 597 | |
| 598 | #[test] |
| 599 | fn send_interval_zero_is_rejected() { |
| 600 | let mut estimator = ProbeEstimator::new(); |
| 601 | let cluster: TwccClusterId = 10.into(); |
| 602 | let config = ProbeClusterConfig::new(cluster, Bitrate::mbps(2), ProbeKind::Initial); |
| 603 | |
| 604 | let base = Instant::now(); |
| 605 | // All packets have the same send time -> send_interval == 0. |
| 606 | let records: Vec<_> = (0..5) |
| 607 | .map(|i| { |
| 608 | let seq: TwccSeq = (4000 + i).into(); |
| 609 | let pid = TwccPacketId::with_cluster(seq, cluster); |
| 610 | crate::rtp_::TwccSendRecord::test_new( |
| 611 | pid, |
| 612 | base, // identical send time for all |
| 613 | 1200, |
| 614 | base + Duration::from_millis(10 + i as u64), |
| 615 | Some(base + Duration::from_millis(20 + i as u64)), |
| 616 | ) |
| 617 | }) |
| 618 | .collect(); |
| 619 | |
| 620 | estimator.probe_start(config, base); |
| 621 | let results: Vec<_> = estimator.update(records.iter()).collect(); |
| 622 | |
| 623 | assert!( |
| 624 | results.is_empty(), |
| 625 | "send_interval == 0 should be rejected, got: {:?}", |
| 626 | results |
| 627 | ); |
| 628 | } |
| 629 | |
| 630 | #[test] |
| 631 | fn stale_probes_cleaned_up_at_capacity() { |
| 632 | let mut estimator = ProbeEstimator::new(); |
| 633 | let base = Instant::now(); |
| 634 | |
| 635 | // Fill up to MAX_ACTIVE_PROBES with probes created at `base` |
| 636 | for i in 0..MAX_ACTIVE_PROBES { |
| 637 | let config = |
| 638 | ProbeClusterConfig::new((i as u64).into(), Bitrate::mbps(2), ProbeKind::Initial); |
| 639 | assert!( |
| 640 | estimator.probe_start(config, base), |
| 641 | "should accept probe {i}" |
| 642 | ); |
| 643 | } |
| 644 | assert_eq!(estimator.states.len(), MAX_ACTIVE_PROBES); |
| 645 | |
| 646 | // Try to add another probe at the same time - should be rejected |
| 647 | let config = ProbeClusterConfig::new(100.into(), Bitrate::mbps(2), ProbeKind::Initial); |
| 648 | assert!( |
| 649 | !estimator.probe_start(config, base), |
| 650 | "should reject probe when at capacity with no stale probes" |
| 651 | ); |
| 652 | assert_eq!(estimator.states.len(), MAX_ACTIVE_PROBES); |
| 653 | |
| 654 | // Now try adding a probe 6 seconds later - all existing probes are stale |
| 655 | let later = base + Duration::from_secs(6); |
| 656 | let config = ProbeClusterConfig::new(101.into(), Bitrate::mbps(2), ProbeKind::Initial); |
| 657 | assert!( |
| 658 | estimator.probe_start(config, later), |
| 659 | "should accept probe after cleaning stale ones" |
| 660 | ); |
| 661 | // All old probes should be cleaned, leaving only the new one |
| 662 | assert_eq!(estimator.states.len(), 1); |
| 663 | } |
| 664 | |
| 665 | #[test] |
| 666 | fn non_stale_probes_preserved_during_cleanup() { |
| 667 | let mut estimator = ProbeEstimator::new(); |
| 668 | let base = Instant::now(); |
| 669 | |
| 670 | // Add some old probes |
| 671 | for i in 0..15 { |
| 672 | let config = |
| 673 | ProbeClusterConfig::new((i as u64).into(), Bitrate::mbps(2), ProbeKind::Initial); |
| 674 | estimator.probe_start(config, base); |
| 675 | } |
| 676 | |
| 677 | // Add some newer probes (3 seconds later, still within threshold) |
| 678 | let mid = base + Duration::from_secs(3); |
| 679 | for i in 15..MAX_ACTIVE_PROBES { |
| 680 | let config = |
| 681 | ProbeClusterConfig::new((i as u64).into(), Bitrate::mbps(2), ProbeKind::Initial); |
| 682 | estimator.probe_start(config, mid); |
| 683 | } |
| 684 | assert_eq!(estimator.states.len(), MAX_ACTIVE_PROBES); |
| 685 | |
| 686 | // Try to add at 6 seconds - old probes are stale, mid probes are not |
| 687 | let later = base + Duration::from_secs(6); |
| 688 | let config = ProbeClusterConfig::new(100.into(), Bitrate::mbps(2), ProbeKind::Initial); |
| 689 | assert!( |
| 690 | estimator.probe_start(config, later), |
| 691 | "should accept after cleaning only stale probes" |
| 692 | ); |
| 693 | // 15 old probes removed, 5 mid probes kept, 1 new probe added = 6 |
| 694 | assert_eq!(estimator.states.len(), 6); |
| 695 | } |
| 696 | } |