Skip to content
File

Blob: firmware/vendor/str0m/tests/stream-lifecycle.rs

rust377 lines
1//! Tests for stream pause/resume and lifecycle events.
2 
3use std::net::Ipv4Addr;
4use std::time::Duration;
5 
6use str0m::format::Codec;
7use str0m::media::{Direction, MediaKind};
8use str0m::{Event, RtcError};
9 
10mod common;
11use common::{Peer, TestRtc, init_crypto_default, init_log, negotiate, progress};
12 
13/// Test that StreamPaused event is emitted after no packets for ~1.5 seconds.
14#[test]
15fn stream_pause_detection_timeout() -> Result<(), RtcError> {
16 init_log();
17 init_crypto_default();
18 
19 let mut l = TestRtc::new(Peer::Left);
20 let mut r = TestRtc::new(Peer::Right);
21 
22 l.add_host_candidate((Ipv4Addr::new(1, 1, 1, 1), 1000).into());
23 r.add_host_candidate((Ipv4Addr::new(2, 2, 2, 2), 2000).into());
24 
25 let mid = negotiate(&mut l, &mut r, |change| {
26 change.add_media(MediaKind::Audio, Direction::SendRecv, None, None, None)
27 });
28 
29 loop {
30 if l.is_connected() && r.is_connected() {
31 break;
32 }
33 progress(&mut l, &mut r)?;
34 }
35 
36 let max = l.last.max(r.last);
37 l.last = max;
38 r.last = max;
39 
40 let params = l.params_opus();
41 assert_eq!(params.spec().codec, Codec::Opus);
42 let pt = params.pt();
43 let data = vec![1_u8; 80];
44 
45 // Send packets for 500ms
46 let send_until = l.duration() + Duration::from_millis(500);
47 loop {
48 if l.duration() >= send_until {
49 break;
50 }
51 let wallclock = l.start + l.duration();
52 let time = l.duration().into();
53 l.writer(mid)
54 .unwrap()
55 .write(pt, wallclock, time, data.clone())?;
56 progress(&mut l, &mut r)?;
57 }
58 
59 // Stop sending and wait for pause detection (>1.5 seconds)
60 let pause_wait = l.duration() + Duration::from_secs(3);
61 loop {
62 if l.duration() >= pause_wait {
63 break;
64 }
65 progress(&mut l, &mut r)?;
66 }
67 
68 // Check for StreamPaused event on receiver
69 let paused_events: Vec<_> = r
70 .events
71 .iter()
72 .filter_map(|(_, e)| {
73 if let Event::StreamPaused(p) = e {
74 Some(p)
75 } else {
76 None
77 }
78 })
79 .collect();
80 
81 assert!(
82 paused_events.iter().any(|p| p.paused && p.mid == mid),
83 "Expected StreamPaused event with paused=true for mid {:?}",
84 mid
85 );
86 
87 Ok(())
88}
89 
90/// Test pause detection followed by resume when packets arrive again.
91#[test]
92fn stream_pause_resume_cycle() -> Result<(), RtcError> {
93 init_log();
94 init_crypto_default();
95 
96 let mut l = TestRtc::new(Peer::Left);
97 let mut r = TestRtc::new(Peer::Right);
98 
99 l.add_host_candidate((Ipv4Addr::new(1, 1, 1, 1), 1000).into());
100 r.add_host_candidate((Ipv4Addr::new(2, 2, 2, 2), 2000).into());
101 
102 let mid = negotiate(&mut l, &mut r, |change| {
103 change.add_media(MediaKind::Audio, Direction::SendRecv, None, None, None)
104 });
105 
106 loop {
107 if l.is_connected() && r.is_connected() {
108 break;
109 }
110 progress(&mut l, &mut r)?;
111 }
112 
113 let max = l.last.max(r.last);
114 l.last = max;
115 r.last = max;
116 
117 let params = l.params_opus();
118 let pt = params.pt();
119 let data = vec![1_u8; 80];
120 
121 // Phase 1: Send packets
122 let send_until = l.duration() + Duration::from_millis(500);
123 loop {
124 if l.duration() >= send_until {
125 break;
126 }
127 let wallclock = l.start + l.duration();
128 let time = l.duration().into();
129 l.writer(mid)
130 .unwrap()
131 .write(pt, wallclock, time, data.clone())?;
132 progress(&mut l, &mut r)?;
133 }
134 
135 // Phase 2: Stop sending and wait for pause
136 let pause_wait = l.duration() + Duration::from_secs(2);
137 loop {
138 if l.duration() >= pause_wait {
139 break;
140 }
141 progress(&mut l, &mut r)?;
142 }
143 
144 // Phase 3: Resume sending
145 let resume_until = l.duration() + Duration::from_millis(500);
146 loop {
147 if l.duration() >= resume_until {
148 break;
149 }
150 let wallclock = l.start + l.duration();
151 let time = l.duration().into();
152 l.writer(mid)
153 .unwrap()
154 .write(pt, wallclock, time, data.clone())?;
155 progress(&mut l, &mut r)?;
156 }
157 
158 // Check for both pause and resume events
159 let paused_states: Vec<_> = r
160 .events
161 .iter()
162 .filter_map(|(_, e)| {
163 if let Event::StreamPaused(p) = e {
164 if p.mid == mid { Some(p.paused) } else { None }
165 } else {
166 None
167 }
168 })
169 .collect();
170 
171 assert!(
172 paused_states.contains(&true),
173 "Expected StreamPaused event with paused=true"
174 );
175 assert!(
176 paused_states.contains(&false),
177 "Expected StreamPaused event with paused=false (resume)"
178 );
179 
180 Ok(())
181}
182 
183/// Test changing media direction to SendOnly.
184#[test]
185fn stream_direction_change_sendonly() -> Result<(), RtcError> {
186 init_log();
187 init_crypto_default();
188 
189 let mut l = TestRtc::new(Peer::Left);
190 let mut r = TestRtc::new(Peer::Right);
191 
192 l.add_host_candidate((Ipv4Addr::new(1, 1, 1, 1), 1000).into());
193 r.add_host_candidate((Ipv4Addr::new(2, 2, 2, 2), 2000).into());
194 
195 // Initial negotiation with SendRecv
196 let mid = negotiate(&mut l, &mut r, |change| {
197 change.add_media(MediaKind::Audio, Direction::SendRecv, None, None, None)
198 });
199 
200 loop {
201 if l.is_connected() && r.is_connected() {
202 break;
203 }
204 progress(&mut l, &mut r)?;
205 }
206 
207 assert_eq!(l.media(mid).unwrap().direction(), Direction::SendRecv);
208 assert_eq!(r.media(mid).unwrap().direction(), Direction::SendRecv);
209 
210 // Change L to SendOnly
211 negotiate(&mut l, &mut r, |change| {
212 change.set_direction(mid, Direction::SendOnly);
213 });
214 
215 assert_eq!(
216 l.media(mid).unwrap().direction(),
217 Direction::SendOnly,
218 "L should be SendOnly"
219 );
220 assert_eq!(
221 r.media(mid).unwrap().direction(),
222 Direction::RecvOnly,
223 "R should be RecvOnly (opposite of L's SendOnly)"
224 );
225 
226 Ok(())
227}
228 
229/// Test changing media direction to RecvOnly.
230#[test]
231fn stream_direction_change_recvonly() -> Result<(), RtcError> {
232 init_log();
233 init_crypto_default();
234 
235 let mut l = TestRtc::new(Peer::Left);
236 let mut r = TestRtc::new(Peer::Right);
237 
238 l.add_host_candidate((Ipv4Addr::new(1, 1, 1, 1), 1000).into());
239 r.add_host_candidate((Ipv4Addr::new(2, 2, 2, 2), 2000).into());
240 
241 let mid = negotiate(&mut l, &mut r, |change| {
242 change.add_media(MediaKind::Audio, Direction::SendRecv, None, None, None)
243 });
244 
245 loop {
246 if l.is_connected() && r.is_connected() {
247 break;
248 }
249 progress(&mut l, &mut r)?;
250 }
251 
252 assert_eq!(l.media(mid).unwrap().direction(), Direction::SendRecv);
253 assert_eq!(r.media(mid).unwrap().direction(), Direction::SendRecv);
254 
255 // Change L to RecvOnly
256 negotiate(&mut l, &mut r, |change| {
257 change.set_direction(mid, Direction::RecvOnly);
258 });
259 
260 assert_eq!(
261 l.media(mid).unwrap().direction(),
262 Direction::RecvOnly,
263 "L should be RecvOnly"
264 );
265 assert_eq!(
266 r.media(mid).unwrap().direction(),
267 Direction::SendOnly,
268 "R should be SendOnly (opposite of L's RecvOnly)"
269 );
270 
271 Ok(())
272}
273 
274/// Test changing media direction to Inactive.
275#[test]
276fn stream_direction_change_inactive() -> Result<(), RtcError> {
277 init_log();
278 init_crypto_default();
279 
280 let mut l = TestRtc::new(Peer::Left);
281 let mut r = TestRtc::new(Peer::Right);
282 
283 l.add_host_candidate((Ipv4Addr::new(1, 1, 1, 1), 1000).into());
284 r.add_host_candidate((Ipv4Addr::new(2, 2, 2, 2), 2000).into());
285 
286 let mid = negotiate(&mut l, &mut r, |change| {
287 change.add_media(MediaKind::Audio, Direction::SendRecv, None, None, None)
288 });
289 
290 loop {
291 if l.is_connected() && r.is_connected() {
292 break;
293 }
294 progress(&mut l, &mut r)?;
295 }
296 
297 assert_eq!(l.media(mid).unwrap().direction(), Direction::SendRecv);
298 assert_eq!(r.media(mid).unwrap().direction(), Direction::SendRecv);
299 
300 // Change L to Inactive
301 negotiate(&mut l, &mut r, |change| {
302 change.set_direction(mid, Direction::Inactive);
303 });
304 
305 assert_eq!(
306 l.media(mid).unwrap().direction(),
307 Direction::Inactive,
308 "L should be Inactive"
309 );
310 assert_eq!(
311 r.media(mid).unwrap().direction(),
312 Direction::Inactive,
313 "R should be Inactive (both sides inactive)"
314 );
315 
316 Ok(())
317}
318 
319/// Test MediaChanged event is generated on direction change.
320#[test]
321fn stream_media_changed_event() -> Result<(), RtcError> {
322 init_log();
323 init_crypto_default();
324 
325 let mut l = TestRtc::new(Peer::Left);
326 let mut r = TestRtc::new(Peer::Right);
327 
328 l.add_host_candidate((Ipv4Addr::new(1, 1, 1, 1), 1000).into());
329 r.add_host_candidate((Ipv4Addr::new(2, 2, 2, 2), 2000).into());
330 
331 let mid = negotiate(&mut l, &mut r, |change| {
332 change.add_media(MediaKind::Audio, Direction::SendRecv, None, None, None)
333 });
334 
335 loop {
336 if l.is_connected() && r.is_connected() {
337 break;
338 }
339 progress(&mut l, &mut r)?;
340 }
341 
342 // Clear previous events
343 l.events.clear();
344 r.events.clear();
345 
346 // Change direction
347 negotiate(&mut l, &mut r, |change| {
348 change.set_direction(mid, Direction::SendOnly);
349 });
350 
351 // Progress to process events
352 for _ in 0..20 {
353 progress(&mut l, &mut r)?;
354 }
355 
356 // Check for MediaChanged event
357 let changed_events: Vec<_> = r
358 .events
359 .iter()
360 .filter_map(|(_, e)| {
361 if let Event::MediaChanged(c) = e {
362 Some(c)
363 } else {
364 None
365 }
366 })
367 .collect();
368 
369 assert!(
370 changed_events.iter().any(|c| c.mid == mid),
371 "Expected MediaChanged event for mid {:?}",
372 mid
373 );
374 
375 Ok(())
376}