Skip to content
File

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

rust438 lines
1use std::mem;
2use std::time::{Duration, Instant};
3 
4use crate::rtp_::TwccSeq;
5 
6use super::super::AckedPacket;
7use super::super::time::{TimeDelta, Timestamp};
8 
9const BURST_TIME_INTERVAL: Duration = Duration::from_millis(5);
10const SEND_TIME_GROUP_LENGTH: Duration = Duration::from_millis(5);
11const MAX_BURST_DURATION: Duration = Duration::from_millis(100);
12 
13#[derive(Debug, Default)]
14pub struct ArrivalGroup {
15 first: Option<(TwccSeq, Instant, Instant)>,
16 last_seq_no: Option<TwccSeq>,
17 last_local_send_time: Option<Instant>,
18 last_remote_recv_time: Option<Instant>,
19 size: usize,
20}
21 
22impl ArrivalGroup {
23 /// Maybe add a packet to the group.
24 ///
25 /// Returns [`true`] if a new group needs to be created and [`false`] otherwise.
26 fn add_packet(&mut self, packet: &AckedPacket) -> bool {
27 match self.belongs_to_group(packet) {
28 Belongs::NewGroup => return true,
29 Belongs::Skipped => return false,
30 Belongs::Yes => {}
31 }
32 
33 if self.first.is_none() {
34 self.first = Some((
35 packet.seq_no,
36 packet.local_send_time,
37 packet.remote_recv_time,
38 ));
39 }
40 
41 self.last_remote_recv_time = self
42 .last_remote_recv_time
43 .max(Some(packet.remote_recv_time));
44 self.last_local_send_time = self.last_local_send_time.max(Some(packet.local_send_time));
45 self.size += 1;
46 self.last_seq_no = self.last_seq_no.max(Some(packet.seq_no));
47 
48 false
49 }
50 
51 fn belongs_to_group(&self, packet: &AckedPacket) -> Belongs {
52 let Some((_, first_local_send_time, first_remote_recv_time)) = self.first else {
53 // Start of the group
54 return Belongs::Yes;
55 };
56 
57 let Some(first_send_delta) = packet
58 .local_send_time
59 .checked_duration_since(first_local_send_time)
60 else {
61 // Out of order
62 return Belongs::Skipped;
63 };
64 
65 let send_time_delta = Timestamp::from(packet.local_send_time) - self.local_send_time();
66 if send_time_delta == TimeDelta::ZERO {
67 return Belongs::Yes;
68 }
69 let arrival_time_delta = Timestamp::from(packet.remote_recv_time) - self.remote_recv_time();
70 
71 let propagation_delta = arrival_time_delta - send_time_delta;
72 if propagation_delta < TimeDelta::ZERO
73 && arrival_time_delta <= BURST_TIME_INTERVAL
74 && packet.remote_recv_time - first_remote_recv_time < MAX_BURST_DURATION
75 {
76 Belongs::Yes
77 } else if first_send_delta > SEND_TIME_GROUP_LENGTH {
78 Belongs::NewGroup
79 } else {
80 Belongs::Yes
81 }
82 }
83 
84 /// Calculate the send time delta between self and a subsequent group.
85 fn departure_delta(&self, other: &Self) -> TimeDelta {
86 Timestamp::from(other.local_send_time()) - self.local_send_time()
87 }
88 
89 /// Calculate the remote receive time delta between self and a subsequent group.
90 fn arrival_delta(&self, other: &Self) -> TimeDelta {
91 Timestamp::from(other.remote_recv_time()) - self.remote_recv_time()
92 }
93 
94 /// The local send time i.e. departure time, for the group.
95 ///
96 /// Panics if the group doesn't have at least one packet.
97 fn local_send_time(&self) -> Instant {
98 self.last_local_send_time
99 .expect("local_send_time to only be called on non-empty groups")
100 }
101 
102 /// The remote receive time i.e. arrival time, for the group.
103 ///
104 /// Panics if the group doesn't have at least one packet.
105 fn remote_recv_time(&self) -> Instant {
106 self.last_remote_recv_time
107 .expect("remote_recv_time to only be called on non-empty groups")
108 }
109}
110 
111/// Whether a given packet is belongs to a group or not.
112#[derive(Debug, Clone, Copy, PartialEq, Eq)]
113enum Belongs {
114 /// The packet is belongs to the group.
115 Yes,
116 /// The packet is does not belong to the group, a new group should be created.
117 NewGroup,
118 /// The packet was skipped and a decision wasn't made.
119 Skipped,
120}
121 
122impl Belongs {
123 #[cfg(test)]
124 fn new_group(&self) -> bool {
125 matches!(self, Self::NewGroup)
126 }
127}
128 
129#[derive(Debug, Default)]
130pub struct ArrivalGroupAccumulator {
131 previous_group: Option<ArrivalGroup>,
132 current_group: ArrivalGroup,
133}
134 
135impl ArrivalGroupAccumulator {
136 ///
137 /// Accumulate a packet.
138 ///
139 /// If adding this packet produced a new delay delta it is returned.
140 pub fn accumulate_packet(&mut self, packet: &AckedPacket) -> Option<InterGroupDelayDelta> {
141 let need_new_group = self.current_group.add_packet(packet);
142 
143 if !need_new_group {
144 return None;
145 }
146 
147 // Variation between previous group and current.
148 let arrival_delta = self.arrival_delta();
149 let send_delta = self.send_delta();
150 let last_remote_recv_time = self.current_group.remote_recv_time();
151 
152 let current_group = mem::take(&mut self.current_group);
153 self.previous_group = Some(current_group);
154 
155 self.current_group.add_packet(packet);
156 
157 Some(InterGroupDelayDelta {
158 send_delta: send_delta?,
159 arrival_delta: arrival_delta?,
160 last_remote_recv_time,
161 })
162 }
163 
164 fn arrival_delta(&self) -> Option<TimeDelta> {
165 self.previous_group
166 .as_ref()
167 .map(|prev| prev.arrival_delta(&self.current_group))
168 }
169 
170 fn send_delta(&self) -> Option<TimeDelta> {
171 self.previous_group
172 .as_ref()
173 .map(|prev| prev.departure_delta(&self.current_group))
174 }
175}
176 
177/// The calculate delay delta between two groups of packets.
178#[derive(Debug, Clone, Copy, PartialEq, Eq)]
179pub struct InterGroupDelayDelta {
180 /// The delta between the send times of the two groups i.e. delta between the last packet sent
181 /// in each group.
182 pub send_delta: TimeDelta,
183 /// The delta between the remote arrival times of the two groups.
184 pub arrival_delta: TimeDelta,
185 /// The reported receive time for the last packet in the first arrival group.
186 pub last_remote_recv_time: Instant,
187}
188 
189#[cfg(test)]
190mod test {
191 use std::time::{Duration, Instant};
192 
193 use crate::rtp_::DataSize;
194 
195 use super::{AckedPacket, ArrivalGroup, ArrivalGroupAccumulator, Belongs, TimeDelta};
196 
197 #[test]
198 fn test_arrival_group_all_packets_belong_to_empty_group() {
199 let now = Instant::now();
200 let group = ArrivalGroup::default();
201 
202 assert_eq!(
203 group.belongs_to_group(&AckedPacket {
204 seq_no: 1.into(),
205 size: DataSize::ZERO,
206 local_send_time: now,
207 remote_recv_time: now + duration_us(10),
208 local_recv_time: now + duration_us(12),
209 }),
210 Belongs::Yes,
211 "Any packet should belong to an empty arrival group"
212 );
213 }
214 
215 #[test]
216 fn test_arrival_group_all_packets_sent_within_burst_interval_belong() {
217 let now = Instant::now();
218 #[allow(clippy::vec_init_then_push)]
219 let packets = {
220 let mut packets = vec![];
221 
222 packets.push(AckedPacket {
223 seq_no: 0.into(),
224 size: DataSize::ZERO,
225 local_send_time: now,
226 remote_recv_time: now + duration_us(150),
227 local_recv_time: now + duration_us(200),
228 });
229 
230 packets.push(AckedPacket {
231 seq_no: 1.into(),
232 size: DataSize::ZERO,
233 local_send_time: now + duration_us(50),
234 remote_recv_time: now + duration_us(225),
235 local_recv_time: now + duration_us(275),
236 });
237 
238 packets.push(AckedPacket {
239 seq_no: 2.into(),
240 size: DataSize::ZERO,
241 local_send_time: now + duration_us(1005),
242 remote_recv_time: now + duration_us(1140),
243 local_recv_time: now + duration_us(1190),
244 });
245 
246 packets.push(AckedPacket {
247 seq_no: 3.into(),
248 size: DataSize::ZERO,
249 local_send_time: now + duration_us(4995),
250 remote_recv_time: now + duration_us(5001),
251 local_recv_time: now + duration_us(5051),
252 });
253 
254 // Should not belong
255 packets.push(AckedPacket {
256 seq_no: 4.into(),
257 size: DataSize::ZERO,
258 local_send_time: now + duration_us(5700),
259 remote_recv_time: now + duration_us(6000),
260 local_recv_time: now + duration_us(5750),
261 });
262 
263 packets
264 };
265 
266 let mut group = ArrivalGroup::default();
267 
268 for p in packets {
269 let need_new_group = group.belongs_to_group(&p).new_group();
270 if !need_new_group {
271 group.add_packet(&p);
272 }
273 }
274 
275 assert_eq!(group.size, 4, "Expected group to contain 4 packets");
276 }
277 
278 #[test]
279 fn test_arrival_group_out_order_arrival_ignored() {
280 let now = Instant::now();
281 #[allow(clippy::vec_init_then_push)]
282 let packets = {
283 let mut packets = vec![];
284 
285 packets.push(AckedPacket {
286 seq_no: 0.into(),
287 size: DataSize::ZERO,
288 local_send_time: now,
289 remote_recv_time: now + duration_us(150),
290 local_recv_time: now + duration_us(200),
291 });
292 
293 packets.push(AckedPacket {
294 seq_no: 1.into(),
295 size: DataSize::ZERO,
296 local_send_time: now + duration_us(50),
297 remote_recv_time: now + duration_us(225),
298 local_recv_time: now + duration_us(275),
299 });
300 
301 packets.push(AckedPacket {
302 seq_no: 2.into(),
303 size: DataSize::ZERO,
304 local_send_time: now + duration_us(1005),
305 remote_recv_time: now + duration_us(1140),
306 local_recv_time: now + duration_us(1190),
307 });
308 
309 packets.push(AckedPacket {
310 seq_no: 3.into(),
311 size: DataSize::ZERO,
312 local_send_time: now + duration_us(4995),
313 remote_recv_time: now + duration_us(5001),
314 local_recv_time: now + duration_us(5051),
315 });
316 
317 // Should be skipped
318 packets.push(AckedPacket {
319 seq_no: 4.into(),
320 size: DataSize::ZERO,
321 local_send_time: now - duration_us(100),
322 remote_recv_time: now + duration_us(5000),
323 local_recv_time: now + duration_us(5050),
324 });
325 
326 // Should not belong
327 packets.push(AckedPacket {
328 seq_no: 5.into(),
329 size: DataSize::ZERO,
330 local_send_time: now + duration_us(5700),
331 remote_recv_time: now + duration_us(6000),
332 local_recv_time: now + duration_us(6050),
333 });
334 
335 packets
336 };
337 
338 let mut group = ArrivalGroup::default();
339 
340 for p in packets {
341 let need_new_group = group.belongs_to_group(&p).new_group();
342 if !need_new_group {
343 group.add_packet(&p);
344 }
345 }
346 
347 assert_eq!(group.size, 4, "Expected group to contain 4 packets");
348 }
349 
350 #[test]
351 fn test_arrival_group_arrival_membership() {
352 let now = Instant::now();
353 #[allow(clippy::vec_init_then_push)]
354 let packets = {
355 let mut packets = vec![];
356 
357 packets.push(AckedPacket {
358 seq_no: 0.into(),
359 size: DataSize::ZERO,
360 local_send_time: now,
361 remote_recv_time: now + duration_us(150),
362 local_recv_time: now + duration_us(200),
363 });
364 
365 packets.push(AckedPacket {
366 seq_no: 1.into(),
367 size: DataSize::ZERO,
368 local_send_time: now + duration_us(50),
369 remote_recv_time: now + duration_us(225),
370 local_recv_time: now + duration_us(275),
371 });
372 
373 packets.push(AckedPacket {
374 seq_no: 2.into(),
375 size: DataSize::ZERO,
376 local_send_time: now + duration_us(5152),
377 // Just less than 5ms inter arrival delta
378 remote_recv_time: now + duration_us(5224),
379 local_recv_time: now + duration_us(5274),
380 });
381 
382 // Should not belong
383 packets.push(AckedPacket {
384 seq_no: 3.into(),
385 size: DataSize::ZERO,
386 local_send_time: now + duration_us(5700),
387 remote_recv_time: now + duration_us(6000),
388 local_recv_time: now + duration_us(6050),
389 });
390 
391 packets
392 };
393 
394 let mut group = ArrivalGroup::default();
395 
396 for p in packets {
397 let need_new_group = group.belongs_to_group(&p).new_group();
398 if !need_new_group {
399 group.add_packet(&p);
400 }
401 }
402 
403 assert_eq!(group.size, 3, "Expected group to contain 4 packets");
404 }
405 
406 #[test]
407 fn group_reorder() {
408 let data = vec![
409 ((Duration::from_millis(0), Duration::from_millis(0)), None),
410 ((Duration::from_millis(60), Duration::from_millis(5)), None),
411 ((Duration::from_millis(40), Duration::from_millis(10)), None),
412 (
413 (Duration::from_millis(70), Duration::from_millis(20)),
414 Some((TimeDelta::from_millis(-20), TimeDelta::from_millis(5))),
415 ),
416 ];
417 
418 let now = Instant::now();
419 let mut aga = ArrivalGroupAccumulator::default();
420 
421 for ((local_send_time, remote_recv_time), deltas) in data {
422 let group_delta = aga.accumulate_packet(&AckedPacket {
423 seq_no: Default::default(),
424 size: Default::default(),
425 local_send_time: now + local_send_time,
426 remote_recv_time: now + remote_recv_time,
427 local_recv_time: Instant::now(), // does not matter
428 });
429 
430 assert_eq!(group_delta.map(|d| (d.send_delta, d.arrival_delta)), deltas);
431 }
432 }
433 
434 fn duration_us(us: u64) -> Duration {
435 Duration::from_micros(us)
436 }
437}