File
Blob: firmware/vendor/str0m/tests/rtp-sequence.rs
| 1 | //! Tests for RTP sequence number and timing edge cases. |
| 2 | |
| 3 | use std::net::Ipv4Addr; |
| 4 | use std::time::Instant; |
| 5 | |
| 6 | use str0m::media::{Direction, MediaKind}; |
| 7 | use str0m::rtp::{RtpWrite, Ssrc}; |
| 8 | use str0m::{Event, Rtc, RtcError}; |
| 9 | |
| 10 | mod common; |
| 11 | use common::{Peer, TestRtc, init_log, negotiate, progress}; |
| 12 | use 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] |
| 16 | fn 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] |
| 130 | fn 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] |
| 230 | fn 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] |
| 320 | fn 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] |
| 403 | fn 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] |
| 460 | fn 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 | } |