File
Blob: firmware/vendor/str0m/tests/change-ssrc-reset-receive.rs
| 1 | use std::collections::VecDeque; |
| 2 | use std::time::Duration; |
| 3 | |
| 4 | use str0m::format::Codec; |
| 5 | use str0m::media::{Frequency, MediaKind, MediaTime}; |
| 6 | use str0m::rtp::{ExtensionValues, RtpWrite, Ssrc}; |
| 7 | use str0m::{Event, RtcError}; |
| 8 | use tracing::info; |
| 9 | |
| 10 | mod common; |
| 11 | use common::{connect_l_r, init_crypto_default, init_log, progress}; |
| 12 | |
| 13 | #[test] |
| 14 | pub fn change_ssrc_reset_receive() -> Result<(), RtcError> { |
| 15 | init_log(); |
| 16 | init_crypto_default(); |
| 17 | |
| 18 | let (mut l, mut r) = connect_l_r(); |
| 19 | |
| 20 | let mid = "aud".into(); |
| 21 | let rid = "hi".into(); |
| 22 | |
| 23 | // Initial SSRC |
| 24 | let ssrc_tx_initial: Ssrc = 42.into(); |
| 25 | // New SSRC to use after reset |
| 26 | let ssrc_tx_new: Ssrc = 84.into(); |
| 27 | |
| 28 | l.direct_api().declare_media(mid, MediaKind::Audio); |
| 29 | l.direct_api() |
| 30 | .declare_stream_tx(ssrc_tx_initial, None, mid, Some(rid)); |
| 31 | |
| 32 | r.direct_api() |
| 33 | .declare_media(mid, MediaKind::Audio) |
| 34 | .expect_rid_rx(rid); |
| 35 | |
| 36 | // Set initial timing |
| 37 | let max = l.last.max(r.last); |
| 38 | l.last = max; |
| 39 | r.last = max; |
| 40 | |
| 41 | let params = l.params_opus(); |
| 42 | let ssrc = l |
| 43 | .direct_api() |
| 44 | .stream_tx_by_mid(mid, Some(rid)) |
| 45 | .unwrap() |
| 46 | .ssrc(); |
| 47 | assert_eq!( |
| 48 | ssrc, ssrc_tx_initial, |
| 49 | "Initial SSRC should match what we set" |
| 50 | ); |
| 51 | assert_eq!(params.spec().codec, Codec::Opus); |
| 52 | let pt = params.pt(); |
| 53 | |
| 54 | // First batch of packets with initial SSRC |
| 55 | // We'll create a batch with several packets to cause timestamp rollover |
| 56 | let first_batch: Vec<&[u8]> = vec![ |
| 57 | &[0x1, 0x2, 0x3, 0x4], // packet at a timestamp close to rollover |
| 58 | &[0x5, 0x6, 0x7, 0x8], // packet at rollover point |
| 59 | &[0x9, 0xa, 0xb, 0xc], // packet after rollover |
| 60 | &[0xd, 0xe, 0xf, 0x0], // another packet after rollover |
| 61 | ]; |
| 62 | |
| 63 | let mut first_batch: VecDeque<_> = first_batch.into(); |
| 64 | |
| 65 | // Create a timestamp sequence that will cross the rollover boundary (close to 2^32) |
| 66 | // 32-bit timestamp max is 4,294,967,295 (0xFFFFFFFF) |
| 67 | // Start at a value very close to rollover |
| 68 | let rollover_threshold: u32 = 0xFFFFFFFF; |
| 69 | let timestamp_increment: u32 = 1000; // standard increment (e.g., for Opus) |
| 70 | |
| 71 | // Start our timeline near the rollover point |
| 72 | let start_timestamp: u32 = rollover_threshold - timestamp_increment * 2; |
| 73 | |
| 74 | // Counts for each packet - they'll increase, causing timestamp to roll over |
| 75 | let mut first_counts: Vec<u64> = vec![0, 1, 2, 3]; |
| 76 | let mut write_at = l.last + Duration::from_millis(300); |
| 77 | |
| 78 | // Send first batch with initial SSRC |
| 79 | loop { |
| 80 | if l.start + l.duration() > write_at { |
| 81 | write_at = l.last + Duration::from_millis(300); |
| 82 | if let Some(packet) = first_batch.pop_front() { |
| 83 | let wallclock = l.start + l.duration(); |
| 84 | let mut direct = l.direct_api(); |
| 85 | let stream = direct.stream_tx_by_mid(mid, Some(rid)).unwrap(); |
| 86 | |
| 87 | let count = first_counts.remove(0); |
| 88 | // Calculate timestamp to create rollover scenario |
| 89 | let time = start_timestamp.wrapping_add(count as u32 * timestamp_increment); |
| 90 | |
| 91 | info!( |
| 92 | "Sending packet {} with timestamp: {} on SSRC: {}", |
| 93 | count, time, ssrc_tx_initial |
| 94 | ); |
| 95 | let seq_no = (47_000 + count).into(); |
| 96 | |
| 97 | let exts = ExtensionValues { |
| 98 | audio_level: Some(-42 - count as i8), |
| 99 | voice_activity: Some(false), |
| 100 | ..Default::default() |
| 101 | }; |
| 102 | |
| 103 | stream.write_rtp(RtpWrite::new(pt, seq_no, time, wallclock, packet).ext_vals(exts)); |
| 104 | } |
| 105 | } |
| 106 | |
| 107 | progress(&mut l, &mut r)?; |
| 108 | |
| 109 | if first_batch.is_empty() && first_counts.is_empty() { |
| 110 | break; |
| 111 | } |
| 112 | } |
| 113 | |
| 114 | // Run a bit longer to ensure first batch of packets arrive |
| 115 | for _ in 0..20 { |
| 116 | progress(&mut l, &mut r)?; |
| 117 | } |
| 118 | |
| 119 | info!("First batch sent with SSRC: {}", ssrc_tx_initial); |
| 120 | |
| 121 | // Reset the SSRC with None for RTX since the stream doesn't use RTX |
| 122 | let mut api = l.direct_api(); |
| 123 | let result = api.reset_stream_tx(mid, Some(rid), ssrc_tx_new, None); |
| 124 | assert!(result.is_some(), "Reset should succeed with valid new SSRC"); |
| 125 | |
| 126 | // Verify the SSRC was changed |
| 127 | let updated_ssrc = l |
| 128 | .direct_api() |
| 129 | .stream_tx_by_mid(mid, Some(rid)) |
| 130 | .unwrap() |
| 131 | .ssrc(); |
| 132 | assert_eq!( |
| 133 | updated_ssrc, ssrc_tx_new, |
| 134 | "SSRC should be updated to new value" |
| 135 | ); |
| 136 | |
| 137 | // Wait to ensure pause detection occurs (default threshold is 1.5s) |
| 138 | let pause_duration = Duration::from_millis(2000); |
| 139 | let pause_end = l.last + pause_duration; |
| 140 | while l.last < pause_end { |
| 141 | progress(&mut l, &mut r)?; |
| 142 | } |
| 143 | |
| 144 | info!( |
| 145 | "Advanced time by {} ms to potentially trigger pause", |
| 146 | pause_duration.as_millis() |
| 147 | ); |
| 148 | |
| 149 | // Check if the stream was paused |
| 150 | let mut seen_paused = false; |
| 151 | for (_, event) in &r.events { |
| 152 | if let Event::StreamPaused(stream_paused) = event { |
| 153 | if stream_paused.paused { |
| 154 | seen_paused = true; |
| 155 | info!("Confirmed stream is paused"); |
| 156 | break; |
| 157 | } |
| 158 | } |
| 159 | } |
| 160 | |
| 161 | // Second batch of packets with new SSRC, but with timestamps that would appear to go |
| 162 | // backward if last_time is not reset during change_ssrc |
| 163 | let second_batch: Vec<&[u8]> = vec![&[0x9, 0xa, 0xb, 0xc], &[0xd, 0xe, 0xf, 0x10]]; |
| 164 | |
| 165 | let mut second_batch: VecDeque<_> = second_batch.into(); |
| 166 | let mut second_counts: Vec<u64> = vec![0, 1]; |
| 167 | write_at = l.last + Duration::from_millis(300); |
| 168 | |
| 169 | // Calculate a timestamp that would appear to go backward if last_time isn't reset |
| 170 | // The first batch used timestamps starting at 47_000_000 |
| 171 | // We'll use a timestamp that's lower, which would be interpreted as going backward |
| 172 | // if the internal state isn't reset |
| 173 | let backwards_base_timestamp: u32 = 1_000_000; |
| 174 | |
| 175 | // Send second batch with new SSRC and potentially problematic timestamps |
| 176 | loop { |
| 177 | if l.start + l.duration() > write_at { |
| 178 | write_at = l.last + Duration::from_millis(300); |
| 179 | if let Some(packet) = second_batch.pop_front() { |
| 180 | let wallclock = l.start + l.duration(); |
| 181 | let mut direct = l.direct_api(); |
| 182 | let stream = direct.stream_tx_by_mid(mid, Some(rid)).unwrap(); |
| 183 | |
| 184 | let count = second_counts.remove(0); |
| 185 | // Use a timestamp that would appear to go backward if last_time isn't reset |
| 186 | let time = (count as u32) * 1000 + backwards_base_timestamp; |
| 187 | let seq_no = (48_000 + count).into(); // Different seq range |
| 188 | |
| 189 | let exts = ExtensionValues { |
| 190 | audio_level: Some(-52 - count as i8), |
| 191 | voice_activity: Some(false), |
| 192 | ..Default::default() |
| 193 | }; |
| 194 | |
| 195 | info!( |
| 196 | "Sending packet with potentially backward timestamp: {} on SSRC: {}", |
| 197 | time, ssrc_tx_new |
| 198 | ); |
| 199 | |
| 200 | stream.write_rtp(RtpWrite::new(pt, seq_no, time, wallclock, packet).ext_vals(exts)); |
| 201 | } |
| 202 | } |
| 203 | |
| 204 | progress(&mut l, &mut r)?; |
| 205 | |
| 206 | if second_batch.is_empty() && second_counts.is_empty() { |
| 207 | // Run a bit longer to ensure packets arrive |
| 208 | for _ in 0..20 { |
| 209 | progress(&mut l, &mut r)?; |
| 210 | } |
| 211 | break; |
| 212 | } |
| 213 | } |
| 214 | |
| 215 | // Collect all received media packets and Media events |
| 216 | let media_packets: Vec<_> = r |
| 217 | .events |
| 218 | .iter() |
| 219 | .filter_map(|(_, e)| { |
| 220 | if let Event::RtpPacket(v) = e { |
| 221 | Some(v) |
| 222 | } else { |
| 223 | None |
| 224 | } |
| 225 | }) |
| 226 | .collect(); |
| 227 | |
| 228 | info!("Collected {} RTP packets", media_packets.len(),); |
| 229 | |
| 230 | // Should have received all 6 packets (4 from first batch + 2 from second batch) |
| 231 | assert_eq!(media_packets.len(), 6, "Should have received all 6 packets"); |
| 232 | |
| 233 | // Verify that we have packets with both SSRCs |
| 234 | let first_ssrc_packets: Vec<_> = media_packets |
| 235 | .iter() |
| 236 | .filter(|p| p.header.ssrc == ssrc_tx_initial) |
| 237 | .collect(); |
| 238 | let second_ssrc_packets: Vec<_> = media_packets |
| 239 | .iter() |
| 240 | .filter(|p| p.header.ssrc == ssrc_tx_new) |
| 241 | .collect(); |
| 242 | |
| 243 | assert_eq!( |
| 244 | first_ssrc_packets.len(), |
| 245 | 4, |
| 246 | "Should have 4 packets with initial SSRC" |
| 247 | ); |
| 248 | assert_eq!( |
| 249 | second_ssrc_packets.len(), |
| 250 | 2, |
| 251 | "Should have 2 packets with new SSRC" |
| 252 | ); |
| 253 | |
| 254 | // Verify the stream was unpaused after receiving the second batch with new SSRC |
| 255 | let mut stream_unpaused_after_new_ssrc = false; |
| 256 | for (_, event) in &r.events { |
| 257 | if let Event::StreamPaused(stream_paused) = event { |
| 258 | if !stream_paused.paused && stream_paused.ssrc == ssrc_tx_new { |
| 259 | stream_unpaused_after_new_ssrc = true; |
| 260 | info!("Confirmed stream unpaused after new SSRC packets"); |
| 261 | break; |
| 262 | } |
| 263 | } |
| 264 | } |
| 265 | |
| 266 | // If we saw a pause event, we should also see an unpause event after the new packets |
| 267 | if seen_paused { |
| 268 | assert!( |
| 269 | stream_unpaused_after_new_ssrc, |
| 270 | "Stream should have unpaused after receiving packets with new SSRC" |
| 271 | ); |
| 272 | } |
| 273 | |
| 274 | // Verify the first batch packets sequence numbers |
| 275 | for (i, packet) in first_ssrc_packets.iter().enumerate() { |
| 276 | assert_eq!( |
| 277 | packet.header.sequence_number, |
| 278 | 47000 + i as u16, |
| 279 | "First batch packet {} should have correct sequence number", |
| 280 | i |
| 281 | ); |
| 282 | } |
| 283 | |
| 284 | // Verify first batch timestamps - they should cross the rollover boundary |
| 285 | // First two packets should be before rollover, second two after rollover |
| 286 | // Calculating expected timestamps for verification |
| 287 | let expected_timestamps: Vec<u32> = (0..4) |
| 288 | .map(|i| start_timestamp.wrapping_add(i as u32 * timestamp_increment)) |
| 289 | .collect(); |
| 290 | |
| 291 | // Verify the rollover pattern |
| 292 | // Based on the logs, we can see that packets 1, 2, and 3 have timestamps close to max, |
| 293 | // and only the 4th packet rolls over |
| 294 | info!("Checking timestamp pattern for packets crossing rollover"); |
| 295 | for (i, ts) in expected_timestamps.iter().enumerate() { |
| 296 | info!("Packet {} timestamp: {}", i, ts); |
| 297 | } |
| 298 | |
| 299 | // Verify the actual timestamps match our expected values |
| 300 | for (i, packet) in first_ssrc_packets.iter().enumerate() { |
| 301 | assert_eq!( |
| 302 | packet.header.timestamp, expected_timestamps[i], |
| 303 | "First batch packet {} should have expected timestamp after rollover", |
| 304 | i |
| 305 | ); |
| 306 | } |
| 307 | |
| 308 | // Verify second batch (new SSRC) sequence numbers |
| 309 | for (i, packet) in second_ssrc_packets.iter().enumerate() { |
| 310 | assert_eq!( |
| 311 | packet.header.sequence_number, |
| 312 | 48000 + i as u16, |
| 313 | "Second batch packet {} should have correct sequence number", |
| 314 | i |
| 315 | ); |
| 316 | } |
| 317 | |
| 318 | // The timestamps for the second batch are much lower, but should be interpreted correctly |
| 319 | // because the internal state was reset by change_ssrc() |
| 320 | assert_eq!( |
| 321 | second_ssrc_packets[0].header.timestamp, backwards_base_timestamp, |
| 322 | "First packet with new SSRC should have expected timestamp" |
| 323 | ); |
| 324 | assert_eq!( |
| 325 | second_ssrc_packets[1].header.timestamp, |
| 326 | backwards_base_timestamp + 1000, |
| 327 | "Second packet with new SSRC should have expected timestamp" |
| 328 | ); |
| 329 | |
| 330 | let mut api = r.direct_api(); |
| 331 | let rx = api.stream_rx_by_mid(mid, Some(rid)).unwrap(); |
| 332 | |
| 333 | // If we don't reset the state properly on SSRC change, this will be |
| 334 | // a crazy value. |
| 335 | assert_eq!( |
| 336 | rx.last_time(), |
| 337 | Some(MediaTime::new(1001000, Frequency::FORTY_EIGHT_KHZ)) |
| 338 | ); |
| 339 | |
| 340 | // This test verifies two key behaviors: |
| 341 | // 1. Timestamp rollover is handled correctly within a single SSRC |
| 342 | // 2. After SSRC change, timestamp history is reset so a low timestamp |
| 343 | // doesn't appear to go backward from the high timestamp of the previous SSRC |
| 344 | |
| 345 | Ok(()) |
| 346 | } |