Skip to content
File

Blob: firmware/vendor/str0m/tests/change-ssrc-reset-receive.rs

rust347 lines
1use std::collections::VecDeque;
2use std::time::Duration;
3 
4use str0m::format::Codec;
5use str0m::media::{Frequency, MediaKind, MediaTime};
6use str0m::rtp::{ExtensionValues, RtpWrite, Ssrc};
7use str0m::{Event, RtcError};
8use tracing::info;
9 
10mod common;
11use common::{connect_l_r, init_crypto_default, init_log, progress};
12 
13#[test]
14pub 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}