File
Blob: firmware/vendor/str0m/tests/mediatime-backwards.rs
| 1 | use std::net::Ipv4Addr; |
| 2 | use std::time::{Duration, Instant}; |
| 3 | |
| 4 | use str0m::format::Codec; |
| 5 | use str0m::media::{Direction, MediaKind, MediaTime}; |
| 6 | use str0m::rtp::RawPacket; |
| 7 | use str0m::{Event, RtcError}; |
| 8 | use tracing::info; |
| 9 | |
| 10 | mod common; |
| 11 | use common::{Peer, TestRtc, init_crypto_default, init_log, progress}; |
| 12 | |
| 13 | #[test] |
| 14 | pub fn mediatime_backwards() -> Result<(), RtcError> { |
| 15 | init_log(); |
| 16 | init_crypto_default(); |
| 17 | |
| 18 | let mut l = TestRtc::new(Peer::Left); |
| 19 | let mut r = TestRtc::new(Peer::Right); |
| 20 | |
| 21 | // Enable raw packets to trace the RTP headers |
| 22 | r.rtc = str0m::Rtc::builder() |
| 23 | .enable_raw_packets(true) |
| 24 | .build(Instant::now()); |
| 25 | |
| 26 | l.add_host_candidate((Ipv4Addr::new(1, 1, 1, 1), 1000).into()); |
| 27 | r.add_host_candidate((Ipv4Addr::new(2, 2, 2, 2), 2000).into()); |
| 28 | |
| 29 | // The change is on the L (sending side) with Direction::SendRecv. |
| 30 | let mut change = l.sdp_api(); |
| 31 | let mid = change.add_media(MediaKind::Audio, Direction::SendRecv, None, None, None); |
| 32 | let (offer, pending) = change.apply().unwrap(); |
| 33 | |
| 34 | let answer = r.rtc.sdp_api().accept_offer(offer)?; |
| 35 | l.rtc.sdp_api().accept_answer(pending, answer)?; |
| 36 | |
| 37 | loop { |
| 38 | if l.is_connected() || r.is_connected() { |
| 39 | break; |
| 40 | } |
| 41 | progress(&mut l, &mut r)?; |
| 42 | } |
| 43 | |
| 44 | let max = l.last.max(r.last); |
| 45 | l.last = max; |
| 46 | r.last = max; |
| 47 | |
| 48 | let params = l.params_opus(); |
| 49 | assert_eq!(params.spec().codec, Codec::Opus); |
| 50 | let pt = params.pt(); |
| 51 | |
| 52 | let data_a = vec![1_u8; 80]; |
| 53 | |
| 54 | // First, send some normal packets with timestamps < 2^31 |
| 55 | let base_timestamp = 1_000_000_u32; // Well below 2^31 |
| 56 | |
| 57 | // Send 10 normal packets |
| 58 | for i in 0..10 { |
| 59 | let wallclock = l.start + l.duration(); |
| 60 | let timestamp = base_timestamp + i * 960; // Standard Opus timestamp increment |
| 61 | let time = MediaTime::from_90khz(timestamp as u64); |
| 62 | |
| 63 | l.rtc |
| 64 | .writer(mid) |
| 65 | .unwrap() |
| 66 | .write(pt, wallclock, time, data_a.clone())?; |
| 67 | |
| 68 | progress(&mut l, &mut r)?; |
| 69 | } |
| 70 | |
| 71 | info!("Sent 10 normal packets with regular timestamps"); |
| 72 | |
| 73 | // Store the last RTP timestamp we observed before the pause |
| 74 | let mut last_ts_before_pause: Option<u32> = None; |
| 75 | |
| 76 | // Get the last timestamp from received packets |
| 77 | for (_, event) in &r.events { |
| 78 | if let Event::RawPacket(raw_packet) = event { |
| 79 | if let RawPacket::RtpRx(header, _) = &**raw_packet { |
| 80 | last_ts_before_pause = Some(header.timestamp); |
| 81 | } |
| 82 | } |
| 83 | } |
| 84 | |
| 85 | info!("Last timestamp before pause: {:?}", last_ts_before_pause); |
| 86 | |
| 87 | // Now advance time by more than 1.5 seconds (default pause threshold) without sending any packets |
| 88 | // This will trigger the pause condition in the receiver |
| 89 | let pause_duration = Duration::from_millis(2000); |
| 90 | |
| 91 | // Advance time using timeouts instead of sleep |
| 92 | let pause_end = l.last + pause_duration; |
| 93 | while l.last < pause_end { |
| 94 | // Just handle timeouts without sending any packets |
| 95 | progress(&mut l, &mut r)?; |
| 96 | } |
| 97 | |
| 98 | info!( |
| 99 | "Advanced time by {} ms to trigger pause", |
| 100 | pause_duration.as_millis() |
| 101 | ); |
| 102 | |
| 103 | // Wait to ensure we get a paused event |
| 104 | let mut seen_paused = false; |
| 105 | let pause_check_start = l.last; |
| 106 | let pause_check_timeout = Duration::from_millis(500); |
| 107 | |
| 108 | while !seen_paused && l.last < pause_check_start + pause_check_timeout { |
| 109 | progress(&mut l, &mut r)?; |
| 110 | |
| 111 | // Check for pause event |
| 112 | for (_, event) in &r.events { |
| 113 | if let Event::StreamPaused(stream_paused) = event { |
| 114 | if stream_paused.paused { |
| 115 | seen_paused = true; |
| 116 | info!("Confirmed stream is paused"); |
| 117 | break; |
| 118 | } |
| 119 | } |
| 120 | } |
| 121 | } |
| 122 | |
| 123 | assert!( |
| 124 | seen_paused, |
| 125 | "Stream did not enter paused state within timeout" |
| 126 | ); |
| 127 | |
| 128 | // Now send a packet with a timestamp that will trigger the bug |
| 129 | // We need timestamp > 2^31 + previous_timestamp to trigger the backwards condition |
| 130 | // Using the threshold from header.rs, HALF = 1 << 31 |
| 131 | let problematic_timestamp = (1u32 << 31) + base_timestamp + 10 * 960; |
| 132 | |
| 133 | info!( |
| 134 | "Sending packet with problematic timestamp: {}", |
| 135 | problematic_timestamp |
| 136 | ); |
| 137 | |
| 138 | let wallclock = l.start + l.duration(); |
| 139 | let time = MediaTime::from_90khz(problematic_timestamp as u64); |
| 140 | |
| 141 | // Save the timestamp index for comparison BEFORE sending |
| 142 | let events_before_problematic = r.events.len(); |
| 143 | |
| 144 | l.rtc |
| 145 | .writer(mid) |
| 146 | .unwrap() |
| 147 | .write(pt, wallclock, time, data_a.clone())?; |
| 148 | |
| 149 | // Process the packet and a few more cycles to ensure it's received |
| 150 | for _ in 0..10 { |
| 151 | progress(&mut l, &mut r)?; |
| 152 | } |
| 153 | |
| 154 | // We're looking for evidence that: |
| 155 | // 1. The stream unpaused (proving packet was received) |
| 156 | // 2. The timestamp value in the received packet |
| 157 | let mut problematic_ts: Option<u32> = None; |
| 158 | let mut stream_unpaused_after_problematic = false; |
| 159 | |
| 160 | // Only look at events that occurred after sending our problematic packet |
| 161 | for (_, (_, event)) in r.events.iter().enumerate().skip(events_before_problematic) { |
| 162 | match event { |
| 163 | Event::RawPacket(raw_packet) => { |
| 164 | if let RawPacket::RtpRx(header, _) = &**raw_packet { |
| 165 | info!("After pause - RTP timestamp: {}", header.timestamp); |
| 166 | problematic_ts = Some(header.timestamp); |
| 167 | } |
| 168 | } |
| 169 | Event::StreamPaused(stream_paused) => { |
| 170 | if !stream_paused.paused { |
| 171 | stream_unpaused_after_problematic = true; |
| 172 | info!("Stream unpaused after receiving problematic packet"); |
| 173 | } |
| 174 | } |
| 175 | _ => {} |
| 176 | } |
| 177 | } |
| 178 | |
| 179 | // The test succeeds if: |
| 180 | // 1. We confirmed the stream was paused |
| 181 | // 2. We received the packet with the problematic timestamp |
| 182 | // 3. The stream unpaused after receiving the problematic packet |
| 183 | // 4. The timestamps continue to move forward after the pause |
| 184 | |
| 185 | assert!(seen_paused, "Stream was never paused"); |
| 186 | assert!( |
| 187 | problematic_ts.is_some(), |
| 188 | "Did not receive packet with problematic timestamp" |
| 189 | ); |
| 190 | assert!( |
| 191 | stream_unpaused_after_problematic, |
| 192 | "Stream did not unpause after receiving problematic packet" |
| 193 | ); |
| 194 | |
| 195 | // Verify that time moves forward even after a pause with a problematic timestamp |
| 196 | // This is the core of the bugfix - ensuring that MediaTime never goes backwards |
| 197 | if let (Some(prev_ts), Some(new_ts)) = (last_ts_before_pause, problematic_ts) { |
| 198 | assert!( |
| 199 | new_ts > prev_ts, |
| 200 | "Time went backwards! Previous timestamp: {}, Current timestamp: {}", |
| 201 | prev_ts, |
| 202 | new_ts |
| 203 | ); |
| 204 | } |
| 205 | |
| 206 | Ok(()) |
| 207 | } |