Skip to content
File

Blob: firmware/vendor/str0m/src/bwe/probe/cluster.rs

rust432 lines
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 
6use std::cmp;
7use std::time::{Duration, Instant};
8 
9use crate::rtp_::{Bitrate, DataSize, TwccClusterId};
10use crate::util::{already_happened, not_happening};
11 
12const 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)]
19pub 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]
46pub enum ProbeKind {
47 Initial,
48 Exponential,
49 IncreaseAlr,
50 PeriodicAlr,
51 LargeDrop,
52 Stagnant,
53}
54 
55impl ProbeKind {
56 pub(crate) fn is_alr(&self) -> bool {
57 matches!(self, ProbeKind::PeriodicAlr | ProbeKind::IncreaseAlr)
58 }
59}
60 
61impl 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)]
134pub 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 
157impl 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)]
316mod 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}