use rtc::statistics::StatsSelector;
use std::sync::Arc;
use std::time::{Duration, Instant};
use webrtc::data_channel::{DataChannel, DataChannelEvent};
use webrtc::peer_connection::{
PeerConnection, PeerConnectionBuilder, PeerConnectionEventHandler, RTCIceGatheringState,
RTCPeerConnectionState,
};
use webrtc::runtime::channel;
mod common;
use common::{block_on, runtime, sleep, timeout};
struct OffererHandler {
gather_complete_tx: webrtc::runtime::Sender<()>,
connected_tx: webrtc::runtime::Sender<()>,
}
#[async_trait::async_trait]
impl PeerConnectionEventHandler for OffererHandler {
async fn on_ice_gathering_state_change(&self, state: RTCIceGatheringState) {
if state == RTCIceGatheringState::Complete {
let _ = self.gather_complete_tx.try_send(());
}
}
async fn on_connection_state_change(&self, state: RTCPeerConnectionState) {
if state == RTCPeerConnectionState::Connected {
let _ = self.connected_tx.try_send(());
}
}
}
struct AnswererHandler {
gather_complete_tx: webrtc::runtime::Sender<()>,
connected_tx: webrtc::runtime::Sender<()>,
data_channel_tx: webrtc::runtime::Sender<Arc<dyn DataChannel>>,
}
#[async_trait::async_trait]
impl PeerConnectionEventHandler for AnswererHandler {
async fn on_ice_gathering_state_change(&self, state: RTCIceGatheringState) {
if state == RTCIceGatheringState::Complete {
let _ = self.gather_complete_tx.try_send(());
}
}
async fn on_connection_state_change(&self, state: RTCPeerConnectionState) {
if state == RTCPeerConnectionState::Connected {
let _ = self.connected_tx.try_send(());
}
}
async fn on_data_channel(&self, data_channel: Arc<dyn DataChannel>) {
let _ = self.data_channel_tx.try_send(data_channel);
}
}
#[test]
fn test_peer_connection_statistics() {
block_on(async {
let runtime = runtime();
let (off_gather_tx, mut off_gather_rx) = channel(1);
let (off_conn_tx, mut off_conn_rx) = channel(1);
let offerer = PeerConnectionBuilder::new()
.with_handler(Arc::new(OffererHandler {
gather_complete_tx: off_gather_tx,
connected_tx: off_conn_tx,
}))
.with_runtime(runtime.clone())
.with_udp_addrs(vec!["127.0.0.1:0".to_owned()])
.build()
.await
.unwrap();
let offerer = Arc::new(offerer);
let (ans_gather_tx, mut ans_gather_rx) = channel(1);
let (ans_conn_tx, mut ans_conn_rx) = channel(1);
let (dc_tx, mut dc_rx) = channel(1);
let answerer = PeerConnectionBuilder::new()
.with_handler(Arc::new(AnswererHandler {
gather_complete_tx: ans_gather_tx,
connected_tx: ans_conn_tx,
data_channel_tx: dc_tx,
}))
.with_runtime(runtime.clone())
.with_udp_addrs(vec!["127.0.0.1:0".to_owned()])
.build()
.await
.unwrap();
let answerer = Arc::new(answerer);
let offer_dc = offerer.create_data_channel("stats-dc", None).await.unwrap();
let offer_dc_clone = offer_dc.clone();
let offer = offerer.create_offer(None).await.unwrap();
offerer.set_local_description(offer).await.unwrap();
timeout(Duration::from_secs(5), off_gather_rx.recv())
.await
.unwrap();
let offer_sdp = offerer.local_description().await.unwrap();
answerer.set_remote_description(offer_sdp).await.unwrap();
let answer = answerer.create_answer(None).await.unwrap();
answerer.set_local_description(answer).await.unwrap();
timeout(Duration::from_secs(5), ans_gather_rx.recv())
.await
.unwrap();
let answer_sdp = answerer.local_description().await.unwrap();
offerer.set_remote_description(answer_sdp).await.unwrap();
timeout(Duration::from_secs(5), off_conn_rx.recv())
.await
.unwrap();
timeout(Duration::from_secs(5), ans_conn_rx.recv())
.await
.unwrap();
let ans_dc = timeout(Duration::from_secs(5), dc_rx.recv())
.await
.unwrap()
.unwrap();
runtime.spawn(Box::pin(async move {
loop {
if let Some(DataChannelEvent::OnOpen) = offer_dc_clone.poll().await {
break;
}
}
offer_dc_clone.send_text("Hello stats!").await.unwrap();
}));
loop {
if let Some(DataChannelEvent::OnMessage(msg)) = ans_dc.poll().await {
assert_eq!(
String::from_utf8(msg.data.to_vec()).unwrap(),
"Hello stats!"
);
break;
}
}
sleep(Duration::from_millis(100)).await;
let offerer_stats = offerer.get_stats(Instant::now(), StatsSelector::None).await;
assert!(
!offerer_stats.is_empty(),
"Offerer stats report should not be empty"
);
let pc_stats = offerer_stats
.peer_connection()
.expect("PeerConnection stats missing");
assert_eq!(pc_stats.data_channels_opened, 1);
let dc_stats_list: Vec<_> = offerer_stats.data_channels().collect();
assert_eq!(
dc_stats_list.len(),
1,
"Expected 1 data channel stats entry"
);
let dc_stats = dc_stats_list[0];
assert_eq!(dc_stats.label, "stats-dc");
assert!(dc_stats.messages_sent > 0, "Expected messages_sent > 0");
let transport_stats = offerer_stats.transport().expect("Transport stats missing");
assert!(
transport_stats.bytes_sent > 0,
"Expected transport bytes_sent > 0"
);
let pair_stats_list: Vec<_> = offerer_stats.candidate_pairs().collect();
assert!(!pair_stats_list.is_empty(), "Candidate pair stats missing");
let answerer_stats = answerer
.get_stats(Instant::now(), StatsSelector::None)
.await;
assert!(
!answerer_stats.is_empty(),
"Answerer stats report should not be empty"
);
let ans_pc_stats = answerer_stats
.peer_connection()
.expect("Answerer PeerConnection stats missing");
assert_eq!(ans_pc_stats.data_channels_opened, 1);
offerer.close().await.unwrap();
answerer.close().await.unwrap();
});
}