use std::ffi::{c_void, CStr, CString};
use std::os::raw::{c_char, c_int};
use std::sync::mpsc::{sync_channel, SyncSender};
use std::sync::Mutex;
use std::time::Duration;
use reactor_webrtc_sys::ReactorStatEntry;
use crate::encoded::FrameTransform;
use crate::media::{MediaKind, Track};
use crate::observer::ObserverState;
use crate::{Error, Result};
const OP_TIMEOUT: Duration = Duration::from_secs(10);
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SdpType {
Offer,
PrAnswer,
Answer,
Rollback,
}
impl SdpType {
fn as_str(self) -> &'static str {
match self {
SdpType::Offer => "offer",
SdpType::PrAnswer => "pranswer",
SdpType::Answer => "answer",
SdpType::Rollback => "rollback",
}
}
fn from_str(s: &str) -> Option<Self> {
match s {
"offer" => Some(SdpType::Offer),
"pranswer" => Some(SdpType::PrAnswer),
"answer" => Some(SdpType::Answer),
"rollback" => Some(SdpType::Rollback),
_ => None,
}
}
}
#[derive(Debug, Clone)]
pub struct SessionDescription {
pub kind: SdpType,
pub sdp: String,
}
#[derive(Debug, Clone)]
pub struct IceCandidate {
pub candidate: String,
pub sdp_mid: Option<String>,
pub sdp_mline_index: Option<u16>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum TransceiverDirection {
SendRecv,
SendOnly,
RecvOnly,
Inactive,
}
impl TransceiverDirection {
fn to_raw(self) -> c_int {
match self {
TransceiverDirection::SendRecv => 0,
TransceiverDirection::SendOnly => 1,
TransceiverDirection::RecvOnly => 2,
TransceiverDirection::Inactive => 3,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum IceGatheringState {
New,
Gathering,
Complete,
}
impl IceGatheringState {
pub(crate) fn from_raw(state: c_int) -> Self {
match state {
1 => IceGatheringState::Gathering,
2 => IceGatheringState::Complete,
_ => IceGatheringState::New,
}
}
}
pub struct Transceiver {
raw: *mut reactor_webrtc_sys::RtpTransceiver,
}
unsafe impl Send for Transceiver {}
unsafe impl Sync for Transceiver {}
impl Transceiver {
pub(crate) fn from_raw(raw: *mut reactor_webrtc_sys::RtpTransceiver) -> Self {
Self { raw }
}
pub fn kind(&self) -> MediaKind {
let k = unsafe { reactor_webrtc_sys::reactor_webrtc_rtp_transceiver_media_kind(self.raw) };
MediaKind::from_raw(k)
}
pub fn mid(&self) -> Option<String> {
let mut buf = [0u8; 256];
let n = unsafe {
reactor_webrtc_sys::reactor_webrtc_rtp_transceiver_mid(
self.raw,
buf.as_mut_ptr() as *mut c_char,
buf.len() as c_int,
)
};
if n < 0 {
None
} else {
Some(
unsafe { CStr::from_ptr(buf.as_ptr() as *const c_char) }
.to_string_lossy()
.into_owned(),
)
}
}
pub fn set_track(&self, track: &Track) -> Result<()> {
let ok = unsafe {
reactor_webrtc_sys::reactor_webrtc_rtp_transceiver_set_track(self.raw, track.raw())
};
if ok == 1 {
Ok(())
} else {
Err(Error::Webrtc("transceiver set_track failed".into()))
}
}
pub fn set_direction(&self, direction: TransceiverDirection) -> Result<()> {
let ok = unsafe {
reactor_webrtc_sys::reactor_webrtc_rtp_transceiver_set_direction(
self.raw,
direction.to_raw() as c_int,
)
};
if ok == 1 {
Ok(())
} else {
Err(Error::Webrtc("transceiver set_direction failed".into()))
}
}
pub fn set_sender_transform(&self, transform: &FrameTransform) -> Result<()> {
let ok = unsafe {
reactor_webrtc_sys::reactor_webrtc_rtp_transceiver_set_sender_transform(
self.raw,
transform.raw(),
)
};
if ok == 1 {
Ok(())
} else {
Err(Error::Webrtc(
"transceiver set_sender_transform failed".into(),
))
}
}
pub fn set_receiver_transform(&self, transform: &FrameTransform) -> Result<()> {
let ok = unsafe {
reactor_webrtc_sys::reactor_webrtc_rtp_transceiver_set_receiver_transform(
self.raw,
transform.raw(),
)
};
if ok == 1 {
Ok(())
} else {
Err(Error::Webrtc(
"transceiver set_receiver_transform failed".into(),
))
}
}
}
impl Drop for Transceiver {
fn drop(&mut self) {
unsafe { reactor_webrtc_sys::reactor_webrtc_rtp_transceiver_destroy(self.raw) }
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PeerConnectionState {
New,
Connecting,
Connected,
Disconnected,
Failed,
Closed,
}
impl PeerConnectionState {
pub(crate) fn from_raw(state: c_int) -> Self {
match state {
0 => PeerConnectionState::New,
1 => PeerConnectionState::Connecting,
2 => PeerConnectionState::Connected,
3 => PeerConnectionState::Disconnected,
4 => PeerConnectionState::Failed,
_ => PeerConnectionState::Closed,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum DataChannelState {
Connecting,
Open,
Closing,
Closed,
}
impl DataChannelState {
fn from_raw(v: c_int) -> Self {
match v {
0 => DataChannelState::Connecting,
1 => DataChannelState::Open,
2 => DataChannelState::Closing,
_ => DataChannelState::Closed,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum IceCandidatePairState {
Waiting,
InProgress,
Failed,
Succeeded,
Cancelled,
}
impl IceCandidatePairState {
fn from_raw(v: c_int) -> Self {
match v {
1 => IceCandidatePairState::InProgress,
2 => IceCandidatePairState::Failed,
3 => IceCandidatePairState::Succeeded,
4 => IceCandidatePairState::Cancelled,
_ => IceCandidatePairState::Waiting,
}
}
}
#[derive(Debug, Clone)]
pub struct InboundRtpStats {
pub ssrc: u32,
pub packets_received: u32,
pub bytes_received: u64,
pub jitter_s: f64,
pub packets_lost: i32,
pub nack_count: u32,
pub total_decode_time_s: f64,
}
#[derive(Debug, Clone)]
pub struct OutboundRtpStats {
pub ssrc: u32,
pub packets_sent: u32,
pub bytes_sent: u64,
pub target_bitrate_bps: f64,
pub round_trip_time_s: f64,
pub retransmitted_packets_sent: u32,
}
#[derive(Debug, Clone)]
pub struct IceCandidatePairStats {
pub current_round_trip_time_s: f64,
pub priority: u64,
pub state: IceCandidatePairState,
}
#[derive(Debug, Clone, Default)]
pub struct StatsReport {
pub inbound_rtp: Vec<InboundRtpStats>,
pub outbound_rtp: Vec<OutboundRtpStats>,
pub candidate_pairs: Vec<IceCandidatePairStats>,
}
type MessageCb = Box<dyn for<'a> FnMut(&'a [u8], bool) + Send>;
type EventCb = Box<dyn FnMut() + Send>;
type StateCb = Box<dyn FnMut(DataChannelState) + Send>;
#[derive(Default)]
struct DcObserverState {
on_message: Option<Mutex<MessageCb>>,
on_state_change: Option<Mutex<StateCb>>,
on_buffered_amount_low: Option<Mutex<EventCb>>,
}
extern "C" fn dc_on_message(ud: *mut c_void, data: *const u8, len: usize, binary: c_int) {
let st = unsafe { &*(ud as *const DcObserverState) };
if let Some(m) = &st.on_message {
let bytes = unsafe { std::slice::from_raw_parts(data, len) };
if let Ok(mut cb) = m.lock() {
cb(bytes, binary != 0);
}
}
}
extern "C" fn dc_on_state_change(ud: *mut c_void, state: c_int) {
let st = unsafe { &*(ud as *const DcObserverState) };
if let Some(m) = &st.on_state_change {
if let Ok(mut cb) = m.lock() {
cb(DataChannelState::from_raw(state));
}
}
}
extern "C" fn dc_on_buffered_amount_low(ud: *mut c_void) {
let st = unsafe { &*(ud as *const DcObserverState) };
if let Some(m) = &st.on_buffered_amount_low {
if let Ok(mut cb) = m.lock() {
cb();
}
}
}
pub struct DataChannel {
raw: *mut reactor_webrtc_sys::DataChannel,
observer: Option<Box<DcObserverState>>,
}
unsafe impl Send for DataChannel {}
unsafe impl Sync for DataChannel {}
impl DataChannel {
pub(crate) fn from_raw(raw: *mut reactor_webrtc_sys::DataChannel) -> Self {
Self {
raw,
observer: None,
}
}
pub fn label(&self) -> String {
let mut buf = [0u8; 256];
let n = unsafe {
reactor_webrtc_sys::reactor_webrtc_data_channel_label(
self.raw,
buf.as_mut_ptr() as *mut c_char,
buf.len() as c_int,
)
};
if n < 0 {
String::new()
} else {
unsafe { CStr::from_ptr(buf.as_ptr() as *const c_char) }
.to_string_lossy()
.into_owned()
}
}
pub fn state(&self) -> DataChannelState {
DataChannelState::from_raw(unsafe {
reactor_webrtc_sys::reactor_webrtc_data_channel_state(self.raw)
})
}
pub fn buffered_amount(&self) -> u64 {
unsafe { reactor_webrtc_sys::reactor_webrtc_data_channel_buffered_amount(self.raw) }
}
pub fn set_buffered_amount_low_threshold(&self, threshold: u64) {
unsafe {
reactor_webrtc_sys::reactor_webrtc_data_channel_set_low_threshold(self.raw, threshold)
}
}
pub fn send(&self, data: &[u8], binary: bool) -> Result<()> {
let ok = unsafe {
reactor_webrtc_sys::reactor_webrtc_data_channel_send(
self.raw,
data.as_ptr(),
data.len(),
binary as c_int,
)
};
if ok == 1 {
Ok(())
} else {
Err(Error::Webrtc("data channel send failed".into()))
}
}
pub fn on_message(&mut self, cb: impl for<'a> FnMut(&'a [u8], bool) + Send + 'static) {
self.observer
.get_or_insert_with(Default::default)
.on_message = Some(Mutex::new(Box::new(cb)));
self.reregister();
}
pub fn on_state_change(&mut self, cb: impl FnMut(DataChannelState) + Send + 'static) {
self.observer
.get_or_insert_with(Default::default)
.on_state_change = Some(Mutex::new(Box::new(cb)));
self.reregister();
}
pub fn on_open(&mut self, cb: impl FnMut() + Send + 'static) {
let mut cb = cb;
self.on_state_change(move |s| {
if s == DataChannelState::Open {
cb();
}
});
}
pub fn on_close(&mut self, cb: impl FnMut() + Send + 'static) {
let mut cb = cb;
self.on_state_change(move |s| {
if s == DataChannelState::Closed {
cb();
}
});
}
pub fn on_buffered_amount_low(&mut self, cb: impl FnMut() + Send + 'static) {
self.observer
.get_or_insert_with(Default::default)
.on_buffered_amount_low = Some(Mutex::new(Box::new(cb)));
self.reregister();
}
fn reregister(&mut self) {
if let Some(state) = &self.observer {
let ud = &**state as *const DcObserverState as *mut c_void;
unsafe {
reactor_webrtc_sys::reactor_webrtc_data_channel_register_observer(
self.raw,
ud,
dc_on_message,
dc_on_state_change,
dc_on_buffered_amount_low,
);
}
}
}
}
impl Drop for DataChannel {
fn drop(&mut self) {
unsafe { reactor_webrtc_sys::reactor_webrtc_data_channel_destroy(self.raw) }
}
}
type SdpTx = SyncSender<Result<SessionDescription>>;
type CompleteTx = SyncSender<Result<()>>;
extern "C" fn sdp_ok(ud: *mut c_void, ty: *const c_char, sdp: *const c_char) {
let tx = unsafe { &*(ud as *const SdpTx) };
let kind = unsafe { CStr::from_ptr(ty) }.to_string_lossy();
let sdp = unsafe { CStr::from_ptr(sdp) }
.to_string_lossy()
.into_owned();
let result = match SdpType::from_str(&kind) {
Some(kind) => Ok(SessionDescription { kind, sdp }),
None => Err(Error::Webrtc(format!("unknown sdp type: {kind}"))),
};
let _ = tx.try_send(result);
}
extern "C" fn sdp_err(ud: *mut c_void, message: *const c_char) {
let tx = unsafe { &*(ud as *const SdpTx) };
let msg = unsafe { CStr::from_ptr(message) }
.to_string_lossy()
.into_owned();
let _ = tx.try_send(Err(Error::Webrtc(msg)));
}
extern "C" fn complete_cb(ud: *mut c_void, error: *const c_char) {
let tx = unsafe { &*(ud as *const CompleteTx) };
let r = if error.is_null() {
Ok(())
} else {
Err(Error::Webrtc(
unsafe { CStr::from_ptr(error) }
.to_string_lossy()
.into_owned(),
))
};
let _ = tx.try_send(r);
}
type StatsTx = SyncSender<StatsReport>;
extern "C" fn stats_cb(ud: *mut c_void, entries: *const ReactorStatEntry, count: c_int) {
let tx = unsafe { &*(ud as *const StatsTx) };
let slice = if entries.is_null() || count <= 0 {
&[][..]
} else {
unsafe { std::slice::from_raw_parts(entries, count as usize) }
};
let mut report = StatsReport::default();
for e in slice {
match e.kind {
0 => report.inbound_rtp.push(InboundRtpStats {
ssrc: e.ssrc,
packets_received: e.packets_received,
bytes_received: e.bytes_received,
jitter_s: e.jitter,
packets_lost: e.packets_lost,
nack_count: e.nack_count,
total_decode_time_s: e.total_decode_time,
}),
1 => report.outbound_rtp.push(OutboundRtpStats {
ssrc: e.ssrc,
packets_sent: e.packets_sent,
bytes_sent: e.bytes_sent,
target_bitrate_bps: e.target_bitrate,
round_trip_time_s: e.round_trip_time,
retransmitted_packets_sent: e.retransmitted_packets_sent,
}),
2 => report.candidate_pairs.push(IceCandidatePairStats {
current_round_trip_time_s: e.current_round_trip_time,
priority: e.priority,
state: IceCandidatePairState::from_raw(e.pair_state),
}),
_ => {}
}
}
let _ = tx.try_send(report);
}
fn run_stats(call: impl FnOnce(*mut c_void)) -> Result<StatsReport> {
let (tx, rx) = sync_channel::<StatsReport>(1);
let p = Box::into_raw(Box::new(tx));
call(p as *mut c_void);
let r = rx.recv_timeout(OP_TIMEOUT);
drop(unsafe { Box::from_raw(p) });
r.map_err(|_| Error::Webrtc("get_stats timed out".into()))
}
fn run_sdp(call: impl FnOnce(*mut c_void)) -> Result<SessionDescription> {
let (tx, rx) = sync_channel::<Result<SessionDescription>>(1);
let p = Box::into_raw(Box::new(tx));
call(p as *mut c_void);
let r = rx.recv_timeout(OP_TIMEOUT);
drop(unsafe { Box::from_raw(p) });
r.map_err(|_| Error::Webrtc("sdp operation timed out".into()))?
}
fn run_complete(call: impl FnOnce(*mut c_void)) -> Result<()> {
let (tx, rx) = sync_channel::<Result<()>>(1);
let p = Box::into_raw(Box::new(tx));
call(p as *mut c_void);
let r = rx.recv_timeout(OP_TIMEOUT);
drop(unsafe { Box::from_raw(p) });
r.map_err(|_| Error::Webrtc("operation timed out".into()))?
}
pub struct PeerConnection {
raw: *mut reactor_webrtc_sys::PeerConnection,
_observer: Box<ObserverState>,
}
unsafe impl Send for PeerConnection {}
unsafe impl Sync for PeerConnection {}
impl PeerConnection {
pub(crate) fn new(
raw: *mut reactor_webrtc_sys::PeerConnection,
observer: Box<ObserverState>,
) -> Self {
Self {
raw,
_observer: observer,
}
}
pub fn create_offer(&self) -> Result<SessionDescription> {
run_sdp(|ud| unsafe {
reactor_webrtc_sys::reactor_webrtc_peer_connection_create_offer(
self.raw, ud, sdp_ok, sdp_err,
)
})
}
pub fn create_answer(&self) -> Result<SessionDescription> {
run_sdp(|ud| unsafe {
reactor_webrtc_sys::reactor_webrtc_peer_connection_create_answer(
self.raw, ud, sdp_ok, sdp_err,
)
})
}
pub fn set_local_description(&self, sdp: &SessionDescription) -> Result<()> {
self.set_description(sdp, true)
}
pub fn set_remote_description(&self, sdp: &SessionDescription) -> Result<()> {
self.set_description(sdp, false)
}
fn set_description(&self, sdp: &SessionDescription, local: bool) -> Result<()> {
let ty = CString::new(sdp.kind.as_str()).unwrap();
let body = CString::new(sdp.sdp.as_str())
.map_err(|_| Error::Webrtc("sdp contains a NUL byte".into()))?;
run_complete(|ud| unsafe {
if local {
reactor_webrtc_sys::reactor_webrtc_peer_connection_set_local_description(
self.raw,
ty.as_ptr(),
body.as_ptr(),
ud,
complete_cb,
)
} else {
reactor_webrtc_sys::reactor_webrtc_peer_connection_set_remote_description(
self.raw,
ty.as_ptr(),
body.as_ptr(),
ud,
complete_cb,
)
}
})
}
pub fn add_ice_candidate(&self, candidate: &IceCandidate) -> Result<()> {
let mid = CString::new(candidate.sdp_mid.clone().unwrap_or_default()).unwrap_or_default();
let cand = CString::new(candidate.candidate.as_str())
.map_err(|_| Error::Webrtc("candidate contains a NUL byte".into()))?;
let idx = candidate.sdp_mline_index.unwrap_or(0) as c_int;
run_complete(|ud| unsafe {
reactor_webrtc_sys::reactor_webrtc_peer_connection_add_ice_candidate(
self.raw,
mid.as_ptr(),
idx,
cand.as_ptr(),
ud,
complete_cb,
)
})
}
pub fn add_track(&self, track: &Track) -> Result<()> {
let ok = unsafe {
reactor_webrtc_sys::reactor_webrtc_peer_connection_add_track(self.raw, track.raw())
};
if ok == 1 {
Ok(())
} else {
Err(Error::Webrtc("add_track failed".into()))
}
}
pub fn add_transceiver(
&self,
kind: MediaKind,
direction: TransceiverDirection,
) -> Result<Transceiver> {
let media_kind = match kind {
MediaKind::Audio => 0,
MediaKind::Video => 1,
MediaKind::Unknown => {
return Err(Error::Webrtc("add_transceiver needs audio or video".into()))
}
};
let raw = unsafe {
reactor_webrtc_sys::reactor_webrtc_peer_connection_add_transceiver(
self.raw,
media_kind,
direction.to_raw(),
)
};
if raw.is_null() {
Err(Error::Webrtc("add_transceiver failed".into()))
} else {
Ok(Transceiver::from_raw(raw))
}
}
pub fn transceivers(&self) -> Vec<Transceiver> {
let n = unsafe {
reactor_webrtc_sys::reactor_webrtc_peer_connection_transceiver_count(self.raw)
};
(0..n)
.filter_map(|i| {
let raw = unsafe {
reactor_webrtc_sys::reactor_webrtc_peer_connection_get_transceiver(self.raw, i)
};
(!raw.is_null()).then(|| Transceiver::from_raw(raw))
})
.collect()
}
pub fn create_data_channel(&self, label: &str) -> Result<DataChannel> {
let label =
CString::new(label).map_err(|_| Error::Webrtc("label has a NUL byte".into()))?;
let raw = unsafe {
reactor_webrtc_sys::reactor_webrtc_peer_connection_create_data_channel(
self.raw,
label.as_ptr(),
)
};
if raw.is_null() {
Err(Error::Webrtc("create_data_channel returned null".into()))
} else {
Ok(DataChannel::from_raw(raw))
}
}
pub fn get_stats(&self) -> Result<StatsReport> {
run_stats(|ud| unsafe {
reactor_webrtc_sys::reactor_webrtc_peer_connection_get_stats(self.raw, ud, stats_cb)
})
}
}
impl Drop for PeerConnection {
fn drop(&mut self) {
unsafe { reactor_webrtc_sys::reactor_webrtc_peer_connection_destroy(self.raw) }
}
}