diff --git a/crates/livekit-webrtc/libwebrtc-sys/include/livekit/data_channel.h b/crates/livekit-webrtc/libwebrtc-sys/include/livekit/data_channel.h index 0c50656..cc2618c 100644 --- a/crates/livekit-webrtc/libwebrtc-sys/include/livekit/data_channel.h +++ b/crates/livekit-webrtc/libwebrtc-sys/include/livekit/data_channel.h @@ -7,8 +7,10 @@ #include #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 data_channel_; }; + std::unique_ptr create_data_channel_init(DataChannelInit init); + static std::unique_ptr _unique_data_channel(){ return nullptr; // Ignore } + } // livekit #endif //CLIENT_SDK_NATIVE_DATA_CHANNEL_H diff --git a/crates/livekit-webrtc/libwebrtc-sys/include/livekit/jsep.h b/crates/livekit-webrtc/libwebrtc-sys/include/livekit/jsep.h index bc08972..ab75789 100644 --- a/crates/livekit-webrtc/libwebrtc-sys/include/livekit/jsep.h +++ b/crates/livekit-webrtc/libwebrtc-sys/include/livekit/jsep.h @@ -40,10 +40,6 @@ namespace livekit { return nullptr; // Ignore } - static std::shared_ptr _shared_session_description(){ - return nullptr; // Ignore - } - // SetCreateSdpObserver class NativeCreateSdpObserver : public webrtc::CreateSessionDescriptionObserver { diff --git a/crates/livekit-webrtc/libwebrtc-sys/include/livekit/peer_connection.h b/crates/livekit-webrtc/libwebrtc-sys/include/livekit/peer_connection.h index ed79f0d..af86759 100644 --- a/crates/livekit-webrtc/libwebrtc-sys/include/livekit/peer_connection.h +++ b/crates/livekit-webrtc/libwebrtc-sys/include/livekit/peer_connection.h @@ -19,11 +19,12 @@ namespace livekit { public: explicit PeerConnection(rtc::scoped_refptr peer_connection, std::unique_ptr observer); - void close(); void create_offer(std::unique_ptr observer, RTCOfferAnswerOptions options); void create_answer(std::unique_ptr observer, RTCOfferAnswerOptions options); void set_local_description(std::unique_ptr desc, std::unique_ptr observer); void set_remote_description(std::unique_ptr desc, std::unique_ptr observer); + std::unique_ptr create_data_channel(rust::String label, std::unique_ptr init); + void close(); private: rtc::scoped_refptr peer_connection_; diff --git a/crates/livekit-webrtc/libwebrtc-sys/include/livekit/rust_types.h b/crates/livekit-webrtc/libwebrtc-sys/include/livekit/rust_types.h index eee4ef9..d6ddf74 100644 --- a/crates/livekit-webrtc/libwebrtc-sys/include/livekit/rust_types.h +++ b/crates/livekit-webrtc/libwebrtc-sys/include/livekit/rust_types.h @@ -18,6 +18,7 @@ namespace livekit { // Shared types struct RTCOfferAnswerOptions; struct RTCError; + struct DataChannelInit; } #endif //RUST_TYPES_H diff --git a/crates/livekit-webrtc/libwebrtc-sys/src/data_channel.cpp b/crates/livekit-webrtc/libwebrtc-sys/src/data_channel.cpp index ba6ffb8..94fbe72 100644 --- a/crates/livekit-webrtc/libwebrtc-sys/src/data_channel.cpp +++ b/crates/livekit-webrtc/libwebrtc-sys/src/data_channel.cpp @@ -2,11 +2,35 @@ // Created by Théo Monnom on 01/09/2022. // +#include + #include "livekit/data_channel.h" +#include "libwebrtc-sys/src/data_channel.rs.h" namespace livekit { - DataChannel::DataChannel(rtc::scoped_refptr data_channel) : data_channel_(data_channel) { + DataChannel::DataChannel(rtc::scoped_refptr data_channel) : data_channel_(std::move(data_channel)) { } + + std::unique_ptr create_data_channel_init(DataChannelInit init) { + auto rtc_init = std::make_unique(); + 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(init.priority); + + return rtc_init; + } + } // livekit \ No newline at end of file diff --git a/crates/livekit-webrtc/libwebrtc-sys/src/data_channel.rs b/crates/livekit-webrtc/libwebrtc-sys/src/data_channel.rs index a349944..dc22f0c 100644 --- a/crates/livekit-webrtc/libwebrtc-sys/src/data_channel.rs +++ b/crates/livekit-webrtc/libwebrtc-sys/src/data_channel.rs @@ -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; fn _unique_data_channel() -> UniquePtr; // Ignore } diff --git a/crates/livekit-webrtc/libwebrtc-sys/src/jsep.rs b/crates/livekit-webrtc/libwebrtc-sys/src/jsep.rs index fc8299a..e691f7a 100644 --- a/crates/livekit-webrtc/libwebrtc-sys/src/jsep.rs +++ b/crates/livekit-webrtc/libwebrtc-sys/src/jsep.rs @@ -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; fn create_native_create_sdp_observer( observer: Box, @@ -45,7 +46,6 @@ pub mod ffi { ) -> UniquePtr; fn _unique_ice_candidate() -> UniquePtr; // Ignore - fn _shared_session_description() -> SharedPtr; // Ignore fn _unique_session_description() -> UniquePtr; // Ignore } } diff --git a/crates/livekit-webrtc/libwebrtc-sys/src/peer_connection.cpp b/crates/livekit-webrtc/libwebrtc-sys/src/peer_connection.cpp index b762444..6ebbf59 100644 --- a/crates/livekit-webrtc/libwebrtc-sys/src/peer_connection.cpp +++ b/crates/livekit-webrtc/libwebrtc-sys/src/peer_connection.cpp @@ -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 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 PeerConnection::create_data_channel(rust::String label, std::unique_ptr 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(result.value()); + } + + void PeerConnection::close() { + peer_connection_->Close(); + } + /* Observer */ NativePeerConnectionObserver::NativePeerConnectionObserver(rust::Box observer) : observer_(std::move(observer)) { diff --git a/crates/livekit-webrtc/libwebrtc-sys/src/peer_connection.rs b/crates/livekit-webrtc/libwebrtc-sys/src/peer_connection.rs index d028187..4d66dae 100644 --- a/crates/livekit-webrtc/libwebrtc-sys/src/peer_connection.rs +++ b/crates/livekit-webrtc/libwebrtc-sys/src/peer_connection.rs @@ -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, ); - fn on_remove_stream( + unsafe fn on_remove_stream( self: &mut PeerConnectionObserverWrapper, stream: UniquePtr, ); - fn on_data_channel( + unsafe fn on_data_channel( self: &mut PeerConnectionObserverWrapper, data_channel: UniquePtr, ); - 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, ); - 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, ); - 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, streams: Vec, ); - fn on_track( + unsafe fn on_track( self: &mut PeerConnectionObserverWrapper, transceiver: UniquePtr, ); - fn on_remove_track( + unsafe fn on_remove_track( self: &mut PeerConnectionObserverWrapper, receiver: UniquePtr, ); - 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>, + observer: *mut dyn PeerConnectionObserver, } impl PeerConnectionObserverWrapper { - pub fn new(observer: Rc>) -> 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) { - self.observer.borrow_mut().on_add_stream(stream); + unsafe fn on_add_stream(&mut self, stream: UniquePtr) { + (*self.observer).on_add_stream(stream); } - fn on_remove_stream(&mut self, stream: UniquePtr) { - self.observer.borrow_mut().on_remove_stream(stream); + unsafe fn on_remove_stream(&mut self, stream: UniquePtr) { + (*self.observer).on_remove_stream(stream); } - fn on_data_channel(&mut self, data_channel: UniquePtr) { - self.observer.borrow_mut().on_data_channel(data_channel); + unsafe fn on_data_channel(&mut self, data_channel: UniquePtr) { + (*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) { - self.observer.borrow_mut().on_ice_candidate(candidate); + unsafe fn on_ice_candidate(&mut self, candidate: UniquePtr) { + (*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) { + unsafe fn on_ice_candidates_removed(&mut self, removed: Vec) { 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, streams: Vec, @@ -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) { - self.observer.borrow_mut().on_track(transceiver); + unsafe fn on_track(&mut self, transceiver: UniquePtr) { + (*self.observer).on_track(transceiver); } - fn on_remove_track(&mut self, receiver: UniquePtr) { - self.observer.borrow_mut().on_remove_track(receiver); + unsafe fn on_remove_track(&mut self, receiver: UniquePtr) { + (*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); } } diff --git a/crates/livekit-webrtc/libwebrtc-sys/src/peer_connection_factory.cpp b/crates/livekit-webrtc/libwebrtc-sys/src/peer_connection_factory.cpp index a3eacd0..ee31583 100644 --- a/crates/livekit-webrtc/libwebrtc-sys/src/peer_connection_factory.cpp +++ b/crates/livekit-webrtc/libwebrtc-sys/src/peer_connection_factory.cpp @@ -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(std::move(result.value()), std::move(observer)); + return std::make_unique(result.value(), std::move(observer)); } std::unique_ptr create_peer_connection_factory() { diff --git a/crates/livekit-webrtc/libwebrtc-sys/src/peer_connection_factory.rs b/crates/livekit-webrtc/libwebrtc-sys/src/peer_connection_factory.rs index 290ac3a..df964b1 100644 --- a/crates/livekit-webrtc/libwebrtc-sys/src/peer_connection_factory.rs +++ b/crates/livekit-webrtc/libwebrtc-sys/src/peer_connection_factory.rs @@ -22,12 +22,12 @@ pub mod ffi { #[derive(Debug, Clone)] pub struct ICEServer { - urls: Vec, - username: String, - password: String, + pub urls: Vec, + pub username: String, + pub password: String, } - #[derive(Debug)] + #[derive(Debug, Clone)] pub struct RTCConfiguration { pub ice_servers: Vec, } @@ -51,137 +51,3 @@ pub mod ffi { ) -> Result>; } } - -/* - - - - - - -#[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) { - todo!() - } - - fn on_remove_stream(&self, stream: UniquePtr) { - todo!() - } - - fn on_data_channel(&self, data_channel: UniquePtr) { - 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) { - 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>) { - 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, streams: Vec>) { - todo!() - } - - fn on_track(&self, transceiver: UniquePtr) { - todo!() - } - - fn on_remove_track(&self, receiver: UniquePtr) { - todo!() - } - - fn on_interesting_usage(&self, usage_pattern: i32) { - todo!() - } - } - - struct SessionObserver { - - } - - impl CreateSdpObserver for SessionObserver { - fn on_success(&self, session_description: UniquePtr) { - info!("on_success"); - } - - fn on_failure(&self, error: UniquePtr) { - 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(); - } - } -} - -*/ diff --git a/crates/livekit-webrtc/libwebrtc-sys/src/rtc_error.cpp b/crates/livekit-webrtc/libwebrtc-sys/src/rtc_error.cpp index f9eb1bc..5dd1d05 100644 --- a/crates/livekit-webrtc/libwebrtc-sys/src/rtc_error.cpp +++ b/crates/livekit-webrtc/libwebrtc-sys/src/rtc_error.cpp @@ -15,7 +15,7 @@ namespace livekit { lk_error.error_detail = static_cast(error.error_detail()); lk_error.error_type = static_cast(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; } diff --git a/crates/livekit-webrtc/src/data_channel.rs b/crates/livekit-webrtc/src/data_channel.rs index 93fcf85..8a0b3df 100644 --- a/crates/livekit-webrtc/src/data_channel.rs +++ b/crates/livekit-webrtc/src/data_channel.rs @@ -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, } + +impl DataChannel { + pub(crate) fn new(cxx_handle: UniquePtr) -> Self { + Self { + cxx_handle, + } + } +} \ No newline at end of file diff --git a/crates/livekit-webrtc/src/jsep.rs b/crates/livekit-webrtc/src/jsep.rs index 3c2c460..51e8e35 100644 --- a/crates/livekit-webrtc/src/jsep.rs +++ b/crates/livekit-webrtc/src/jsep.rs @@ -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 + cxx_handle: UniquePtr, } impl SessionDescription { pub(crate) fn new(cxx_handle: UniquePtr) -> Self { - Self { - cxx_handle - } + Self { cxx_handle } } - pub(crate) fn release(self) -> UniquePtr{ + pub(crate) fn release(self) -> UniquePtr { self.cxx_handle } +} + + +impl Clone for SessionDescription { + fn clone(&self) -> Self { + SessionDescription::new(self.cxx_handle.clone()) + } } \ No newline at end of file diff --git a/crates/livekit-webrtc/src/lib.rs b/crates/livekit-webrtc/src/lib.rs index e99192a..c4403af 100644 --- a/crates/livekit-webrtc/src/lib.rs +++ b/crates/livekit-webrtc/src/lib.rs @@ -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; diff --git a/crates/livekit-webrtc/src/media_stream.rs b/crates/livekit-webrtc/src/media_stream.rs index 4f62fac..9a693e7 100644 --- a/crates/livekit-webrtc/src/media_stream.rs +++ b/crates/livekit-webrtc/src/media_stream.rs @@ -1,5 +1,2 @@ - #[derive(Debug)] -pub struct MediaStream { - -} +pub struct MediaStream {} diff --git a/crates/livekit-webrtc/src/peer_connection.rs b/crates/livekit-webrtc/src/peer_connection.rs index fee7c0a..b1e48e9 100644 --- a/crates/livekit-webrtc/src/peer_connection.rs +++ b/crates/livekit-webrtc/src/peer_connection.rs @@ -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, - observer: InternalObserver, + observer: Box, } impl PeerConnection { - pub(crate) fn new(cxx_handle: UniquePtr) -> Self { + pub(crate) fn new( + cxx_handle: UniquePtr, + observer: Box, + ) -> 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, ) { - 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, ) { + 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, ) { + 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, ) { + 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) { + 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>, ) { + 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, streams: Vec>, ) { + 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, ) { + 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, ) { + 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(); } } diff --git a/crates/livekit-webrtc/src/peer_connection_factory.rs b/crates/livekit-webrtc/src/peer_connection_factory.rs index 7ac5504..ed17686 100644 --- a/crates/livekit-webrtc/src/peer_connection_factory.rs +++ b/crates/livekit-webrtc/src/peer_connection_factory.rs @@ -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, ) -> Result { 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, 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 + } } } } diff --git a/crates/livekit-webrtc/src/rtp_receiver.rs b/crates/livekit-webrtc/src/rtp_receiver.rs index e15433d..d7f431b 100644 --- a/crates/livekit-webrtc/src/rtp_receiver.rs +++ b/crates/livekit-webrtc/src/rtp_receiver.rs @@ -1,5 +1,2 @@ - #[derive(Debug)] -pub struct RtpReceiver { - -} +pub struct RtpReceiver {} diff --git a/crates/livekit-webrtc/src/rtp_transceiver.rs b/crates/livekit-webrtc/src/rtp_transceiver.rs index 96783c0..bb7598c 100644 --- a/crates/livekit-webrtc/src/rtp_transceiver.rs +++ b/crates/livekit-webrtc/src/rtp_transceiver.rs @@ -1,5 +1,2 @@ - #[derive(Debug)] -pub struct RtpTransceiver { - -} +pub struct RtpTransceiver {}