Skip to content
File

Blob: firmware/vendor/str0m/tests/rtp-sequence.rs

rust547 lines
1//! Tests for RTP sequence number and timing edge cases.
2 
3use std::net::Ipv4Addr;
4use std::time::Instant;
5 
6use str0m::media::{Direction, MediaKind};
7use str0m::rtp::{RtpWrite, Ssrc};
8use str0m::{Event, Rtc, RtcError};
9 
10mod common;
11use common::{Peer, TestRtc, init_log, negotiate, progress};
12use common::{connect_l_r_with_rtc, init_crypto_default};
13 
14/// Test handling of packets crossing the u16 sequence number boundary (65535 -> 0).
15#[test]
16fn rtp_sequence_number_near_boundary() -> Result<(), RtcError> {
17 init_log();
18 init_crypto_default();
19 
20 // Use RTP mode to control sequence numbers directly
21 let now = Instant::now();
22 let rtc1 = Rtc::builder().set_rtp_mode(true).build(now);
23 let rtc2 = Rtc::builder().build(now);
24 
25 let (mut l, mut r) = connect_l_r_with_rtc(rtc1, rtc2);
26 
27 let mid = "audio".into();
28 let ssrc_tx: Ssrc = 1337.into();
29 
30 l.direct_api().declare_media(mid, MediaKind::Audio);
31 l.direct_api().declare_stream_tx(ssrc_tx, None, mid, None);
32 r.direct_api().declare_media(mid, MediaKind::Audio);
33 
34 let max = l.last.max(r.last);
35 l.last = max;
36 r.last = max;
37 
38 let params = l.params_opus();
39 let pt = params.pt();
40 let ssrc = l.direct_api().stream_tx_by_mid(mid, None).unwrap().ssrc();
41 
42 // Start sequence numbers near the u16 boundary (65535)
43 // Send packets: 65530, 65531, 65532, 65533, 65534, 65535, 0, 1, 2, 3, 4, 5
44 let start_seq: u64 = 65530;
45 let packet_count = 12;
46 
47 for i in 0..packet_count {
48 let seq_no = start_seq + i;
49 let time = (i * 960) as u32; // Opus uses 48kHz, 20ms = 960 samples
50 let wallclock = l.start + l.duration();
51 
52 {
53 let mut direct = l.direct_api();
54 let stream = direct.stream_tx(&ssrc).unwrap();
55 
56 stream.write_rtp(RtpWrite::new(
57 pt,
58 seq_no.into(),
59 time,
60 wallclock,
61 [1_u8; 80],
62 ));
63 }
64 progress(&mut l, &mut r)?;
65 }
66 
67 // Final progress to deliver remaining packets
68 for _ in 0..50 {
69 progress(&mut l, &mut r)?;
70 }
71 
72 // Collect received sequence numbers from seq_range
73 let received_seqs: Vec<u64> = r
74 .events
75 .iter()
76 .filter_map(|(_, e)| {
77 if let Event::MediaData(data) = e {
78 // Get the start of the sequence range (SeqNo derefs to u64)
79 Some(**data.seq_range.start())
80 } else {
81 None
82 }
83 })
84 .collect();
85 
86 // Verify we received packets across the boundary
87 assert!(
88 received_seqs.len() >= 10,
89 "Should receive most packets, got {}",
90 received_seqs.len()
91 );
92 
93 // The library uses extended 64-bit sequence numbers internally.
94 // Wire seq 65535 -> internal 65535, wire seq 0 (after wrap) -> internal 65536
95 // Verify sequence numbers before boundary (65530-65535) were received
96 let before_boundary: Vec<_> = received_seqs
97 .iter()
98 .filter(|&&s| (65530..=65535).contains(&s))
99 .collect();
100 assert!(
101 !before_boundary.is_empty(),
102 "Should receive packets before boundary (65530-65535), got seq_nos: {:?}",
103 received_seqs
104 );
105 
106 // Verify sequence numbers after boundary (65536+ = wire 0+) were received
107 // Internal 65536 = wire seq 0, 65537 = wire seq 1, etc.
108 let after_boundary: Vec<_> = received_seqs.iter().filter(|&&s| s >= 65536).collect();
109 assert!(
110 !after_boundary.is_empty(),
111 "Should receive packets after boundary (65536+ = wire 0+), got seq_nos: {:?}",
112 received_seqs
113 );
114 
115 // Verify sequence numbers are monotonically increasing (no gaps from wrap-around)
116 for i in 1..received_seqs.len() {
117 assert!(
118 received_seqs[i] > received_seqs[i - 1],
119 "Sequence numbers should be monotonically increasing across boundary: {} should be > {}",
120 received_seqs[i],
121 received_seqs[i - 1]
122 );
123 }
124 
125 Ok(())
126}
127 
128/// Test reordering buffer with audio - send packets out of order and verify reordering.
129#[test]
130fn rtp_reordering_buffer_audio() -> Result<(), RtcError> {
131 init_log();
132 init_crypto_default();
133 
134 // Sender uses RTP mode to control sequence numbers
135 let now = Instant::now();
136 let rtc1 = Rtc::builder().set_rtp_mode(true).build(now);
137 // Receiver has reordering buffer enabled (default is 15 for audio)
138 let rtc2 = Rtc::builder().set_reordering_size_audio(15).build(now);
139 
140 let (mut l, mut r) = connect_l_r_with_rtc(rtc1, rtc2);
141 
142 let mid = "audio".into();
143 let ssrc_tx: Ssrc = 1337.into();
144 
145 l.direct_api().declare_media(mid, MediaKind::Audio);
146 l.direct_api().declare_stream_tx(ssrc_tx, None, mid, None);
147 r.direct_api().declare_media(mid, MediaKind::Audio);
148 
149 let max = l.last.max(r.last);
150 l.last = max;
151 r.last = max;
152 
153 let params = l.params_opus();
154 let pt = params.pt();
155 let ssrc = l.direct_api().stream_tx_by_mid(mid, None).unwrap().ssrc();
156 
157 // Send packets OUT OF ORDER: 1, 3, 2, 5, 4, 7, 6, 9, 8, 10
158 // The reordering buffer should fix this
159 let send_order = [1u64, 3, 2, 5, 4, 7, 6, 9, 8, 10];
160 
161 for &seq in &send_order {
162 let time = (seq * 960) as u32; // Opus 48kHz, 20ms frames
163 let wallclock = l.start + l.duration();
164 {
165 let mut direct = l.direct_api();
166 let stream = direct.stream_tx(&ssrc).unwrap();
167 
168 stream.write_rtp(RtpWrite::new(
169 pt,
170 seq.into(),
171 time,
172 wallclock,
173 [seq as u8; 80],
174 ));
175 }
176 progress(&mut l, &mut r)?;
177 }
178 
179 // Final progress to flush reordering buffer
180 for _ in 0..50 {
181 progress(&mut l, &mut r)?;
182 }
183 
184 // Collect received sequence number ranges
185 let received_ranges: Vec<(u64, u64)> = r
186 .events
187 .iter()
188 .filter_map(|(_, e)| {
189 if let Event::MediaData(data) = e {
190 Some((**data.seq_range.start(), **data.seq_range.end()))
191 } else {
192 None
193 }
194 })
195 .collect();
196 
197 assert!(!received_ranges.is_empty(), "Should receive some packets");
198 
199 // Verify ranges are non-decreasing (reordering buffer should fix order)
200 // Each range's start should be >= previous range's start
201 for i in 1..received_ranges.len() {
202 assert!(
203 received_ranges[i].0 >= received_ranges[i - 1].0,
204 "Packets should be reordered: range {:?} should come after {:?}, all ranges: {:?}",
205 received_ranges[i],
206 received_ranges[i - 1],
207 received_ranges
208 );
209 }
210 
211 // Verify we received packets spanning our sent range (1-10)
212 let min_received = received_ranges.iter().map(|r| r.0).min().unwrap();
213 let max_received = received_ranges.iter().map(|r| r.1).max().unwrap();
214 assert!(
215 min_received <= 2,
216 "Should receive early packets, min was {}",
217 min_received
218 );
219 assert!(
220 max_received >= 9,
221 "Should receive late packets, max was {}",
222 max_received
223 );
224 
225 Ok(())
226}
227 
228/// Test reordering buffer with video - send packets out of order and verify reordering.
229#[test]
230fn rtp_reordering_buffer_video() -> Result<(), RtcError> {
231 init_log();
232 init_crypto_default();
233 
234 // Sender uses RTP mode to control sequence numbers
235 let now = Instant::now();
236 let rtc1 = Rtc::builder().set_rtp_mode(true).build(now);
237 // Receiver has reordering buffer for video
238 let rtc2 = Rtc::builder().set_reordering_size_video(30).build(now);
239 
240 let (mut l, mut r) = connect_l_r_with_rtc(rtc1, rtc2);
241 
242 let mid = "video".into();
243 let ssrc_tx: Ssrc = 1337.into();
244 
245 l.direct_api().declare_media(mid, MediaKind::Video);
246 l.direct_api().declare_stream_tx(ssrc_tx, None, mid, None);
247 r.direct_api().declare_media(mid, MediaKind::Video);
248 
249 let max = l.last.max(r.last);
250 l.last = max;
251 r.last = max;
252 
253 let params = l.params_vp8();
254 let pt = params.pt();
255 let ssrc = l.direct_api().stream_tx_by_mid(mid, None).unwrap().ssrc();
256 
257 // Send video packets OUT OF ORDER: 1, 3, 2, 5, 4, 7, 6, 9, 8, 10
258 let send_order = [1u64, 3, 2, 5, 4, 7, 6, 9, 8, 10];
259 
260 for &seq in &send_order {
261 let time = (seq * 3000) as u32; // 90kHz video clock, ~33ms frames
262 let wallclock = l.start + l.duration();
263 {
264 let mut direct = l.direct_api();
265 let stream = direct.stream_tx(&ssrc).unwrap();
266 
267 // VP8 keyframe header
268 stream.write_rtp(
269 RtpWrite::new(
270 pt,
271 seq.into(),
272 time,
273 wallclock,
274 [0x10, 0x00, 0x00, seq as u8],
275 )
276 .marker(true),
277 );
278 }
279 progress(&mut l, &mut r)?;
280 }
281 
282 // Final progress to flush reordering buffer
283 for _ in 0..50 {
284 progress(&mut l, &mut r)?;
285 }
286 
287 // Collect received sequence number ranges
288 let received_ranges: Vec<(u64, u64)> = r
289 .events
290 .iter()
291 .filter_map(|(_, e)| {
292 if let Event::MediaData(data) = e {
293 Some((**data.seq_range.start(), **data.seq_range.end()))
294 } else {
295 None
296 }
297 })
298 .collect();
299 
300 assert!(
301 !received_ranges.is_empty(),
302 "Should receive some video packets"
303 );
304 
305 // Verify ranges are non-decreasing (reordering buffer should fix order)
306 for i in 1..received_ranges.len() {
307 assert!(
308 received_ranges[i].0 >= received_ranges[i - 1].0,
309 "Video packets should be reordered: {:?} should come after {:?}",
310 received_ranges[i],
311 received_ranges[i - 1]
312 );
313 }
314 
315 Ok(())
316}
317 
318/// Test custom reordering buffer size - small buffer handles small gaps.
319#[test]
320fn rtp_reordering_buffer_custom_size() -> Result<(), RtcError> {
321 init_log();
322 init_crypto_default();
323 
324 // Sender uses RTP mode
325 let now = Instant::now();
326 let rtc1 = Rtc::builder().set_rtp_mode(true).build(now);
327 // Receiver has small reordering buffer (5 packets)
328 let rtc2 = Rtc::builder().set_reordering_size_audio(5).build(now);
329 
330 let (mut l, mut r) = connect_l_r_with_rtc(rtc1, rtc2);
331 
332 let mid = "audio".into();
333 let ssrc_tx: Ssrc = 1337.into();
334 
335 l.direct_api().declare_media(mid, MediaKind::Audio);
336 l.direct_api().declare_stream_tx(ssrc_tx, None, mid, None);
337 r.direct_api().declare_media(mid, MediaKind::Audio);
338 
339 let max = l.last.max(r.last);
340 l.last = max;
341 r.last = max;
342 
343 let params = l.params_opus();
344 let pt = params.pt();
345 let ssrc = l.direct_api().stream_tx_by_mid(mid, None).unwrap().ssrc();
346 
347 // Send packets with small gaps (within buffer size of 5)
348 // Order: 1, 3, 2, 4, 6, 5, 7, 9, 8, 10 (max gap of 2)
349 let send_order = [1u64, 3, 2, 4, 6, 5, 7, 9, 8, 10];
350 
351 for &seq in &send_order {
352 let time = (seq * 960) as u32;
353 let wallclock = l.start + l.duration();
354 {
355 let mut direct = l.direct_api();
356 let stream = direct.stream_tx(&ssrc).unwrap();
357 
358 stream.write_rtp(RtpWrite::new(
359 pt,
360 seq.into(),
361 time,
362 wallclock,
363 [seq as u8; 80],
364 ));
365 }
366 progress(&mut l, &mut r)?;
367 }
368 
369 for _ in 0..50 {
370 progress(&mut l, &mut r)?;
371 }
372 
373 let received_ranges: Vec<(u64, u64)> = r
374 .events
375 .iter()
376 .filter_map(|(_, e)| {
377 if let Event::MediaData(data) = e {
378 Some((**data.seq_range.start(), **data.seq_range.end()))
379 } else {
380 None
381 }
382 })
383 .collect();
384 
385 assert!(
386 !received_ranges.is_empty(),
387 "Should receive packets with small reordering buffer"
388 );
389 
390 // Verify ordering is correct
391 for i in 1..received_ranges.len() {
392 assert!(
393 received_ranges[i].0 >= received_ranges[i - 1].0,
394 "Packets should be reordered with small buffer"
395 );
396 }
397 
398 Ok(())
399}
400 
401/// Test media time increases correctly.
402#[test]
403fn rtp_media_time_increasing() -> Result<(), RtcError> {
404 init_log();
405 init_crypto_default();
406 
407 let mut l = TestRtc::new(Peer::Left);
408 let mut r = TestRtc::new(Peer::Right);
409 
410 l.add_host_candidate((Ipv4Addr::new(1, 1, 1, 1), 1000).into());
411 r.add_host_candidate((Ipv4Addr::new(2, 2, 2, 2), 2000).into());
412 
413 let mid = negotiate(&mut l, &mut r, |change| {
414 change.add_media(MediaKind::Audio, Direction::SendRecv, None, None, None)
415 });
416 
417 loop {
418 if l.is_connected() && r.is_connected() {
419 break;
420 }
421 progress(&mut l, &mut r)?;
422 }
423 
424 let max = l.last.max(r.last);
425 l.last = max;
426 r.last = max;
427 
428 let params = l.params_opus();
429 let pt = params.pt();
430 let data = vec![1_u8; 80];
431 
432 // Send packets with increasing timestamps
433 for _ in 0..50 {
434 let wallclock = l.start + l.duration();
435 let time = l.duration().into();
436 l.writer(mid)
437 .unwrap()
438 .write(pt, wallclock, time, data.clone())?;
439 progress(&mut l, &mut r)?;
440 }
441 
442 // Verify we received packets (media time handling worked)
443 let received_count = r
444 .events
445 .iter()
446 .filter(|(_, e)| matches!(e, Event::MediaData(_)))
447 .count();
448 
449 assert!(
450 received_count > 20,
451 "Should receive packets with increasing timestamps, got {}",
452 received_count
453 );
454 
455 Ok(())
456}
457 
458/// Test large reordering buffer - handles larger out-of-order gaps.
459#[test]
460fn rtp_reordering_buffer_large() -> Result<(), RtcError> {
461 init_log();
462 init_crypto_default();
463 
464 // Sender uses RTP mode
465 let now = Instant::now();
466 let rtc1 = Rtc::builder().set_rtp_mode(true).build(now);
467 // Receiver has large reordering buffer (50 packets)
468 let rtc2 = Rtc::builder().set_reordering_size_audio(50).build(now);
469 
470 let (mut l, mut r) = connect_l_r_with_rtc(rtc1, rtc2);
471 
472 let mid = "audio".into();
473 let ssrc_tx: Ssrc = 1337.into();
474 
475 l.direct_api().declare_media(mid, MediaKind::Audio);
476 l.direct_api().declare_stream_tx(ssrc_tx, None, mid, None);
477 r.direct_api().declare_media(mid, MediaKind::Audio);
478 
479 let max = l.last.max(r.last);
480 l.last = max;
481 r.last = max;
482 
483 let params = l.params_opus();
484 let pt = params.pt();
485 let ssrc = l.direct_api().stream_tx_by_mid(mid, None).unwrap().ssrc();
486 
487 // Send packets with larger gaps that require big buffer
488 // Send 1, 10, 2, 11, 3, 12, 4, 13... (interleaved with gap of ~9)
489 // Shuffle to create out-of-order: 1, 11, 2, 12, 3, 13...
490 let mut send_order = Vec::new();
491 for i in 0_u64..10 {
492 send_order.push(1 + i);
493 send_order.push(11 + i);
494 }
495 
496 for seq in send_order {
497 let time = (seq * 960) as u32;
498 let wallclock = l.start + l.duration();
499 {
500 let mut direct = l.direct_api();
501 let stream = direct.stream_tx(&ssrc).unwrap();
502 
503 stream.write_rtp(RtpWrite::new(
504 pt,
505 seq.into(),
506 time,
507 wallclock,
508 [seq as u8; 80],
509 ));
510 }
511 progress(&mut l, &mut r)?;
512 }
513 
514 for _ in 0..50 {
515 progress(&mut l, &mut r)?;
516 }
517 
518 let received_ranges: Vec<(u64, u64)> = r
519 .events
520 .iter()
521 .filter_map(|(_, e)| {
522 if let Event::MediaData(data) = e {
523 Some((**data.seq_range.start(), **data.seq_range.end()))
524 } else {
525 None
526 }
527 })
528 .collect();
529 
530 assert!(
531 !received_ranges.is_empty(),
532 "Should receive packets with large reordering buffer"
533 );
534 
535 // Verify ordering is correct despite large gaps
536 for i in 1..received_ranges.len() {
537 assert!(
538 received_ranges[i].0 >= received_ranges[i - 1].0,
539 "Packets should be reordered with large buffer: {:?} vs {:?}",
540 received_ranges[i],
541 received_ranges[i - 1]
542 );
543 }
544 
545 Ok(())
546}