File
Blob: firmware/vendor/str0m/src/bwe/probe/cluster.rs
| 1 | //! Probe cluster data structures for bandwidth estimation. |
| 2 | //! |
| 3 | //! Probe clusters are short bursts of packets sent at specific bitrates to test |
| 4 | //! network capacity. |
| 5 | |
| 6 | use std::cmp; |
| 7 | use std::time::{Duration, Instant}; |
| 8 | |
| 9 | use crate::rtp_::{Bitrate, DataSize, TwccClusterId}; |
| 10 | use crate::util::{already_happened, not_happening}; |
| 11 | |
| 12 | const MAX_PADDING_PACKET_SIZE: DataSize = DataSize::bytes(240); |
| 13 | |
| 14 | /// Configuration for a probe cluster (the plan). |
| 15 | /// |
| 16 | /// This represents the immutable blueprint for a bandwidth probe: what bitrate |
| 17 | /// to test, for how long, and with what constraints. |
| 18 | #[derive(Debug, Clone, Copy)] |
| 19 | pub struct ProbeClusterConfig { |
| 20 | /// Unique identifier for this probe cluster |
| 21 | cluster: TwccClusterId, |
| 22 | |
| 23 | /// Target bitrate to probe at (e.g., 3 Mbps) |
| 24 | target_bitrate: Bitrate, |
| 25 | |
| 26 | /// How long to sustain the target bitrate (e.g., 15ms) |
| 27 | target_duration: Duration, |
| 28 | |
| 29 | /// Minimum number of packets to send (e.g., 5) |
| 30 | /// This ensures statistical validity even for short bursts. |
| 31 | min_packet_count: usize, |
| 32 | |
| 33 | /// Delta time between sent bursts of packets during probe. |
| 34 | /// |
| 35 | /// Mirrors WebRTC's `ProbeClusterConfig.min_probe_delta`. The pacer/probe sender must not |
| 36 | /// schedule probe packets more frequently than this. |
| 37 | min_probe_delta: Duration, |
| 38 | |
| 39 | /// The kind of probe this is. |
| 40 | kind: ProbeKind, |
| 41 | } |
| 42 | |
| 43 | /// Kind of probe |
| 44 | #[derive(Debug, Clone, Copy, PartialEq, Eq)] |
| 45 | #[non_exhaustive] |
| 46 | pub enum ProbeKind { |
| 47 | Initial, |
| 48 | Exponential, |
| 49 | IncreaseAlr, |
| 50 | PeriodicAlr, |
| 51 | LargeDrop, |
| 52 | Stagnant, |
| 53 | } |
| 54 | |
| 55 | impl ProbeKind { |
| 56 | pub(crate) fn is_alr(&self) -> bool { |
| 57 | matches!(self, ProbeKind::PeriodicAlr | ProbeKind::IncreaseAlr) |
| 58 | } |
| 59 | } |
| 60 | |
| 61 | impl ProbeClusterConfig { |
| 62 | /// Create a new probe cluster configuration with standard defaults: |
| 63 | /// - 15ms duration (enough to get meaningful feedback without excessive delay) |
| 64 | /// - 5 minimum packets (statistical significance for BWE analysis) |
| 65 | pub fn new(cluster: TwccClusterId, target_bitrate: Bitrate, kind: ProbeKind) -> Self { |
| 66 | Self { |
| 67 | cluster, |
| 68 | target_bitrate, |
| 69 | // WebRTC defaults |
| 70 | target_duration: Duration::from_millis(15), |
| 71 | min_packet_count: 5, |
| 72 | // WebRTC default for general probing (not initial/network-state): 2ms. |
| 73 | min_probe_delta: Duration::from_millis(2), |
| 74 | kind, |
| 75 | } |
| 76 | } |
| 77 | |
| 78 | /// Set a custom target duration for this probe. |
| 79 | /// WebRTC uses 100ms for initial probes to allow time for media to start. |
| 80 | pub fn with_duration(mut self, duration: Duration) -> Self { |
| 81 | self.target_duration = duration; |
| 82 | self |
| 83 | } |
| 84 | |
| 85 | /// Set a custom minimum packet count for this probe. |
| 86 | pub fn with_min_packet_count(mut self, min_packet_count: usize) -> Self { |
| 87 | self.min_packet_count = min_packet_count; |
| 88 | self |
| 89 | } |
| 90 | |
| 91 | /// Set a custom minimum probe delta (spacing constraint) for this probe. |
| 92 | pub fn with_min_probe_delta(mut self, min_probe_delta: Duration) -> Self { |
| 93 | self.min_probe_delta = min_probe_delta; |
| 94 | self |
| 95 | } |
| 96 | |
| 97 | /// Check if this probe was created during ALR. |
| 98 | pub fn is_alr_probe(&self) -> bool { |
| 99 | self.kind.is_alr() |
| 100 | } |
| 101 | |
| 102 | /// Get the probe cluster ID. |
| 103 | pub fn cluster(&self) -> TwccClusterId { |
| 104 | self.cluster |
| 105 | } |
| 106 | |
| 107 | /// Get the target bitrate. |
| 108 | pub fn target_bitrate(&self) -> Bitrate { |
| 109 | self.target_bitrate |
| 110 | } |
| 111 | |
| 112 | /// Get the minimum packet count required for a valid probe. |
| 113 | pub fn min_packet_count(&self) -> usize { |
| 114 | self.min_packet_count |
| 115 | } |
| 116 | |
| 117 | /// Get the minimum probe delta (spacing constraint). |
| 118 | pub fn min_probe_delta(&self) -> Duration { |
| 119 | self.min_probe_delta |
| 120 | } |
| 121 | |
| 122 | /// Calculate the target bytes for this probe. |
| 123 | /// This is how much data we expect to send at target_bitrate for target_duration. |
| 124 | pub fn target_bytes(&self) -> DataSize { |
| 125 | self.target_bitrate * self.target_duration |
| 126 | } |
| 127 | } |
| 128 | |
| 129 | /// Runtime state of an active probe cluster (the execution). |
| 130 | /// |
| 131 | /// This tracks what's actually happening as we send probe packets: how much |
| 132 | /// we've sent, how many packets, and when we started. |
| 133 | #[derive(Debug)] |
| 134 | pub struct ProbeClusterState { |
| 135 | /// The immutable plan |
| 136 | config: ProbeClusterConfig, |
| 137 | |
| 138 | /// Total bytes sent so far in this probe |
| 139 | bytes_sent: DataSize, |
| 140 | |
| 141 | /// Total packets sent so far in this probe |
| 142 | packets_sent: usize, |
| 143 | |
| 144 | /// When the first packet was sent |
| 145 | /// None = probe hasn't started yet (still queued) |
| 146 | started_at: Option<Instant>, |
| 147 | |
| 148 | /// When the last packet was sent |
| 149 | /// Used to calculate actual send duration (first to last packet) |
| 150 | last_packet_at: Option<Instant>, |
| 151 | |
| 152 | /// When we last created a padding request |
| 153 | /// Used to prevent creating multiple padding packets for the same instant |
| 154 | last_padding_at: Instant, |
| 155 | } |
| 156 | |
| 157 | impl ProbeClusterState { |
| 158 | /// Create a new probe cluster state from a config. |
| 159 | pub fn new(config: ProbeClusterConfig) -> Self { |
| 160 | Self { |
| 161 | config, |
| 162 | bytes_sent: DataSize::ZERO, |
| 163 | packets_sent: 0, |
| 164 | started_at: None, |
| 165 | last_packet_at: None, |
| 166 | last_padding_at: already_happened(), |
| 167 | } |
| 168 | } |
| 169 | |
| 170 | /// Get the probe configuration. |
| 171 | pub fn config(&self) -> &ProbeClusterConfig { |
| 172 | &self.config |
| 173 | } |
| 174 | |
| 175 | /// Calculates the next probe time based on total bytes sent and target bitrate. |
| 176 | /// This naturally handles variable packet sizes (media packets can be much larger |
| 177 | /// than padding packets). Next packet can be sent when: |
| 178 | /// `now >= start_time + (bytes_sent / target_bitrate)` |
| 179 | pub fn next_probe_time(&self) -> Instant { |
| 180 | let Some(started_at) = self.started_at else { |
| 181 | return already_happened(); |
| 182 | }; |
| 183 | |
| 184 | if self.config.target_bitrate == Bitrate::ZERO { |
| 185 | return not_happening(); |
| 186 | } |
| 187 | |
| 188 | // Packet-level pacing: schedule based on bytes sent at the target bitrate. |
| 189 | // next_time = start_time + (bytes_sent / target_bitrate) |
| 190 | let send_duration = self.bytes_sent / self.config.target_bitrate; |
| 191 | let probe_time = started_at + send_duration; |
| 192 | |
| 193 | let min_burst_time = self.last_padding_at + self.config.min_probe_delta(); |
| 194 | probe_time.max(min_burst_time) |
| 195 | } |
| 196 | |
| 197 | /// Check if the probe cluster is complete. |
| 198 | /// |
| 199 | /// A probe is complete when BOTH conditions are met (matching WebRTC): |
| 200 | /// 1. Sent bytes >= target_bytes (target_bitrate * target_duration) |
| 201 | /// 2. Sent packets >= min_packet_count |
| 202 | /// |
| 203 | /// Note: WebRTC does NOT check duration for completion. Duration is only used |
| 204 | /// to calculate the target bytes threshold. This ensures probes complete even |
| 205 | /// when all packets are sent at the same instant (e.g., in tests or when the |
| 206 | /// pacer generates bursts faster than wall-clock time advances). |
| 207 | pub fn is_complete(&self, _now: Instant) -> bool { |
| 208 | if self.started_at.is_none() { |
| 209 | return false; // Not started yet |
| 210 | } |
| 211 | |
| 212 | // Match WebRTC's BitrateProber::ProbeSent() logic: |
| 213 | // if (sent_bytes >= probe_cluster_min_bytes && sent_probes >= probe_cluster_min_probes) |
| 214 | let bytes_met = self.bytes_sent >= self.config.target_bytes(); |
| 215 | let packets_met = self.packets_sent >= self.config.min_packet_count; |
| 216 | |
| 217 | bytes_met && packets_met |
| 218 | } |
| 219 | |
| 220 | /// Check if it's time to send the next probe packet. |
| 221 | /// |
| 222 | /// Returns `true` if `now >= next_probe_time()`, meaning we should send a packet now. |
| 223 | pub fn should_send_now(&self, now: Instant) -> bool { |
| 224 | let next_time = self.next_probe_time(); |
| 225 | let min_burst_time = self.last_padding_at + self.config.min_probe_delta(); |
| 226 | now >= next_time && now >= min_burst_time |
| 227 | } |
| 228 | |
| 229 | /// Calculate how much padding should be generated for the next probe packet. |
| 230 | /// |
| 231 | /// WebRTC probing uses **bursts** of packets, with a minimum delta between bursts |
| 232 | /// (`min_probe_delta`). In str0m, a `PaddingRequest` is expressed as a number of bytes |
| 233 | /// to queue (which may result in multiple RTP padding packets), so we request a burst |
| 234 | /// sized approximately as: |
| 235 | /// |
| 236 | /// `burst_bytes = target_bitrate * min_probe_delta` |
| 237 | /// |
| 238 | /// This avoids unintentionally capping probe throughput at `MAX_PADDING_PACKET_SIZE / |
| 239 | /// min_probe_delta` (e.g. 240 bytes / 2ms = 960 kbit/s), which would be far below the |
| 240 | /// intended probe bitrate. |
| 241 | /// |
| 242 | /// Returns `None` if it's not time to send yet, or if we've already created |
| 243 | /// a padding packet for this instant. |
| 244 | pub fn next_packet(&mut self, now: Instant) -> Option<DataSize> { |
| 245 | if self.started_at.is_none() { |
| 246 | self.started_at = Some(now); |
| 247 | } |
| 248 | |
| 249 | // Enforce min burst spacing (`min_probe_delta`) between successive padding requests. |
| 250 | // (The very first burst is not gated because `last_padding_at` starts at already_happened()). |
| 251 | if now < self.last_padding_at + self.config.min_probe_delta() { |
| 252 | return None; |
| 253 | } |
| 254 | |
| 255 | if now < self.next_probe_time() { |
| 256 | return None; |
| 257 | } |
| 258 | |
| 259 | // Check if we've already created a padding packet for this exact instant |
| 260 | // This prevents creating multiple padding packets when handle_timeout() is |
| 261 | // called multiple times with the same `now` value |
| 262 | if self.last_padding_at >= now { |
| 263 | return None; |
| 264 | } |
| 265 | |
| 266 | // Mark that we've created padding for this instant |
| 267 | self.last_padding_at = now; |
| 268 | |
| 269 | // It's time to send a new probe burst. |
| 270 | // Calculate recommended probe size as target_bitrate * min_probe_delta. |
| 271 | let min_delta = self.config.min_probe_delta(); |
| 272 | let recommended_probe_size = self.config.target_bitrate * min_delta; |
| 273 | |
| 274 | // Ensure we request at least one padding packet worth of bytes even if min_delta is 0. |
| 275 | let recommended_probe_size = DataSize::bytes(cmp::max( |
| 276 | recommended_probe_size.as_bytes_i64(), |
| 277 | MAX_PADDING_PACKET_SIZE.as_bytes_i64(), |
| 278 | )); |
| 279 | |
| 280 | // Calculate remaining bytes needed to complete the probe cluster. |
| 281 | let bytes_remaining = self.config.target_bytes().saturating_sub(self.bytes_sent); |
| 282 | |
| 283 | // Return the minimum of bytes_remaining and recommended_probe_size. |
| 284 | // When bytes_remaining is zero, this returns None (no more padding). |
| 285 | let request_bytes = cmp::min(bytes_remaining, recommended_probe_size); |
| 286 | |
| 287 | if request_bytes == DataSize::ZERO { |
| 288 | None |
| 289 | } else { |
| 290 | Some(request_bytes) |
| 291 | } |
| 292 | } |
| 293 | |
| 294 | /// Record a packet that was sent as part of this probe (media or padding). |
| 295 | /// |
| 296 | /// This updates the probe's tracking state so that `next_probe_time()` correctly |
| 297 | /// calculates when the next packet should be sent based on the target bitrate. |
| 298 | /// |
| 299 | /// Should be called by the pacer whenever ANY packet is sent during an active probe, |
| 300 | /// not just padding packets generated by `next_packet()`. |
| 301 | pub fn record_packet(&mut self, now: Instant, size: DataSize) { |
| 302 | // If a probe is active and we observe a media packet before any probe-generated padding, |
| 303 | // treat that as the start of the probe. This keeps timing semantics consistent and ensures |
| 304 | // `next_probe_time()` enforces `min_probe_delta` from the first packet onwards. |
| 305 | if self.started_at.is_none() { |
| 306 | self.started_at = Some(now); |
| 307 | } |
| 308 | |
| 309 | self.bytes_sent += size; |
| 310 | self.packets_sent += 1; |
| 311 | self.last_packet_at = Some(now); |
| 312 | } |
| 313 | } |
| 314 | |
| 315 | #[cfg(test)] |
| 316 | mod test { |
| 317 | use super::*; |
| 318 | |
| 319 | // Test helper to directly set state for testing |
| 320 | impl ProbeClusterState { |
| 321 | fn test_set_state(&mut self, bytes_sent: DataSize, packets_sent: usize, now: Instant) { |
| 322 | self.bytes_sent = bytes_sent; |
| 323 | self.packets_sent = packets_sent; |
| 324 | if self.started_at.is_none() { |
| 325 | self.started_at = Some(now); |
| 326 | } |
| 327 | self.last_packet_at = Some(now); |
| 328 | } |
| 329 | } |
| 330 | |
| 331 | #[test] |
| 332 | fn probe_cluster_not_complete_before_start() { |
| 333 | let now = Instant::now(); |
| 334 | let config = ProbeClusterConfig::new(1.into(), Bitrate::mbps(3), ProbeKind::Initial); |
| 335 | let state = ProbeClusterState::new(config); |
| 336 | |
| 337 | assert!(!state.is_complete(now)); |
| 338 | } |
| 339 | |
| 340 | #[test] |
| 341 | fn probe_cluster_not_complete_after_packets_but_not_bytes() { |
| 342 | let now = Instant::now(); |
| 343 | let config = ProbeClusterConfig::new(1.into(), Bitrate::mbps(3), ProbeKind::Initial); |
| 344 | // target_bytes = 3 Mbps * 15ms = 45,000 bits = 5,625 bytes |
| 345 | let mut state = ProbeClusterState::new(config); |
| 346 | |
| 347 | // Send 5 packets (meets packet count) but only small packets (doesn't meet bytes) |
| 348 | // 5 * 100 = 500 bytes, which is < 5,625 bytes |
| 349 | state.test_set_state(DataSize::bytes(500), 5, now); |
| 350 | |
| 351 | // Not complete: packets met but bytes not met |
| 352 | assert!(!state.is_complete(now)); |
| 353 | } |
| 354 | |
| 355 | #[test] |
| 356 | fn probe_cluster_not_complete_after_bytes_but_not_packets() { |
| 357 | let now = Instant::now(); |
| 358 | let config = ProbeClusterConfig::new(1.into(), Bitrate::mbps(3), ProbeKind::Initial); |
| 359 | // target_bytes = 3 Mbps * 15ms = 45,000 bits = 5,625 bytes |
| 360 | let mut state = ProbeClusterState::new(config); |
| 361 | |
| 362 | // Send 2 large packets (meets bytes) but not enough packets |
| 363 | // 2 * 3000 = 6000 bytes > 5,625 bytes, but only 2 packets < 5 packets |
| 364 | state.test_set_state(DataSize::bytes(6000), 2, now); |
| 365 | |
| 366 | // Not complete: bytes met but packets not met |
| 367 | assert!(!state.is_complete(now)); |
| 368 | } |
| 369 | |
| 370 | #[test] |
| 371 | fn probe_cluster_complete_when_both_criteria_met() { |
| 372 | let now = Instant::now(); |
| 373 | let config = ProbeClusterConfig::new(1.into(), Bitrate::mbps(3), ProbeKind::Initial); |
| 374 | // target_bytes = 3 Mbps * 15ms = 45,000 bits = 5,625 bytes |
| 375 | let mut state = ProbeClusterState::new(config); |
| 376 | |
| 377 | // Send 5 packets of 1200 bytes each = 6000 bytes |
| 378 | // Meets both: 6000 >= 5625 bytes AND 5 >= 5 packets |
| 379 | state.test_set_state(DataSize::bytes(6000), 5, now); |
| 380 | |
| 381 | // Complete: both bytes and packets met, even at same instant |
| 382 | assert!(state.is_complete(now)); |
| 383 | } |
| 384 | |
| 385 | #[test] |
| 386 | fn probe_cluster_complete_even_with_zero_duration() { |
| 387 | let now = Instant::now(); |
| 388 | let config = ProbeClusterConfig::new(1.into(), Bitrate::mbps(3), ProbeKind::Initial); |
| 389 | let mut state = ProbeClusterState::new(config); |
| 390 | |
| 391 | // Send all packets at exactly the same instant (duration = 0) |
| 392 | state.test_set_state(DataSize::bytes(6000), 5, now); |
| 393 | |
| 394 | // Should still complete (this is the bug fix - no duration requirement) |
| 395 | assert!(state.is_complete(now)); |
| 396 | } |
| 397 | |
| 398 | #[test] |
| 399 | fn probe_cluster_tracks_bytes_and_packets() { |
| 400 | let now = Instant::now(); |
| 401 | let config = ProbeClusterConfig::new(1.into(), Bitrate::mbps(3), ProbeKind::Initial); |
| 402 | let mut state = ProbeClusterState::new(config); |
| 403 | |
| 404 | state.test_set_state(DataSize::bytes(1200), 1, now); |
| 405 | assert_eq!(state.bytes_sent, DataSize::bytes(1200)); |
| 406 | assert_eq!(state.packets_sent, 1); |
| 407 | |
| 408 | state.test_set_state(DataSize::bytes(2200), 2, now + Duration::from_millis(1)); |
| 409 | assert_eq!(state.bytes_sent, DataSize::bytes(2200)); |
| 410 | assert_eq!(state.packets_sent, 2); |
| 411 | } |
| 412 | |
| 413 | #[test] |
| 414 | fn min_probe_delta_is_enforced_in_next_probe_time() { |
| 415 | let now = Instant::now(); |
| 416 | let config = ProbeClusterConfig::new(1.into(), Bitrate::mbps(3), ProbeKind::Initial) |
| 417 | .with_min_probe_delta(Duration::from_millis(20)); |
| 418 | let mut state = ProbeClusterState::new(config); |
| 419 | assert_eq!(state.config().min_probe_delta(), Duration::from_millis(20)); |
| 420 | |
| 421 | // Start a burst at `now`, then verify that a second burst is blocked until >= now+20ms. |
| 422 | assert!(state.next_packet(now).is_some()); |
| 423 | assert!(state.next_packet(now + Duration::from_millis(19)).is_none()); |
| 424 | assert!(state.next_packet(now + Duration::from_millis(20)).is_some()); |
| 425 | let next = state.last_padding_at; |
| 426 | assert!( |
| 427 | next >= now + Duration::from_millis(20), |
| 428 | "next_probe_time {next:?} must be >= now + min_probe_delta" |
| 429 | ); |
| 430 | } |
| 431 | } |