use std::sync::Arc;
use std::sync::atomic::Ordering;
use rtc::peer_connection::configuration::media_engine::MIME_TYPE_VP8;
use rtc::rtp_transceiver::RTCRtpTransceiverDirection;
use webrtc::peer_connection::RTCSdpType;
mod common;
use common::{HOST, Peer, SIGNAL_PORT, SendTrack};
const VERIFY_COUNT: usize = 50;
#[tokio::test]
async fn test_rtp_uni_direction_0sendonly_1recvonly() -> anyhow::Result<()> {
let room_id = sfu::RoomId::new_v4();
let mut publisher = common::connect(HOST, SIGNAL_PORT, room_id, 0).await?;
let mut subscriber = common::connect(HOST, SIGNAL_PORT, room_id, 1).await?;
let send_track = Arc::new(
publisher
.add_track(
MIME_TYPE_VP8,
"video_track",
RTCRtpTransceiverDirection::Sendonly,
)
.await?,
);
publisher.renegotiate().await?;
assert_eq!(RTCSdpType::Answer, publisher.next_sdp().await?.sdp_type);
assert_eq!(RTCSdpType::Offer, subscriber.next_sdp().await?.sdp_type);
let (writer, stop) = send_track.spawn_writer();
let remote = subscriber.next_track().await?;
common::verify_rtp_flow(&remote, VERIFY_COUNT).await?;
stop.store(true, Ordering::Relaxed);
let _ = writer.await;
publisher.close().await?;
subscriber.close().await?;
Ok(())
}
async fn test_rtp_bi_direction_sendrecv(endpoint_count: u64) -> anyhow::Result<()> {
let room_id = sfu::RoomId::new_v4();
let mut peers = Vec::new();
for client_id in 0..endpoint_count {
peers.push(common::connect(HOST, SIGNAL_PORT, room_id, client_id).await?);
}
let mut send_tracks = Vec::new();
for peer in &mut peers {
let track_id = format!("video_track_{}", peer.client_id);
let send_track = Arc::new(
peer.add_track(
MIME_TYPE_VP8,
&track_id,
RTCRtpTransceiverDirection::Sendonly,
)
.await?,
);
peer.renegotiate().await?;
wait_for_answer(peer).await?;
send_tracks.push(send_track);
tokio::time::sleep(std::time::Duration::from_millis(500)).await;
}
let mut writers = Vec::new();
let mut stops = Vec::new();
for send_track in &send_tracks {
let (writer, stop) = SendTrack::spawn_writer(send_track.clone());
writers.push(writer);
stops.push(stop);
}
for peer in &mut peers {
for _ in 0..(endpoint_count - 1) {
let remote = peer.next_track().await?;
common::verify_rtp_flow(&remote, VERIFY_COUNT).await?;
}
}
for stop in &stops {
stop.store(true, Ordering::Relaxed);
}
for writer in writers {
let _ = writer.await;
}
for peer in peers {
peer.close().await?;
}
Ok(())
}
async fn wait_for_answer(peer: &mut Peer) -> anyhow::Result<()> {
loop {
if peer.next_sdp().await?.sdp_type == RTCSdpType::Answer {
return Ok(());
}
}
}
#[tokio::test]
async fn test_rtp_2p_bi_direction_sendrecv() -> anyhow::Result<()> {
test_rtp_bi_direction_sendrecv(2).await
}
#[tokio::test]
async fn test_rtp_3p_bi_direction_sendrecv() -> anyhow::Result<()> {
test_rtp_bi_direction_sendrecv(3).await
}
#[tokio::test]
async fn test_multiple_concurrent_media_rooms() -> anyhow::Result<()> {
let mut room_handles = Vec::new();
for room_idx in 0..4 {
let room_id = sfu::RoomId::new_v4();
let handle = tokio::spawn(async move {
let mut publisher = common::connect(HOST, SIGNAL_PORT, room_id, 0).await?;
let mut subscriber = common::connect(HOST, SIGNAL_PORT, room_id, 1).await?;
let track_id = format!("video_track_room_{}_{}", room_idx, room_id);
let send_track = Arc::new(
publisher
.add_track(
MIME_TYPE_VP8,
&track_id,
RTCRtpTransceiverDirection::Sendonly,
)
.await?,
);
publisher.renegotiate().await?;
wait_for_answer(&mut publisher).await?;
assert_eq!(RTCSdpType::Offer, subscriber.next_sdp().await?.sdp_type);
let (writer, stop) = SendTrack::spawn_writer(send_track);
let remote = subscriber.next_track().await?;
common::verify_rtp_flow(&remote, VERIFY_COUNT).await?;
stop.store(true, Ordering::Relaxed);
let _ = writer.await;
publisher.close().await?;
subscriber.close().await?;
anyhow::Result::<()>::Ok(())
});
room_handles.push(handle);
}
for handle in room_handles {
handle.await??;
}
Ok(())
}