Skip to content
File

Blob: firmware/vendor/str0m/tests/data-channel.rs

rust661 lines
1use std::net::Ipv4Addr;
2use std::time::{Duration, Instant};
3 
4use netem::NetemConfig;
5use str0m::channel::ChannelConfig;
6use str0m::{Event, Input, Output, RtcError};
7 
8mod common;
9use common::{Peer, TestRtc, connect_l_r, init_crypto_default, init_log, progress};
10 
11/// Poll one peer while deliberately withholding all of its network output.
12///
13/// This models an application that keeps driving its local RTC while packets
14/// to the peer are delayed. Public events are returned to the test; transmits
15/// are dropped instead of being handed to the remote `TestRtc`.
16fn poll_without_delivering_network(
17 rtc: &mut TestRtc,
18 now: Instant,
19) -> Result<Vec<Event>, RtcError> {
20 rtc.rtc.handle_input(Input::Timeout(now))?;
21 
22 let mut events = vec![];
23 loop {
24 match rtc.rtc.poll_output()? {
25 Output::Event(event) => events.push(event),
26 Output::Transmit(_) => {}
27 Output::Timeout(_) => return Ok(events),
28 }
29 }
30}
31 
32#[test]
33pub fn data_channel() -> Result<(), RtcError> {
34 init_log();
35 init_crypto_default();
36 
37 let mut l = TestRtc::new(Peer::Left);
38 let mut r = TestRtc::new(Peer::Right);
39 
40 l.add_host_candidate((Ipv4Addr::new(1, 1, 1, 1), 1000).into());
41 r.add_host_candidate((Ipv4Addr::new(2, 2, 2, 2), 2000).into());
42 
43 let mut change = l.sdp_api();
44 let cid = change.add_channel("My little channel".into());
45 change.add_channel("My little channel 2".into());
46 let (offer, pending) = change.apply().unwrap();
47 
48 let answer = r.rtc.sdp_api().accept_offer(offer)?;
49 l.rtc.sdp_api().accept_answer(pending, answer)?;
50 
51 loop {
52 if l.is_connected() || r.is_connected() {
53 break;
54 }
55 progress(&mut l, &mut r)?;
56 }
57 
58 let max = l.last.max(r.last);
59 l.last = max;
60 r.last = max;
61 
62 loop {
63 if let Some(mut chan) = l.channel(cid) {
64 chan.write(false, "Hello world! ".as_bytes())
65 .expect("to write string");
66 }
67 
68 progress(&mut l, &mut r)?;
69 
70 if l.duration() > Duration::from_secs(10) {
71 break;
72 }
73 }
74 
75 assert!(r.events.len() > 120);
76 
77 Ok(())
78}
79 
80/// Closing a data channel must propagate to the remote peer (via the SCTP
81/// stream reset handshake), and once the handshake completes the freed stream
82/// id must be reusable by a new in-band channel.
83#[test]
84pub fn data_channel_close_reopen() -> Result<(), RtcError> {
85 init_log();
86 init_crypto_default();
87 
88 let mut l = TestRtc::new(Peer::Left);
89 let mut r = TestRtc::new(Peer::Right);
90 
91 l.add_host_candidate((Ipv4Addr::new(1, 1, 1, 1), 1000).into());
92 r.add_host_candidate((Ipv4Addr::new(2, 2, 2, 2), 2000).into());
93 
94 let mut change = l.sdp_api();
95 let cid = change.add_channel("churn".into());
96 let (offer, pending) = change.apply().unwrap();
97 
98 let answer = r.rtc.sdp_api().accept_offer(offer)?;
99 l.rtc.sdp_api().accept_answer(pending, answer)?;
100 
101 loop {
102 if l.is_connected() || r.is_connected() {
103 break;
104 }
105 progress(&mut l, &mut r)?;
106 }
107 
108 let max = l.last.max(r.last);
109 l.last = max;
110 r.last = max;
111 
112 // Wait until both sides see the channel open.
113 loop {
114 progress(&mut l, &mut r)?;
115 
116 let l_open = l
117 .events
118 .iter()
119 .any(|(_, e)| matches!(e, Event::ChannelOpen(id, _) if *id == cid));
120 let r_open = r
121 .events
122 .iter()
123 .any(|(_, e)| matches!(e, Event::ChannelOpen(_, _)));
124 
125 if l_open && r_open {
126 break;
127 }
128 assert!(
129 l.duration() < Duration::from_secs(10),
130 "first channel should open on both sides"
131 );
132 }
133 
134 let stream_id = l
135 .direct_api()
136 .sctp_stream_id_by_channel_id(cid)
137 .expect("stream id for open channel");
138 
139 // Close locally. The reset handshake must inform the remote, which
140 // previously never received ChannelClose.
141 l.direct_api().close_data_channel(cid);
142 
143 loop {
144 progress(&mut l, &mut r)?;
145 
146 let l_closed = l
147 .events
148 .iter()
149 .any(|(_, e)| matches!(e, Event::ChannelClose(id) if *id == cid));
150 let r_closed = r
151 .events
152 .iter()
153 .any(|(_, e)| matches!(e, Event::ChannelClose(_)));
154 
155 if l_closed && r_closed {
156 break;
157 }
158 assert!(
159 l.duration() < Duration::from_secs(20),
160 "both sides should see ChannelClose"
161 );
162 }
163 
164 // The reset handshake finishes with round-trips (reciprocal reset and
165 // RECONFIG-RESPONSEs) that carry no public events, so there is nothing to
166 // wait on here. Creating the next channel immediately is fine: its stream
167 // id allocation happens on a later timeout, by which time the handshake
168 // rounds have been ferried through. The stream id assertion below fails
169 // loudly if the allocator did not release the id in time.
170 let cid2 = l.direct_api().create_data_channel(ChannelConfig {
171 label: "churn2".into(),
172 ..Default::default()
173 });
174 assert_ne!(cid, cid2);
175 
176 loop {
177 progress(&mut l, &mut r)?;
178 
179 let l_open = l.events.iter().any(
180 |(_, e)| matches!(e, Event::ChannelOpen(id, label) if *id == cid2 && label == "churn2"),
181 );
182 let r_open = r
183 .events
184 .iter()
185 .any(|(_, e)| matches!(e, Event::ChannelOpen(_, label) if label == "churn2"));
186 
187 if l_open && r_open {
188 break;
189 }
190 assert!(
191 l.duration() < Duration::from_secs(30),
192 "reopened channel should open on both sides"
193 );
194 }
195 
196 // The freed stream id must have been reused, proving the allocator
197 // released it when the reset handshake completed.
198 assert_eq!(
199 l.direct_api().sctp_stream_id_by_channel_id(cid2),
200 Some(stream_id),
201 "reopened channel should reuse the freed stream id"
202 );
203 
204 // Data flows on the reopened channel.
205 loop {
206 if let Some(mut chan) = l.channel(cid2) {
207 chan.write(false, b"hello again").expect("write to succeed");
208 }
209 
210 progress(&mut l, &mut r)?;
211 
212 let got_data = r.events.iter().any(
213 |(_, e)| matches!(e, Event::ChannelData(d) if d.data.as_slice() == b"hello again"),
214 );
215 if got_data {
216 break;
217 }
218 assert!(
219 l.duration() < Duration::from_secs(40),
220 "data should flow on the reopened channel"
221 );
222 }
223 
224 Ok(())
225}
226 
227#[test]
228pub fn negotiated_reuse_waits_for_reset_before_channel_open() -> Result<(), RtcError> {
229 init_log();
230 init_crypto_default();
231 
232 let (mut l, mut r) = connect_l_r();
233 let stream_id = 10;
234 let config = ChannelConfig {
235 label: "old-generation".into(),
236 negotiated: Some(stream_id),
237 ..Default::default()
238 };
239 let old_l = l.direct_api().create_data_channel(config.clone());
240 let _old_r = r.direct_api().create_data_channel(config);
241 
242 loop {
243 progress(&mut l, &mut r)?;
244 let left_open = l
245 .events
246 .iter()
247 .any(|(_, event)| matches!(event, Event::ChannelOpen(id, _) if *id == old_l));
248 let right_open = r
249 .events
250 .iter()
251 .any(|(_, event)| matches!(event, Event::ChannelOpen(_, _)));
252 if left_open && right_open {
253 break;
254 }
255 assert!(
256 l.duration() < Duration::from_secs(10),
257 "negotiated channel should open on both peers"
258 );
259 }
260 
261 l.events.clear();
262 l.direct_api().close_data_channel(old_l);
263 
264 // Process the local close, but withhold its RE-CONFIG packets. The old
265 // SCTP stream therefore still exists and its reset cannot be complete.
266 let close_at = l.last;
267 let close_events = poll_without_delivering_network(&mut l, close_at)?;
268 assert!(
269 close_events
270 .iter()
271 .any(|event| matches!(event, Event::ChannelClose(id) if *id == old_l)),
272 "the old local channel should close"
273 );
274 
275 let replacement = l.direct_api().create_data_channel(ChannelConfig {
276 label: "replacement".into(),
277 negotiated: Some(stream_id),
278 ..Default::default()
279 });
280 let replacement_events =
281 poll_without_delivering_network(&mut l, close_at + Duration::from_millis(1))?;
282 
283 assert!(
284 !replacement_events
285 .iter()
286 .any(|event| matches!(event, Event::ChannelOpen(id, _) if *id == replacement)),
287 "a negotiated replacement must not open before the old reset completes"
288 );
289 
290 Ok(())
291}
292 
293#[test]
294pub fn unconfirmed_reset_keeps_stream_id_reserved() -> Result<(), RtcError> {
295 init_log();
296 init_crypto_default();
297 
298 let (mut l, mut r) = connect_l_r();
299 let old = l.direct_api().create_data_channel(ChannelConfig::default());
300 
301 loop {
302 progress(&mut l, &mut r)?;
303 if l.events
304 .iter()
305 .any(|(_, event)| matches!(event, Event::ChannelOpen(id, _) if *id == old))
306 {
307 break;
308 }
309 assert!(
310 l.duration() < Duration::from_secs(10),
311 "first channel should open"
312 );
313 }
314 
315 let old_stream_id = l
316 .direct_api()
317 .sctp_stream_id_by_channel_id(old)
318 .expect("old channel should have a stream ID");
319 l.direct_api().close_data_channel(old);
320 
321 // Drop all local reset traffic and advance well past the close. No peer response
322 // has made the old ID safe during this interval, so it stays reserved.
323 let close_at = l.last;
324 let close_events = poll_without_delivering_network(&mut l, close_at)?;
325 assert!(
326 close_events
327 .iter()
328 .any(|event| matches!(event, Event::ChannelClose(id) if *id == old)),
329 "the old local channel should close"
330 );
331 poll_without_delivering_network(&mut l, close_at + Duration::from_secs(31))?;
332 
333 let replacement = l.direct_api().create_data_channel(ChannelConfig::default());
334 poll_without_delivering_network(
335 &mut l,
336 close_at + Duration::from_secs(31) + Duration::from_millis(1),
337 )?;
338 
339 let replacement_stream_id = l
340 .direct_api()
341 .sctp_stream_id_by_channel_id(replacement)
342 .expect("replacement should have a stream ID");
343 assert_ne!(
344 replacement_stream_id, old_stream_id,
345 "elapsed time cannot make an unconfirmed reset safe; another free ID should be used"
346 );
347 
348 Ok(())
349}
350 
351#[test]
352pub fn data_channel_flood() -> Result<(), RtcError> {
353 init_log();
354 init_crypto_default();
355 
356 let mut l = TestRtc::new(Peer::Left);
357 let mut r = TestRtc::new(Peer::Right);
358 
359 l.add_host_candidate((Ipv4Addr::new(1, 1, 1, 1), 1000).into());
360 r.add_host_candidate((Ipv4Addr::new(2, 2, 2, 2), 2000).into());
361 
362 let mut change = l.sdp_api();
363 let cid = change.add_channel("My little channel".into());
364 let (offer, pending) = change.apply().unwrap();
365 
366 let answer = r.rtc.sdp_api().accept_offer(offer)?;
367 l.rtc.sdp_api().accept_answer(pending, answer)?;
368 
369 loop {
370 if l.is_connected() || r.is_connected() {
371 break;
372 }
373 progress(&mut l, &mut r)?;
374 }
375 
376 let max = l.last.max(r.last);
377 l.last = max;
378 r.last = max;
379 
380 while l.channel(cid).is_none() {
381 progress(&mut l, &mut r)?;
382 }
383 
384 r.set_netem(NetemConfig::new().latency(Duration::from_millis(1000)));
385 
386 let mut count = 0;
387 
388 for _ in 0..10_000 {
389 let mut chan = l.channel(cid).unwrap();
390 let did_write = chan.write(true, &[0u8; 1400]).expect("to write string");
391 if did_write {
392 count += 1;
393 }
394 progress(&mut l, &mut r)?;
395 }
396 
397 loop {
398 progress(&mut l, &mut r)?;
399 
400 if l.duration() > Duration::from_secs(10) {
401 break;
402 }
403 }
404 assert!(count > 9000, "Too few events: {}", count);
405 
406 Ok(())
407}
408 
409#[test]
410pub fn channel_config_inband() -> Result<(), RtcError> {
411 init_log();
412 init_crypto_default();
413 
414 let mut l = TestRtc::new(Peer::Left);
415 let mut r = TestRtc::new(Peer::Right);
416 
417 l.add_host_candidate((Ipv4Addr::new(1, 1, 1, 1), 1000).into());
418 r.add_host_candidate((Ipv4Addr::new(2, 2, 2, 2), 2000).into());
419 
420 // Create in-band negotiated channel (DCEP)
421 let mut change = l.sdp_api();
422 let cid = change.add_channel("DCEP Channel".into());
423 let (offer, pending) = change.apply().unwrap();
424 
425 let answer = r.rtc.sdp_api().accept_offer(offer)?;
426 l.rtc.sdp_api().accept_answer(pending, answer)?;
427 
428 // Wait for connection
429 loop {
430 if l.is_connected() && r.is_connected() {
431 break;
432 }
433 progress(&mut l, &mut r)?;
434 }
435 
436 let max = l.last.max(r.last);
437 l.last = max;
438 r.last = max;
439 
440 let mut l_channel_opened = false;
441 let mut r_channel_opened = false;
442 let mut l_config_available_on_open = false;
443 let mut r_config_available_on_open = false;
444 
445 // Process events and verify config availability immediately when ChannelOpen is fired
446 loop {
447 progress(&mut l, &mut r)?;
448 
449 // Check L side events and collect channel ID if found
450 let mut l_found_id = None;
451 for (_, event) in &l.events {
452 if let Event::ChannelOpen(id, label) = event {
453 if *id == cid && label == "DCEP Channel" {
454 l_channel_opened = true;
455 l_found_id = Some(*id);
456 break;
457 }
458 }
459 }
460 
461 // Check R side events and collect channel ID if found
462 let mut r_found_id = None;
463 for (_, event) in &r.events {
464 if let Event::ChannelOpen(id, label) = event {
465 if label == "DCEP Channel" {
466 r_channel_opened = true;
467 r_found_id = Some(*id);
468 break;
469 }
470 }
471 }
472 
473 // Verify config is available immediately when ChannelOpen is emitted
474 if let Some(id) = l_found_id {
475 if let Some(channel) = l.channel(id) {
476 l_config_available_on_open = channel.config().is_some();
477 }
478 }
479 
480 if let Some(id) = r_found_id {
481 if let Some(channel) = r.channel(id) {
482 r_config_available_on_open = channel.config().is_some();
483 }
484 }
485 
486 if (l_channel_opened && r_channel_opened) || l.duration() > Duration::from_secs(10) {
487 break;
488 }
489 }
490 
491 assert!(l_channel_opened, "L side should receive ChannelOpen event");
492 assert!(r_channel_opened, "R side should receive ChannelOpen event");
493 assert!(
494 l_config_available_on_open,
495 "L side config should be available on ChannelOpen"
496 );
497 assert!(
498 r_config_available_on_open,
499 "R side config should be available on ChannelOpen"
500 );
501 
502 Ok(())
503}
504 
505#[test]
506pub fn channel_config_outband_local() -> Result<(), RtcError> {
507 init_log();
508 init_crypto_default();
509 
510 let mut l = TestRtc::new(Peer::Left);
511 let mut r = TestRtc::new(Peer::Right);
512 
513 l.add_host_candidate((Ipv4Addr::new(1, 1, 1, 1), 1000).into());
514 r.add_host_candidate((Ipv4Addr::new(2, 2, 2, 2), 2000).into());
515 
516 // Enable SCTP by adding a temporary channel (will be removed)
517 let mut change_l = l.sdp_api();
518 let _temp_cid = change_l.add_channel("temp".into());
519 let (offer, pending) = change_l.apply().unwrap();
520 
521 let answer = r.rtc.sdp_api().accept_offer(offer)?;
522 l.rtc.sdp_api().accept_answer(pending, answer)?;
523 
524 // Wait for connection
525 loop {
526 if l.is_connected() && r.is_connected() {
527 break;
528 }
529 progress(&mut l, &mut r)?;
530 }
531 
532 let max = l.last.max(r.last);
533 l.last = max;
534 r.last = max;
535 
536 // Wait for SCTP to be established first
537 loop {
538 progress(&mut l, &mut r)?;
539 
540 // Check for SCTP connection via any channel events
541 let connected = l
542 .events
543 .iter()
544 .any(|(_, e)| matches!(e, Event::ChannelOpen(_, _)))
545 || r.events
546 .iter()
547 .any(|(_, e)| matches!(e, Event::ChannelOpen(_, _)));
548 
549 if connected || l.duration() > Duration::from_secs(5) {
550 break;
551 }
552 }
553 
554 // Create out-of-band negotiated channel on both sides
555 let config = ChannelConfig {
556 negotiated: Some(10),
557 label: "OutOfBand Local".into(),
558 ..Default::default()
559 };
560 
561 let cid_l = l.direct_api().create_data_channel(config.clone());
562 let cid_r = r.direct_api().create_data_channel(config);
563 
564 // Allow some time for channels to be established
565 for _ in 0..10 {
566 progress(&mut l, &mut r)?;
567 }
568 
569 // Verify config is immediately available for locally created out-of-band channels
570 let l_channel = l.channel(cid_l).expect("L channel should be available");
571 let r_channel = r.channel(cid_r).expect("R channel should be available");
572 
573 assert!(
574 l_channel.config().is_some(),
575 "L side config should be immediately available for local out-of-band channel"
576 );
577 assert!(
578 r_channel.config().is_some(),
579 "R side config should be immediately available for local out-of-band channel"
580 );
581 
582 let l_config = l_channel.config().unwrap();
583 let r_config = r_channel.config().unwrap();
584 
585 assert_eq!(l_config.label, "OutOfBand Local");
586 assert_eq!(r_config.label, "OutOfBand Local");
587 assert_eq!(l_config.negotiated, Some(10));
588 assert_eq!(r_config.negotiated, Some(10));
589 
590 Ok(())
591}
592 
593#[test]
594pub fn channel_config_with_protocol() -> Result<(), RtcError> {
595 init_log();
596 init_crypto_default();
597 
598 let mut l = TestRtc::new(Peer::Left);
599 let mut r = TestRtc::new(Peer::Right);
600 
601 l.add_host_candidate((Ipv4Addr::new(1, 1, 1, 1), 1000).into());
602 r.add_host_candidate((Ipv4Addr::new(2, 2, 2, 2), 2000).into());
603 
604 let mut change = l.sdp_api();
605 let _temp_cid = change.add_channel("temp".into());
606 let (offer, pending) = change.apply().unwrap();
607 
608 let answer = r.rtc.sdp_api().accept_offer(offer)?;
609 l.rtc.sdp_api().accept_answer(pending, answer)?;
610 
611 // Wait for connection
612 loop {
613 if l.is_connected() && r.is_connected() {
614 break;
615 }
616 progress(&mut l, &mut r)?;
617 }
618 
619 let max = l.last.max(r.last);
620 l.last = max;
621 r.last = max;
622 
623 // Wait for SCTP to be established
624 loop {
625 progress(&mut l, &mut r)?;
626 let connected = l
627 .events
628 .iter()
629 .any(|(_, e)| matches!(e, Event::ChannelOpen(_, _)));
630 if connected || l.duration() > Duration::from_secs(5) {
631 break;
632 }
633 }
634 
635 // Create channels with custom protocol
636 let custom_protocol = "my-custom-protocol";
637 let config = ChannelConfig {
638 negotiated: Some(20),
639 protocol: custom_protocol.into(),
640 ..Default::default()
641 };
642 
643 let cid_l = l.direct_api().create_data_channel(config.clone());
644 let cid_r = r.direct_api().create_data_channel(config);
645 
646 for _ in 0..10 {
647 progress(&mut l, &mut r)?;
648 }
649 
650 // Verify protocol is correctly set on both sides
651 let l_channel = l.channel(cid_l).unwrap();
652 let r_channel = r.channel(cid_r).unwrap();
653 let l_config = l_channel.config().unwrap();
654 let r_config = r_channel.config().unwrap();
655 
656 assert_eq!(l_config.protocol, custom_protocol);
657 assert_eq!(r_config.protocol, custom_protocol);
658 
659 Ok(())
660}