This commit is contained in:
Théo Monnom
2022-09-16 11:45:39 +02:00
parent f643f458f8
commit 06587d154f
20 changed files with 261 additions and 267 deletions
@@ -7,8 +7,10 @@
#include <memory>
#include "api/data_channel_interface.h"
#include "rust_types.h"
namespace livekit {
using NativeDataChannelInit = webrtc::DataChannelInit;
class DataChannel {
public:
@@ -18,9 +20,12 @@ namespace livekit {
rtc::scoped_refptr<webrtc::DataChannelInterface> data_channel_;
};
std::unique_ptr<NativeDataChannelInit> create_data_channel_init(DataChannelInit init);
static std::unique_ptr<DataChannel> _unique_data_channel(){
return nullptr; // Ignore
}
} // livekit
#endif //CLIENT_SDK_NATIVE_DATA_CHANNEL_H
@@ -40,10 +40,6 @@ namespace livekit {
return nullptr; // Ignore
}
static std::shared_ptr<SessionDescription> _shared_session_description(){
return nullptr; // Ignore
}
// SetCreateSdpObserver
class NativeCreateSdpObserver : public webrtc::CreateSessionDescriptionObserver {
@@ -19,11 +19,12 @@ namespace livekit {
public:
explicit PeerConnection(rtc::scoped_refptr<webrtc::PeerConnectionInterface> peer_connection, std::unique_ptr<NativePeerConnectionObserver> observer);
void close();
void create_offer(std::unique_ptr<NativeCreateSdpObserverHandle> observer, RTCOfferAnswerOptions options);
void create_answer(std::unique_ptr<NativeCreateSdpObserverHandle> observer, RTCOfferAnswerOptions options);
void set_local_description(std::unique_ptr<SessionDescription> desc, std::unique_ptr<NativeSetLocalSdpObserverHandle> observer);
void set_remote_description(std::unique_ptr<SessionDescription> desc, std::unique_ptr<NativeSetRemoteSdpObserverHandle> observer);
std::unique_ptr<DataChannel> create_data_channel(rust::String label, std::unique_ptr<NativeDataChannelInit> init);
void close();
private:
rtc::scoped_refptr<webrtc::PeerConnectionInterface> peer_connection_;
@@ -18,6 +18,7 @@ namespace livekit {
// Shared types
struct RTCOfferAnswerOptions;
struct RTCError;
struct DataChannelInit;
}
#endif //RUST_TYPES_H
@@ -2,11 +2,35 @@
// Created by Théo Monnom on 01/09/2022.
//
#include <utility>
#include "livekit/data_channel.h"
#include "libwebrtc-sys/src/data_channel.rs.h"
namespace livekit {
DataChannel::DataChannel(rtc::scoped_refptr<webrtc::DataChannelInterface> data_channel) : data_channel_(data_channel) {
DataChannel::DataChannel(rtc::scoped_refptr<webrtc::DataChannelInterface> data_channel) : data_channel_(std::move(data_channel)) {
}
std::unique_ptr<NativeDataChannelInit> create_data_channel_init(DataChannelInit init) {
auto rtc_init = std::make_unique<webrtc::DataChannelInit>();
rtc_init->id = init.id;
rtc_init->negotiated = init.negotiated;
rtc_init->ordered = init.ordered;
rtc_init->protocol = init.protocol.c_str();
rtc_init->reliable = init.reliable;
if(init.has_max_retransmit_time)
rtc_init->maxRetransmitTime = init.max_retransmit_time;
if(init.has_max_retransmits)
rtc_init->maxRetransmits = init.max_retransmits;
if(init.has_priority)
rtc_init->priority = static_cast<webrtc::Priority>(init.priority);
return rtc_init;
}
} // livekit
@@ -2,10 +2,40 @@ use cxx::UniquePtr;
#[cxx::bridge(namespace = "livekit")]
pub mod ffi {
#[derive(Debug)]
#[repr(u32)]
enum Priority {
VeryLow,
Low,
Medium,
High,
}
#[derive(Debug)]
//#[allow(deprecated)]
pub struct DataChannelInit {
#[deprecated]
reliable: bool,
ordered: bool,
has_max_retransmit_time: bool,
max_retransmit_time: i32,
has_max_retransmits: bool,
max_retransmits: i32,
protocol: String,
negotiated: bool,
id: i32,
has_priority: bool,
priority: Priority
}
unsafe extern "C++" {
include!("livekit/data_channel.h");
type DataChannel;
type NativeDataChannelInit;
fn create_data_channel_init(init: DataChannelInit) -> UniquePtr<NativeDataChannelInit>;
fn _unique_data_channel() -> UniquePtr<DataChannel>; // Ignore
}
@@ -1,5 +1,5 @@
use std::fmt::{Debug, Formatter};
use cxx::UniquePtr;
use std::fmt::{Debug, Formatter};
use crate::rtc_error::ffi::RTCError;
@@ -33,6 +33,7 @@ pub mod ffi {
type NativeSetRemoteSdpObserverHandle;
fn stringify(self: &SessionDescription) -> String;
fn clone(self: &SessionDescription) -> UniquePtr<SessionDescription>;
fn create_native_create_sdp_observer(
observer: Box<CreateSdpObserverWrapper>,
@@ -45,7 +46,6 @@ pub mod ffi {
) -> UniquePtr<NativeSetRemoteSdpObserverHandle>;
fn _unique_ice_candidate() -> UniquePtr<IceCandidate>; // Ignore
fn _shared_session_description() -> SharedPtr<SessionDescription>; // Ignore
fn _unique_session_description() -> UniquePtr<SessionDescription>; // Ignore
}
}
@@ -4,6 +4,7 @@
#include "livekit/peer_connection.h"
#include "libwebrtc-sys/src/peer_connection.rs.h"
#include "livekit/rtc_error.h"
namespace livekit {
@@ -26,10 +27,6 @@ namespace livekit {
}
void PeerConnection::close() {
peer_connection_->Close();
}
void PeerConnection::create_offer(std::unique_ptr<NativeCreateSdpObserverHandle> observer_handle, RTCOfferAnswerOptions options) {
peer_connection_->CreateOffer(observer_handle->observer.get(), toNativeOfferAnswerOptions(options));
}
@@ -46,6 +43,20 @@ namespace livekit {
peer_connection_->SetRemoteDescription(desc->clone()->release(), observer->observer);
}
std::unique_ptr<DataChannel> PeerConnection::create_data_channel(rust::String label, std::unique_ptr<NativeDataChannelInit> init) {
auto result = peer_connection_->CreateDataChannelOrError(label.c_str(), init.get());
if(!result.ok()) {
throw std::runtime_error(serialize_error(to_error(result.error())));
}
return std::make_unique<DataChannel>(result.value());
}
void PeerConnection::close() {
peer_connection_->Close();
}
/* Observer */
NativePeerConnectionObserver::NativePeerConnectionObserver(rust::Box<PeerConnectionObserverWrapper> observer) : observer_(std::move(observer)) {
@@ -1,5 +1,3 @@
use std::cell::RefCell;
use std::rc::Rc;
use crate::candidate::ffi::Candidate;
use crate::data_channel::ffi::DataChannel;
use crate::jsep::ffi::IceCandidate;
@@ -143,42 +141,45 @@ pub mod ffi {
extern "Rust" {
type PeerConnectionObserverWrapper;
fn on_signaling_change(self: &mut PeerConnectionObserverWrapper, new_state: SignalingState);
fn on_add_stream(
unsafe fn on_signaling_change(
self: &mut PeerConnectionObserverWrapper,
new_state: SignalingState,
);
unsafe fn on_add_stream(
self: &mut PeerConnectionObserverWrapper,
stream: UniquePtr<MediaStreamInterface>,
);
fn on_remove_stream(
unsafe fn on_remove_stream(
self: &mut PeerConnectionObserverWrapper,
stream: UniquePtr<MediaStreamInterface>,
);
fn on_data_channel(
unsafe fn on_data_channel(
self: &mut PeerConnectionObserverWrapper,
data_channel: UniquePtr<DataChannel>,
);
fn on_renegotiation_needed(self: &mut PeerConnectionObserverWrapper);
fn on_negotiation_needed_event(self: &mut PeerConnectionObserverWrapper, event: u32);
fn on_ice_connection_change(
unsafe fn on_renegotiation_needed(self: &mut PeerConnectionObserverWrapper);
unsafe fn on_negotiation_needed_event(self: &mut PeerConnectionObserverWrapper, event: u32);
unsafe fn on_ice_connection_change(
self: &mut PeerConnectionObserverWrapper,
new_state: IceConnectionState,
);
fn on_standardized_ice_connection_change(
unsafe fn on_standardized_ice_connection_change(
self: &mut PeerConnectionObserverWrapper,
new_state: IceConnectionState,
);
fn on_connection_change(
unsafe fn on_connection_change(
self: &mut PeerConnectionObserverWrapper,
new_state: PeerConnectionState,
);
fn on_ice_gathering_change(
unsafe fn on_ice_gathering_change(
self: &mut PeerConnectionObserverWrapper,
new_state: IceGatheringState,
);
fn on_ice_candidate(
unsafe fn on_ice_candidate(
self: &mut PeerConnectionObserverWrapper,
candidate: UniquePtr<IceCandidate>,
);
fn on_ice_candidate_error(
unsafe fn on_ice_candidate_error(
self: &mut PeerConnectionObserverWrapper,
address: String,
port: i32,
@@ -186,32 +187,35 @@ pub mod ffi {
error_code: i32,
error_text: String,
);
fn on_ice_candidates_removed(
unsafe fn on_ice_candidates_removed(
self: &mut PeerConnectionObserverWrapper,
removed: Vec<CandidatePtr>,
);
fn on_ice_connection_receiving_change(
unsafe fn on_ice_connection_receiving_change(
self: &mut PeerConnectionObserverWrapper,
receiving: bool,
);
fn on_ice_selected_candidate_pair_changed(
unsafe fn on_ice_selected_candidate_pair_changed(
self: &mut PeerConnectionObserverWrapper,
event: CandidatePairChangeEvent,
);
fn on_add_track(
unsafe fn on_add_track(
self: &mut PeerConnectionObserverWrapper,
receiver: UniquePtr<RtpReceiver>,
streams: Vec<MediaStreamPtr>,
);
fn on_track(
unsafe fn on_track(
self: &mut PeerConnectionObserverWrapper,
transceiver: UniquePtr<RtpTransceiver>,
);
fn on_remove_track(
unsafe fn on_remove_track(
self: &mut PeerConnectionObserverWrapper,
receiver: UniquePtr<RtpReceiver>,
);
fn on_interesting_usage(self: &mut PeerConnectionObserverWrapper, usage_pattern: i32);
unsafe fn on_interesting_usage(
self: &mut PeerConnectionObserverWrapper,
usage_pattern: i32,
);
}
}
@@ -273,60 +277,63 @@ pub trait PeerConnectionObserver: Send + Sync {
fn on_interesting_usage(&mut self, usage_pattern: i32);
}
// Thread safety is handled inside PeerConnectionObserver
pub struct PeerConnectionObserverWrapper {
observer: Rc<RefCell<dyn PeerConnectionObserver>>,
observer: *mut dyn PeerConnectionObserver,
}
impl PeerConnectionObserverWrapper {
pub fn new(observer: Rc<RefCell<dyn PeerConnectionObserver>>) -> Self {
/// SAFETY
/// PeerConnectionObserver must lives as long as PeerConnectionObserverWrapper does
pub unsafe fn new(observer: *mut dyn PeerConnectionObserver) -> Self {
Self { observer }
}
fn on_signaling_change(&mut self, new_state: ffi::SignalingState) {
self.observer.borrow_mut().on_signaling_change(new_state);
unsafe fn on_signaling_change(&mut self, new_state: ffi::SignalingState) {
(*self.observer).on_signaling_change(new_state);
}
fn on_add_stream(&mut self, stream: UniquePtr<MediaStreamInterface>) {
self.observer.borrow_mut().on_add_stream(stream);
unsafe fn on_add_stream(&mut self, stream: UniquePtr<MediaStreamInterface>) {
(*self.observer).on_add_stream(stream);
}
fn on_remove_stream(&mut self, stream: UniquePtr<MediaStreamInterface>) {
self.observer.borrow_mut().on_remove_stream(stream);
unsafe fn on_remove_stream(&mut self, stream: UniquePtr<MediaStreamInterface>) {
(*self.observer).on_remove_stream(stream);
}
fn on_data_channel(&mut self, data_channel: UniquePtr<DataChannel>) {
self.observer.borrow_mut().on_data_channel(data_channel);
unsafe fn on_data_channel(&mut self, data_channel: UniquePtr<DataChannel>) {
(*self.observer).on_data_channel(data_channel);
}
fn on_renegotiation_needed(&mut self) {
self.observer.borrow_mut().on_renegotiation_needed();
unsafe fn on_renegotiation_needed(&mut self) {
(*self.observer).on_renegotiation_needed();
}
fn on_negotiation_needed_event(&mut self, event: u32) {
self.observer.borrow_mut().on_negotiation_needed_event(event);
unsafe fn on_negotiation_needed_event(&mut self, event: u32) {
(*self.observer).on_negotiation_needed_event(event);
}
fn on_ice_connection_change(&mut self, new_state: ffi::IceConnectionState) {
self.observer.borrow_mut().on_ice_connection_change(new_state);
unsafe fn on_ice_connection_change(&mut self, new_state: ffi::IceConnectionState) {
(*self.observer).on_ice_connection_change(new_state);
}
fn on_standardized_ice_connection_change(&mut self, new_state: ffi::IceConnectionState) {
self.observer.borrow_mut().on_standardized_ice_connection_change(new_state);
unsafe fn on_standardized_ice_connection_change(&mut self, new_state: ffi::IceConnectionState) {
(*self.observer).on_standardized_ice_connection_change(new_state);
}
fn on_connection_change(&mut self, new_state: ffi::PeerConnectionState) {
self.observer.borrow_mut().on_connection_change(new_state);
unsafe fn on_connection_change(&mut self, new_state: ffi::PeerConnectionState) {
(*self.observer).on_connection_change(new_state);
}
fn on_ice_gathering_change(&mut self, new_state: ffi::IceGatheringState) {
self.observer.borrow_mut().on_ice_gathering_change(new_state);
unsafe fn on_ice_gathering_change(&mut self, new_state: ffi::IceGatheringState) {
(*self.observer).on_ice_gathering_change(new_state);
}
fn on_ice_candidate(&mut self, candidate: UniquePtr<IceCandidate>) {
self.observer.borrow_mut().on_ice_candidate(candidate);
unsafe fn on_ice_candidate(&mut self, candidate: UniquePtr<IceCandidate>) {
(*self.observer).on_ice_candidate(candidate);
}
fn on_ice_candidate_error(
unsafe fn on_ice_candidate_error(
&mut self,
address: String,
port: i32,
@@ -334,29 +341,31 @@ impl PeerConnectionObserverWrapper {
error_code: i32,
error_text: String,
) {
self.observer
.borrow_mut().on_ice_candidate_error(address, port, url, error_code, error_text);
(*self.observer).on_ice_candidate_error(address, port, url, error_code, error_text);
}
fn on_ice_candidates_removed(&mut self, removed: Vec<ffi::CandidatePtr>) {
unsafe fn on_ice_candidates_removed(&mut self, removed: Vec<ffi::CandidatePtr>) {
let mut vec = Vec::new();
for v in removed {
vec.push(v.ptr);
}
self.observer.borrow_mut().on_ice_candidates_removed(vec);
(*self.observer).on_ice_candidates_removed(vec);
}
fn on_ice_connection_receiving_change(&mut self, receiving: bool) {
self.observer.borrow_mut().on_ice_connection_receiving_change(receiving);
unsafe fn on_ice_connection_receiving_change(&mut self, receiving: bool) {
(*self.observer).on_ice_connection_receiving_change(receiving);
}
fn on_ice_selected_candidate_pair_changed(&mut self, event: ffi::CandidatePairChangeEvent) {
self.observer.borrow_mut().on_ice_selected_candidate_pair_changed(event);
unsafe fn on_ice_selected_candidate_pair_changed(
&mut self,
event: ffi::CandidatePairChangeEvent,
) {
(*self.observer).on_ice_selected_candidate_pair_changed(event);
}
fn on_add_track(
unsafe fn on_add_track(
&mut self,
receiver: UniquePtr<RtpReceiver>,
streams: Vec<ffi::MediaStreamPtr>,
@@ -367,18 +376,18 @@ impl PeerConnectionObserverWrapper {
vec.push(v.ptr);
}
self.observer.borrow_mut().on_add_track(receiver, vec);
(*self.observer).on_add_track(receiver, vec);
}
fn on_track(&mut self, transceiver: UniquePtr<RtpTransceiver>) {
self.observer.borrow_mut().on_track(transceiver);
unsafe fn on_track(&mut self, transceiver: UniquePtr<RtpTransceiver>) {
(*self.observer).on_track(transceiver);
}
fn on_remove_track(&mut self, receiver: UniquePtr<RtpReceiver>) {
self.observer.borrow_mut().on_remove_track(receiver);
unsafe fn on_remove_track(&mut self, receiver: UniquePtr<RtpReceiver>) {
(*self.observer).on_remove_track(receiver);
}
fn on_interesting_usage(&mut self, usage_pattern: i32) {
self.observer.borrow_mut().on_interesting_usage(usage_pattern);
unsafe fn on_interesting_usage(&mut self, usage_pattern: i32) {
(*self.observer).on_interesting_usage(usage_pattern);
}
}
@@ -57,11 +57,11 @@ namespace livekit{
webrtc::PeerConnectionDependencies deps{observer.get()};
auto result = peer_factory_->CreatePeerConnectionOrError(*config, std::move(deps));
if(!result.ok()){
if(!result.ok()) {
throw std::runtime_error(serialize_error(to_error(result.error())));
}
return std::make_unique<PeerConnection>(std::move(result.value()), std::move(observer));
return std::make_unique<PeerConnection>(result.value(), std::move(observer));
}
std::unique_ptr<PeerConnectionFactory> create_peer_connection_factory() {
@@ -22,12 +22,12 @@ pub mod ffi {
#[derive(Debug, Clone)]
pub struct ICEServer {
urls: Vec<String>,
username: String,
password: String,
pub urls: Vec<String>,
pub username: String,
pub password: String,
}
#[derive(Debug)]
#[derive(Debug, Clone)]
pub struct RTCConfiguration {
pub ice_servers: Vec<ICEServer>,
}
@@ -51,137 +51,3 @@ pub mod ffi {
) -> Result<UniquePtr<PeerConnection>>;
}
}
/*
#[cfg(test)]
mod test {
struct TestObserver {
}
impl PeerConnectionObserver for TestObserver {
fn on_signaling_change(&self, new_state: SignalingState) {
log::debug!("Signaling state changed: {:?}", new_state);
}
fn on_add_stream(&self, stream: UniquePtr<MediaStreamInterface>) {
todo!()
}
fn on_remove_stream(&self, stream: UniquePtr<MediaStreamInterface>) {
todo!()
}
fn on_data_channel(&self, data_channel: UniquePtr<DataChannel>) {
todo!()
}
fn on_renegotiation_needed(&self) {
todo!()
}
fn on_negotiation_needed_event(&self, event: u32) {
todo!()
}
fn on_ice_connection_change(&self, new_state: IceConnectionState) {
log::debug!("ICE connection state changed: {:?}", new_state);
}
fn on_standardized_ice_connection_change(&self, new_state: IceConnectionState) {
todo!()
}
fn on_connection_change(&self, new_state: PeerConnectionState) {
log::debug!("PeerConnection state changed: {:?}", new_state);
}
fn on_ice_gathering_change(&self, new_state: IceGatheringState) {
todo!()
}
fn on_ice_candidate(&self, candidate: UniquePtr<IceCandidate>) {
todo!()
}
fn on_ice_candidate_error(&self, address: String, port: i32, url: String, error_code: i32, error_text: String) {
todo!()
}
fn on_ice_candidates_removed(&self, removed: Vec<UniquePtr<Candidate>>) {
todo!()
}
fn on_ice_connection_receiving_change(&self, receiving: bool) {
todo!()
}
fn on_ice_selected_candidate_pair_changed(&self, event: CandidatePairChangeEvent) {
todo!()
}
fn on_add_track(&self, receiver: UniquePtr<RtpReceiver>, streams: Vec<UniquePtr<MediaStreamInterface>>) {
todo!()
}
fn on_track(&self, transceiver: UniquePtr<RtpTransceiver>) {
todo!()
}
fn on_remove_track(&self, receiver: UniquePtr<RtpReceiver>) {
todo!()
}
fn on_interesting_usage(&self, usage_pattern: i32) {
todo!()
}
}
struct SessionObserver {
}
impl CreateSdpObserver for SessionObserver {
fn on_success(&self, session_description: UniquePtr<crate::jsep::ffi::SessionDescription>) {
info!("on_success");
}
fn on_failure(&self, error: UniquePtr<crate::rtc_error::ffi::RTCError>) {
info!("on_failure");
}
}
#[test]
fn create_pc_test() {
env_logger::init();
let factory = ffi::create_peer_connection_factory(); // Default factory config is defined on the c++ side atm
unsafe {
let mut pc = factory.create_peer_connection(ffi::create_rtc_configuration(ffi::RTCConfiguration {
ice_servers: vec![ffi::ICEServer {
urls: vec!["stun:stun.l.google.com:19302".to_string()],
username: "".to_string(),
password: "".to_string(),
}],
}), peer_connection::ffi::create_native_peer_connection_observer(Box::new(peer_connection::PeerConnectionObserverWrapper::new(Box::new(TestObserver{}))))).unwrap();
let options = peer_connection::ffi::RTCOfferAnswerOptions::default();
let sdp_observer = jsep::ffi::create_native_create_sdp_observer(Box::new(jsep::CreateSdpObserverWrapper::new(Box::new(SessionObserver{}))));
pc.pin_mut().create_offer(sdp_observer, options);
sleep(Duration::from_secs(2));
pc.pin_mut().close();
}
}
}
*/
@@ -15,7 +15,7 @@ namespace livekit {
lk_error.error_detail = static_cast<RTCErrorDetailType>(error.error_detail());
lk_error.error_type = static_cast<RTCErrorType>(error.type());
lk_error.has_sctp_cause_code = error.sctp_cause_code().has_value();
lk_error.sctp_cause_code = error.sctp_cause_code().value();
lk_error.sctp_cause_code = error.sctp_cause_code().value_or(0);
lk_error.message = error.message();
return lk_error;
}
+12 -2
View File
@@ -1,4 +1,14 @@
#[derive(Debug)]
pub struct DataChannel {
use cxx::UniquePtr;
use libwebrtc_sys::data_channel as sys_dc;
pub struct DataChannel {
cxx_handle: UniquePtr<sys_dc::ffi::DataChannel>,
}
impl DataChannel {
pub(crate) fn new(cxx_handle: UniquePtr<sys_dc::ffi::DataChannel>) -> Self {
Self {
cxx_handle,
}
}
}
+11 -8
View File
@@ -2,23 +2,26 @@ use cxx::{SharedPtr, UniquePtr};
use libwebrtc_sys::jsep as sys_jsep;
#[derive(Debug)]
pub struct IceCandidate {
}
pub struct IceCandidate {}
#[derive(Debug)]
pub struct SessionDescription {
cxx_handle: UniquePtr<sys_jsep::ffi::SessionDescription>
cxx_handle: UniquePtr<sys_jsep::ffi::SessionDescription>,
}
impl SessionDescription {
pub(crate) fn new(cxx_handle: UniquePtr<sys_jsep::ffi::SessionDescription>) -> Self {
Self {
cxx_handle
}
Self { cxx_handle }
}
pub(crate) fn release(self) -> UniquePtr<sys_jsep::ffi::SessionDescription>{
pub(crate) fn release(self) -> UniquePtr<sys_jsep::ffi::SessionDescription> {
self.cxx_handle
}
}
impl Clone for SessionDescription {
fn clone(&self) -> Self {
SessionDescription::new(self.cxx_handle.clone())
}
}
+1 -1
View File
@@ -1,8 +1,8 @@
pub mod data_channel;
pub mod jsep;
pub mod media_stream;
pub mod peer_connection;
pub mod peer_connection_factory;
pub mod rtc_error;
pub mod rtp_receiver;
pub mod rtp_transceiver;
pub mod jsep;
+1 -4
View File
@@ -1,5 +1,2 @@
#[derive(Debug)]
pub struct MediaStream {
}
pub struct MediaStream {}
+62 -16
View File
@@ -2,15 +2,16 @@ use cxx::UniquePtr;
use libwebrtc_sys::jsep as sys_jsep;
use libwebrtc_sys::peer_connection as sys_pc;
use std::sync::{Arc, Mutex};
use log::trace;
use thiserror::Error;
use tokio::sync::{mpsc, oneshot};
use crate::data_channel::DataChannel;
use crate::jsep::{IceCandidate, SessionDescription};
use crate::media_stream::MediaStream;
use crate::rtc_error::RTCError;
use crate::rtp_receiver::RtpReceiver;
use crate::rtp_transceiver::RtpTransceiver;
use crate::jsep::{SessionDescription, IceCandidate};
pub use libwebrtc_sys::peer_connection::ffi::IceConnectionState;
pub use libwebrtc_sys::peer_connection::ffi::IceGatheringState;
@@ -28,14 +29,17 @@ pub enum SdpError {
pub struct PeerConnection {
cxx_handle: UniquePtr<sys_pc::ffi::PeerConnection>,
observer: InternalObserver,
observer: Box<InternalObserver>,
}
impl PeerConnection {
pub(crate) fn new(cxx_handle: UniquePtr<sys_pc::ffi::PeerConnection>) -> Self {
pub(crate) fn new(
cxx_handle: UniquePtr<sys_pc::ffi::PeerConnection>,
observer: Box<InternalObserver>,
) -> Self {
Self {
cxx_handle,
observer: InternalObserver::default()
observer,
}
}
@@ -98,8 +102,11 @@ impl PeerConnection {
) -> Result<(), SdpError> {
let (tx, mut rx) = mpsc::channel(1);
let wrapper =
sys_jsep::SetRemoteSdpObserverWrapper::new(Box::new(InternalSetRemoteSdpObserver { tx }));
let native_wrapper = sys_jsep::ffi::create_native_set_remote_sdp_observer(Box::new(wrapper));
sys_jsep::SetRemoteSdpObserverWrapper::new(Box::new(InternalSetRemoteSdpObserver {
tx,
}));
let native_wrapper =
sys_jsep::ffi::create_native_set_remote_sdp_observer(Box::new(wrapper));
self.cxx_handle
.pin_mut()
@@ -111,6 +118,10 @@ impl PeerConnection {
}
}
pub fn close(&mut self) {
self.cxx_handle.pin_mut().close();
}
pub fn on_signaling_change(&mut self, handler: OnSignalingChangeHandler) {
*self.observer.on_signaling_change_handler.lock().unwrap() = Some(handler);
}
@@ -232,7 +243,9 @@ impl sys_jsep::CreateSdpObserver for InternalCreateSdpObserver {
&self,
session_description: UniquePtr<libwebrtc_sys::jsep::ffi::SessionDescription>,
) {
self.tx.blocking_send(Ok(SessionDescription::new(session_description))).unwrap();
self.tx
.blocking_send(Ok(SessionDescription::new(session_description)))
.unwrap();
}
fn on_failure(&self, error: RTCError) {
@@ -346,6 +359,7 @@ impl Default for InternalObserver {
// Observers are being called on the Signaling Thread
impl sys_pc::PeerConnectionObserver for InternalObserver {
fn on_signaling_change(&mut self, new_state: SignalingState) {
trace!("on_signaling_change, {:?}", new_state);
let mut handler = self.on_signaling_change_handler.lock().unwrap();
if let Some(f) = handler.as_mut() {
f(new_state);
@@ -356,6 +370,7 @@ impl sys_pc::PeerConnectionObserver for InternalObserver {
&mut self,
stream: UniquePtr<libwebrtc_sys::media_stream_interface::ffi::MediaStreamInterface>,
) {
trace!("on_add_stream");
let mut handler = self.on_add_stream_handler.lock().unwrap();
if let Some(f) = handler.as_mut() {
// TODO(theomonnom)
@@ -366,6 +381,7 @@ impl sys_pc::PeerConnectionObserver for InternalObserver {
&mut self,
stream: UniquePtr<libwebrtc_sys::media_stream_interface::ffi::MediaStreamInterface>,
) {
trace!("on_remove_stream");
let mut handler = self.on_remove_stream_handler.lock().unwrap();
if let Some(f) = handler.as_mut() {
// TODO(theomonnom)
@@ -376,6 +392,7 @@ impl sys_pc::PeerConnectionObserver for InternalObserver {
&mut self,
data_channel: UniquePtr<libwebrtc_sys::data_channel::ffi::DataChannel>,
) {
trace!("on_data_channel");
let mut handler = self.on_data_channel_handler.lock().unwrap();
if let Some(f) = handler.as_mut() {
// TODO(theomonnom)
@@ -383,6 +400,7 @@ impl sys_pc::PeerConnectionObserver for InternalObserver {
}
fn on_renegotiation_needed(&mut self) {
trace!("on_renegotiation_needed");
let mut handler = self.on_renegotiation_needed_handler.lock().unwrap();
if let Some(f) = handler.as_mut() {
f();
@@ -390,6 +408,7 @@ impl sys_pc::PeerConnectionObserver for InternalObserver {
}
fn on_negotiation_needed_event(&mut self, event: u32) {
trace!("on_negotiation_needed_event");
let mut handler = self.on_negotiation_needed_event_handler.lock().unwrap();
if let Some(f) = handler.as_mut() {
f(event);
@@ -397,6 +416,7 @@ impl sys_pc::PeerConnectionObserver for InternalObserver {
}
fn on_ice_connection_change(&mut self, new_state: IceConnectionState) {
trace!("on_ice_connection_change");
let mut handler = self.on_ice_connection_change_handler.lock().unwrap();
if let Some(f) = handler.as_mut() {
f(new_state);
@@ -404,6 +424,7 @@ impl sys_pc::PeerConnectionObserver for InternalObserver {
}
fn on_standardized_ice_connection_change(&mut self, new_state: IceConnectionState) {
trace!("on_standardized_ice_connection_change");
let mut handler = self
.on_standardized_ice_connection_change_handler
.lock()
@@ -414,6 +435,7 @@ impl sys_pc::PeerConnectionObserver for InternalObserver {
}
fn on_connection_change(&mut self, new_state: PeerConnectionState) {
trace!("on_connection_change");
let mut handler = self.on_connection_change_handler.lock().unwrap();
if let Some(f) = handler.as_mut() {
f(new_state);
@@ -421,6 +443,7 @@ impl sys_pc::PeerConnectionObserver for InternalObserver {
}
fn on_ice_gathering_change(&mut self, new_state: IceGatheringState) {
trace!("on_ice_gathering_change");
let mut handler = self.on_ice_gathering_change_handler.lock().unwrap();
if let Some(f) = handler.as_mut() {
f(new_state);
@@ -428,6 +451,7 @@ impl sys_pc::PeerConnectionObserver for InternalObserver {
}
fn on_ice_candidate(&mut self, candidate: UniquePtr<libwebrtc_sys::jsep::ffi::IceCandidate>) {
trace!("on_ice_candidate");
let mut handler = self.on_ice_candidate_handler.lock().unwrap();
if let Some(f) = handler.as_mut() {
// TODO(theomonnom)
@@ -442,6 +466,7 @@ impl sys_pc::PeerConnectionObserver for InternalObserver {
error_code: i32,
error_text: String,
) {
trace!("on_ice_candidate_error");
let mut handler = self.on_ice_candidate_error_handler.lock().unwrap();
if let Some(f) = handler.as_mut() {
f(address, port, url, error_code, error_text);
@@ -452,6 +477,7 @@ impl sys_pc::PeerConnectionObserver for InternalObserver {
&mut self,
removed: Vec<UniquePtr<libwebrtc_sys::candidate::ffi::Candidate>>,
) {
trace!("on_ice_candidates_removed");
let mut handler = self.on_ice_candidates_removed_handler.lock().unwrap();
if let Some(f) = handler.as_mut() {
// TODO(theomonnom)
@@ -459,6 +485,7 @@ impl sys_pc::PeerConnectionObserver for InternalObserver {
}
fn on_ice_connection_receiving_change(&mut self, receiving: bool) {
trace!("on_ice_connection_receiving_change");
let mut handler = self
.on_ice_connection_receiving_change_handler
.lock()
@@ -472,6 +499,7 @@ impl sys_pc::PeerConnectionObserver for InternalObserver {
&mut self,
event: libwebrtc_sys::peer_connection::ffi::CandidatePairChangeEvent,
) {
trace!("on_ice_selected_candidate_pair_changed");
let mut handler = self
.on_ice_selected_candidate_pair_changed_handler
.lock()
@@ -486,6 +514,7 @@ impl sys_pc::PeerConnectionObserver for InternalObserver {
receiver: UniquePtr<libwebrtc_sys::rtp_receiver::ffi::RtpReceiver>,
streams: Vec<UniquePtr<libwebrtc_sys::media_stream_interface::ffi::MediaStreamInterface>>,
) {
trace!("on_add_track");
let mut handler = self.on_add_track_handler.lock().unwrap();
if let Some(f) = handler.as_mut() {
// TODO(theomonnom)
@@ -496,6 +525,7 @@ impl sys_pc::PeerConnectionObserver for InternalObserver {
&mut self,
transceiver: UniquePtr<libwebrtc_sys::rtp_transceiver::ffi::RtpTransceiver>,
) {
trace!("on_track");
let mut handler = self.on_track_handler.lock().unwrap();
if let Some(f) = handler.as_mut() {
// TODO(theomonnom)
@@ -506,6 +536,7 @@ impl sys_pc::PeerConnectionObserver for InternalObserver {
&mut self,
receiver: UniquePtr<libwebrtc_sys::rtp_receiver::ffi::RtpReceiver>,
) {
trace!("on_remove_track");
let mut handler = self.on_remove_track_handler.lock().unwrap();
if let Some(f) = handler.as_mut() {
// TODO(theomonnom)
@@ -513,6 +544,7 @@ impl sys_pc::PeerConnectionObserver for InternalObserver {
}
fn on_interesting_usage(&mut self, usage_pattern: i32) {
trace!("on_interesting_usage");
let mut handler = self.on_interesting_usage_handler.lock().unwrap();
if let Some(f) = handler.as_mut() {
f(usage_pattern);
@@ -522,19 +554,33 @@ impl sys_pc::PeerConnectionObserver for InternalObserver {
#[cfg(test)]
mod tests {
use libwebrtc_sys::peer_connection_factory::ffi::RTCConfiguration;
use crate::peer_connection_factory::PeerConnectionFactory;
use libwebrtc_sys::peer_connection_factory::ffi::{RTCConfiguration};
fn init_log() {
let _ = env_logger::builder().is_test(true).try_init();
}
#[tokio::test]
async fn create_pc() {
let factory = PeerConnectionFactory::new();
let mut pc = factory.create_peer_connection(
RTCConfiguration {
ice_servers: vec![],
},
Box::new(()),
).unwrap();
init_log();
let offer = pc.create_offer().await.unwrap();
let factory = PeerConnectionFactory::new();
let config = RTCConfiguration {
ice_servers: vec![],
};
let mut bob = factory.create_peer_connection(config.clone()).unwrap();
let mut alice = factory.create_peer_connection(config.clone()).unwrap();
let offer = bob.create_offer().await.unwrap();
bob.set_local_description(offer.clone()).await.unwrap();
alice.set_remote_description(offer).await.unwrap();
let answer = alice.create_answer().await.unwrap();
alice.set_local_description(answer.clone()).await.unwrap();
bob.set_remote_description(answer).await.unwrap();
alice.close();
bob.close();
}
}
@@ -2,7 +2,7 @@ use cxx::UniquePtr;
use libwebrtc_sys::peer_connection as sys_pc;
use libwebrtc_sys::peer_connection_factory as sys_factory;
use crate::peer_connection::PeerConnection;
use crate::peer_connection::{InternalObserver, PeerConnection};
use crate::rtc_error::RTCError;
pub use sys_factory::ffi::{ICEServer, RTCConfiguration};
@@ -21,22 +21,23 @@ impl PeerConnectionFactory {
pub fn create_peer_connection(
&self,
config: RTCConfiguration,
observer: Box<dyn sys_pc::PeerConnectionObserver>,
) -> Result<PeerConnection, RTCError> {
let native_config = sys_factory::ffi::create_rtc_configuration(config);
let native_observer = sys_pc::ffi::create_native_peer_connection_observer(Box::new(
sys_pc::PeerConnectionObserverWrapper::new(observer),
));
let pc_result: Result<UniquePtr<sys_pc::ffi::PeerConnection>, cxx::Exception> = unsafe {
self.cxx_handle
.create_peer_connection(native_config, native_observer)
};
unsafe {
let mut observer = Box::new(InternalObserver::default());
let observer_wrapper = sys_pc::PeerConnectionObserverWrapper::new(&mut *observer);
let native_observer =
sys_pc::ffi::create_native_peer_connection_observer(Box::new(observer_wrapper));
let res = self
.cxx_handle
.create_peer_connection(native_config, native_observer);
match pc_result {
Ok(cxx_handle) => Ok(PeerConnection::new(cxx_handle)),
Err(e) => {
Err(unsafe { RTCError::from(e.what()) }) // TODO
match res {
Ok(cxx_handle) => Ok(PeerConnection::new(cxx_handle, observer)),
Err(e) => {
Err(RTCError::from(e.what())) // TODO
}
}
}
}
+1 -4
View File
@@ -1,5 +1,2 @@
#[derive(Debug)]
pub struct RtpReceiver {
}
pub struct RtpReceiver {}
+1 -4
View File
@@ -1,5 +1,2 @@
#[derive(Debug)]
pub struct RtpTransceiver {
}
pub struct RtpTransceiver {}