From 8380ffadefcf280664f00e362222e940df6469b2 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Th=C3=A9o=20Monnom?= Date: Mon, 19 Sep 2022 01:11:18 +0200 Subject: [PATCH] DataChannel progress --- crates/livekit-webrtc/libwebrtc-sys/build.rs | 2 + .../libwebrtc-sys/include/livekit/jsep.h | 2 + .../include/livekit/peer_connection.h | 14 +- .../include/livekit/rust_types.h | 1 + .../libwebrtc-sys/include/livekit/webrtc.h | 32 ++++ .../libwebrtc-sys/src/data_channel.rs | 34 ++-- .../livekit-webrtc/libwebrtc-sys/src/jsep.cpp | 4 + .../livekit-webrtc/libwebrtc-sys/src/jsep.rs | 8 + .../livekit-webrtc/libwebrtc-sys/src/lib.rs | 1 + .../libwebrtc-sys/src/peer_connection.cpp | 22 ++- .../libwebrtc-sys/src/peer_connection.rs | 157 +++++++++++------- .../src/peer_connection_factory.cpp | 1 - .../libwebrtc-sys/src/rtc_error.cpp | 1 - .../libwebrtc-sys/src/webrtc.cpp | 22 +++ .../libwebrtc-sys/src/webrtc.rs | 12 ++ crates/livekit-webrtc/src/data_channel.rs | 48 ++++++ crates/livekit-webrtc/src/jsep.rs | 14 +- crates/livekit-webrtc/src/lib.rs | 1 + crates/livekit-webrtc/src/peer_connection.rs | 82 ++++++++- .../src/peer_connection_factory.rs | 4 +- crates/livekit-webrtc/src/webrtc.rs | 14 ++ 21 files changed, 389 insertions(+), 87 deletions(-) create mode 100644 crates/livekit-webrtc/libwebrtc-sys/include/livekit/webrtc.h create mode 100644 crates/livekit-webrtc/libwebrtc-sys/src/webrtc.cpp create mode 100644 crates/livekit-webrtc/libwebrtc-sys/src/webrtc.rs create mode 100644 crates/livekit-webrtc/src/webrtc.rs diff --git a/crates/livekit-webrtc/libwebrtc-sys/build.rs b/crates/livekit-webrtc/libwebrtc-sys/build.rs index 536833e..909f297 100644 --- a/crates/livekit-webrtc/libwebrtc-sys/build.rs +++ b/crates/livekit-webrtc/libwebrtc-sys/build.rs @@ -78,6 +78,7 @@ fn main() { "src/rtp_receiver.rs", "src/rtp_transceiver.rs", "src/rtc_error.rs", + "src/webrtc.rs", ]); builder.file("src/peer_connection.cpp"); @@ -89,6 +90,7 @@ fn main() { builder.file("src/rtp_receiver.cpp"); builder.file("src/rtp_transceiver.cpp"); builder.file("src/rtc_error.cpp"); + builder.file("src/webrtc.cpp"); for include in includes { builder.include(include); diff --git a/crates/livekit-webrtc/libwebrtc-sys/include/livekit/jsep.h b/crates/livekit-webrtc/libwebrtc-sys/include/livekit/jsep.h index ab75789..8ca39ce 100644 --- a/crates/livekit-webrtc/libwebrtc-sys/include/livekit/jsep.h +++ b/crates/livekit-webrtc/libwebrtc-sys/include/livekit/jsep.h @@ -16,6 +16,8 @@ namespace livekit { class IceCandidate { public: explicit IceCandidate(std::unique_ptr ice_candidate); + + std::unique_ptr release(); private: std::unique_ptr ice_candidate_; }; 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 91abb27..ea47b61 100644 --- a/crates/livekit-webrtc/libwebrtc-sys/include/livekit/peer_connection.h +++ b/crates/livekit-webrtc/libwebrtc-sys/include/livekit/peer_connection.h @@ -13,7 +13,7 @@ #include "rust_types.h" namespace livekit { - class NativePeerConnectionObserver; + class NativeAddIceCandidateObserver; class PeerConnection { public: @@ -24,6 +24,7 @@ namespace livekit { void set_local_description(std::unique_ptr desc, NativeSetLocalSdpObserverHandle &observer); void set_remote_description(std::unique_ptr desc, NativeSetRemoteSdpObserverHandle &observer); std::unique_ptr create_data_channel(rust::String label, std::unique_ptr init); + void add_ice_candidate(std::unique_ptr candidate, NativeAddIceCandidateObserver &observer); void close(); private: @@ -34,6 +35,17 @@ namespace livekit { return nullptr; // Ignore } + class NativeAddIceCandidateObserver { + public: + explicit NativeAddIceCandidateObserver(rust::Box observer); + + void OnComplete(const RTCError &error); + private: + rust::Box observer_; + }; + + std::unique_ptr create_native_add_ice_candidate_observer(rust::Box observer); + class NativePeerConnectionObserver : public webrtc::PeerConnectionObserver { public: explicit NativePeerConnectionObserver(rust::Box observer); 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 e03438a..d9334bf 100644 --- a/crates/livekit-webrtc/libwebrtc-sys/include/livekit/rust_types.h +++ b/crates/livekit-webrtc/libwebrtc-sys/include/livekit/rust_types.h @@ -15,6 +15,7 @@ namespace livekit { struct SetLocalSdpObserverWrapper; struct SetRemoteSdpObserverWrapper; struct DataChannelObserverWrapper; + struct AddIceCandidateObserverWrapper; // Shared types struct RTCOfferAnswerOptions; diff --git a/crates/livekit-webrtc/libwebrtc-sys/include/livekit/webrtc.h b/crates/livekit-webrtc/libwebrtc-sys/include/livekit/webrtc.h new file mode 100644 index 0000000..2dc3120 --- /dev/null +++ b/crates/livekit-webrtc/libwebrtc-sys/include/livekit/webrtc.h @@ -0,0 +1,32 @@ +// +// Created by theom on 18/09/2022. +// + +#ifndef LIVEKIT_WEBRTC_WEBRTC_H +#define LIVEKIT_WEBRTC_WEBRTC_H + +#include "rtc_base/ssl_adapter.h" +#include "rtc_base/physical_socket_server.h" + +#ifdef WEBRTC_WIN +#include "rtc_base/win32_socket_init.h" +#endif + +namespace livekit { + + class RTCRuntime { + public: + RTCRuntime(); + ~RTCRuntime(); + + RTCRuntime(const RTCRuntime&) = delete; + RTCRuntime& operator=(const RTCRuntime&) = delete; + private: + rtc::WinsockInitializer winsock_; + }; + + std::unique_ptr create_rtc_runtime(); + +} // livekit + +#endif //LIVEKIT_WEBRTC_WEBRTC_H diff --git a/crates/livekit-webrtc/libwebrtc-sys/src/data_channel.rs b/crates/livekit-webrtc/libwebrtc-sys/src/data_channel.rs index 4b3ac48..5b0fa35 100644 --- a/crates/livekit-webrtc/libwebrtc-sys/src/data_channel.rs +++ b/crates/livekit-webrtc/libwebrtc-sys/src/data_channel.rs @@ -6,7 +6,7 @@ pub mod ffi { #[derive(Debug)] #[repr(u32)] - enum Priority { + pub enum Priority { VeryLow, Low, Medium, @@ -14,20 +14,20 @@ pub mod ffi { } #[derive(Debug)] - #[allow(deprecated)] pub struct DataChannelInit { + #[allow(deprecated)] #[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, + pub reliable: bool, + pub ordered: bool, + pub has_max_retransmit_time: bool, + pub max_retransmit_time: i32, + pub has_max_retransmits: bool, + pub max_retransmits: i32, + pub protocol: String, + pub negotiated: bool, + pub id: i32, + pub has_priority: bool, + pub priority: Priority, } #[derive(Debug)] @@ -60,6 +60,14 @@ pub mod ffi { type NativeDataChannelInit; type NativeDataChannelObserver; + /// SAFETY + /// The observer must live as the datachannel uses it + unsafe fn register_observer( + self: Pin<&mut DataChannel>, + observer: Pin<&mut NativeDataChannelObserver>, + ); + + fn unregister_observer(self: Pin<&mut DataChannel>); fn close(self: Pin<&mut DataChannel>); fn create_data_channel_init(init: DataChannelInit) -> UniquePtr; diff --git a/crates/livekit-webrtc/libwebrtc-sys/src/jsep.cpp b/crates/livekit-webrtc/libwebrtc-sys/src/jsep.cpp index 4617325..b607eb4 100644 --- a/crates/livekit-webrtc/libwebrtc-sys/src/jsep.cpp +++ b/crates/livekit-webrtc/libwebrtc-sys/src/jsep.cpp @@ -15,6 +15,10 @@ namespace livekit { } + std::unique_ptr IceCandidate::release() { + return std::move(ice_candidate_); + } + SessionDescription::SessionDescription(std::unique_ptr session_description) : session_description_(std::move(session_description)){ } diff --git a/crates/livekit-webrtc/libwebrtc-sys/src/jsep.rs b/crates/livekit-webrtc/libwebrtc-sys/src/jsep.rs index e691f7a..1388e8c 100644 --- a/crates/livekit-webrtc/libwebrtc-sys/src/jsep.rs +++ b/crates/livekit-webrtc/libwebrtc-sys/src/jsep.rs @@ -58,6 +58,14 @@ impl Debug for ffi::SessionDescription { unsafe impl Send for ffi::SessionDescription {} +impl Debug for ffi::IceCandidate { + fn fmt(&self, f: &mut Formatter) -> std::fmt::Result { + write!(f, "TODO") // TODO(theomonnom) + } +} + +unsafe impl Send for ffi::IceCandidate {} + // CreateSdpObserver pub trait CreateSdpObserver: Send { diff --git a/crates/livekit-webrtc/libwebrtc-sys/src/lib.rs b/crates/livekit-webrtc/libwebrtc-sys/src/lib.rs index db03eaa..9959239 100644 --- a/crates/livekit-webrtc/libwebrtc-sys/src/lib.rs +++ b/crates/livekit-webrtc/libwebrtc-sys/src/lib.rs @@ -7,3 +7,4 @@ pub mod peer_connection_factory; pub mod rtc_error; pub mod rtp_receiver; pub mod rtp_transceiver; +pub mod webrtc; diff --git a/crates/livekit-webrtc/libwebrtc-sys/src/peer_connection.cpp b/crates/livekit-webrtc/libwebrtc-sys/src/peer_connection.cpp index a7a0bf3..732f52b 100644 --- a/crates/livekit-webrtc/libwebrtc-sys/src/peer_connection.cpp +++ b/crates/livekit-webrtc/libwebrtc-sys/src/peer_connection.cpp @@ -51,11 +51,31 @@ namespace livekit { return std::make_unique(result.value()); } + void PeerConnection::add_ice_candidate(std::unique_ptr candidate, NativeAddIceCandidateObserver &observer){ + peer_connection_->AddIceCandidate(candidate->release(), [&](const webrtc::RTCError& err){ + observer.OnComplete(to_error(err)); + }); + } + void PeerConnection::close() { peer_connection_->Close(); } - /* Observer */ + // AddIceCandidateObserver + + NativeAddIceCandidateObserver::NativeAddIceCandidateObserver(rust::Box observer) : observer_(std::move(observer)) { + + } + + void NativeAddIceCandidateObserver::OnComplete(const RTCError &error) { + observer_->on_complete(error); + } + + std::unique_ptr create_native_add_ice_candidate_observer(rust::Box observer) { + return std::make_unique(std::move(observer)); + } + + // PeerConnectionObserver 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 b39e2e0..ae6516a 100644 --- a/crates/livekit-webrtc/libwebrtc-sys/src/peer_connection.rs +++ b/crates/livekit-webrtc/libwebrtc-sys/src/peer_connection.rs @@ -5,6 +5,7 @@ use crate::media_stream_interface::ffi::MediaStreamInterface; use crate::rtp_receiver::ffi::RtpReceiver; use crate::rtp_transceiver::ffi::RtpTransceiver; use cxx::UniquePtr; +use crate::rtc_error::ffi::RTCError; #[cxx::bridge(namespace = "livekit")] pub mod ffi { @@ -94,7 +95,9 @@ pub mod ffi { include!("livekit/rtp_transceiver.h"); include!("livekit/media_stream_interface.h"); include!("livekit/candidate.h"); + include!("libwebrtc-sys/src/rtc_error.rs.h"); + type RTCError = crate::rtc_error::ffi::RTCError; type Candidate = crate::candidate::ffi::Candidate; type IceCandidate = crate::jsep::ffi::IceCandidate; type DataChannel = crate::data_channel::ffi::DataChannel; @@ -104,8 +107,10 @@ pub mod ffi { type NativeCreateSdpObserverHandle = crate::jsep::ffi::NativeCreateSdpObserverHandle; type NativeSetLocalSdpObserverHandle = crate::jsep::ffi::NativeSetLocalSdpObserverHandle; type NativeSetRemoteSdpObserverHandle = crate::jsep::ffi::NativeSetRemoteSdpObserverHandle; + type NativeDataChannelInit = crate::data_channel::ffi::NativeDataChannelInit; type SessionDescription = crate::jsep::ffi::SessionDescription; + type NativeAddIceCandidateObserver; type NativePeerConnectionObserver; type PeerConnection; @@ -141,57 +146,80 @@ pub mod ffi { observer: Pin<&mut NativeSetRemoteSdpObserverHandle>, ); + fn create_data_channel( + self: Pin<&mut PeerConnection>, + label: String, + init: UniquePtr, + ) -> Result>; + + fn add_ice_candidate( + self: Pin<&mut PeerConnection>, + candidate: UniquePtr, + observer: Pin<&mut NativeAddIceCandidateObserver>, + ); + fn close(self: Pin<&mut PeerConnection>); fn create_native_peer_connection_observer( observer: Box, ) -> UniquePtr; + fn create_native_add_ice_candidate_observer( + observer: Box, + ) -> UniquePtr; + fn _unique_peer_connection() -> UniquePtr; // Ignore } extern "Rust" { + type AddIceCandidateObserverWrapper; + + fn on_complete( + self: &AddIceCandidateObserverWrapper, + error: RTCError, + ); + type PeerConnectionObserverWrapper; - unsafe fn on_signaling_change( + fn on_signaling_change( self: &PeerConnectionObserverWrapper, new_state: SignalingState, ); - unsafe fn on_add_stream( + fn on_add_stream( self: &PeerConnectionObserverWrapper, stream: UniquePtr, ); - unsafe fn on_remove_stream( + fn on_remove_stream( self: &PeerConnectionObserverWrapper, stream: UniquePtr, ); - unsafe fn on_data_channel( + fn on_data_channel( self: &PeerConnectionObserverWrapper, data_channel: UniquePtr, ); - unsafe fn on_renegotiation_needed(self: &PeerConnectionObserverWrapper); - unsafe fn on_negotiation_needed_event(self: &PeerConnectionObserverWrapper, event: u32); - unsafe fn on_ice_connection_change( + fn on_renegotiation_needed(self: &PeerConnectionObserverWrapper); + fn on_negotiation_needed_event(self: &PeerConnectionObserverWrapper, event: u32); + fn on_ice_connection_change( self: &PeerConnectionObserverWrapper, new_state: IceConnectionState, ); - unsafe fn on_standardized_ice_connection_change( + fn on_standardized_ice_connection_change( self: &PeerConnectionObserverWrapper, new_state: IceConnectionState, ); - unsafe fn on_connection_change( + fn on_connection_change( self: &PeerConnectionObserverWrapper, new_state: PeerConnectionState, ); - unsafe fn on_ice_gathering_change( + fn on_ice_gathering_change( self: &PeerConnectionObserverWrapper, new_state: IceGatheringState, ); - unsafe fn on_ice_candidate( + fn on_ice_candidate( self: &PeerConnectionObserverWrapper, candidate: UniquePtr, ); - unsafe fn on_ice_candidate_error( + fn on_ice_candidate_error( self: &PeerConnectionObserverWrapper, address: String, port: i32, @@ -199,32 +227,32 @@ pub mod ffi { error_code: i32, error_text: String, ); - unsafe fn on_ice_candidates_removed( + fn on_ice_candidates_removed( self: &PeerConnectionObserverWrapper, removed: Vec, ); - unsafe fn on_ice_connection_receiving_change( + fn on_ice_connection_receiving_change( self: &PeerConnectionObserverWrapper, receiving: bool, ); - unsafe fn on_ice_selected_candidate_pair_changed( + fn on_ice_selected_candidate_pair_changed( self: &PeerConnectionObserverWrapper, event: CandidatePairChangeEvent, ); - unsafe fn on_add_track( + fn on_add_track( self: &PeerConnectionObserverWrapper, receiver: UniquePtr, streams: Vec, ); - unsafe fn on_track( + fn on_track( self: &PeerConnectionObserverWrapper, transceiver: UniquePtr, ); - unsafe fn on_remove_track( + fn on_remove_track( self: &PeerConnectionObserverWrapper, receiver: UniquePtr, ); - unsafe fn on_interesting_usage(self: &PeerConnectionObserverWrapper, usage_pattern: i32); + fn on_interesting_usage(self: &PeerConnectionObserverWrapper, usage_pattern: i32); } } @@ -253,6 +281,21 @@ impl Default for ffi::RTCOfferAnswerOptions { } } + +pub struct AddIceCandidateObserverWrapper { + observer: Box +} + +impl AddIceCandidateObserverWrapper { + pub fn new(observer: Box) -> Self { + Self { observer } + } + + fn on_complete(&self, error: RTCError){ + (self.observer)(error); + } +} + pub trait PeerConnectionObserver: Send + Sync { fn on_signaling_change(&self, new_state: ffi::SignalingState); fn on_add_stream(&self, stream: UniquePtr); @@ -298,51 +341,51 @@ impl PeerConnectionObserverWrapper { Self { observer } } - unsafe fn on_signaling_change(&self, new_state: ffi::SignalingState) { - (*self.observer).on_signaling_change(new_state); + fn on_signaling_change(&self, new_state: ffi::SignalingState) { + unsafe { (*self.observer).on_signaling_change(new_state); } } - unsafe fn on_add_stream(&self, stream: UniquePtr) { - (*self.observer).on_add_stream(stream); + fn on_add_stream(&self, stream: UniquePtr) { + unsafe { (*self.observer).on_add_stream(stream); } } - unsafe fn on_remove_stream(&self, stream: UniquePtr) { - (*self.observer).on_remove_stream(stream); + fn on_remove_stream(&self, stream: UniquePtr) { + unsafe { (*self.observer).on_remove_stream(stream); } } - unsafe fn on_data_channel(&self, data_channel: UniquePtr) { - (*self.observer).on_data_channel(data_channel); + fn on_data_channel(&self, data_channel: UniquePtr) { + unsafe { (*self.observer).on_data_channel(data_channel); } } - unsafe fn on_renegotiation_needed(&self) { - (*self.observer).on_renegotiation_needed(); + fn on_renegotiation_needed(&self) { + unsafe { (*self.observer).on_renegotiation_needed(); } } - unsafe fn on_negotiation_needed_event(&self, event: u32) { - (*self.observer).on_negotiation_needed_event(event); + fn on_negotiation_needed_event(&self, event: u32) { + unsafe { (*self.observer).on_negotiation_needed_event(event); } } - unsafe fn on_ice_connection_change(&self, new_state: ffi::IceConnectionState) { - (*self.observer).on_ice_connection_change(new_state); + fn on_ice_connection_change(&self, new_state: ffi::IceConnectionState) { + unsafe { (*self.observer).on_ice_connection_change(new_state); } } - unsafe fn on_standardized_ice_connection_change(&self, new_state: ffi::IceConnectionState) { - (*self.observer).on_standardized_ice_connection_change(new_state); + fn on_standardized_ice_connection_change(&self, new_state: ffi::IceConnectionState) { + unsafe { (*self.observer).on_standardized_ice_connection_change(new_state); } } - unsafe fn on_connection_change(&self, new_state: ffi::PeerConnectionState) { - (*self.observer).on_connection_change(new_state); + fn on_connection_change(&self, new_state: ffi::PeerConnectionState) { + unsafe { (*self.observer).on_connection_change(new_state); } } - unsafe fn on_ice_gathering_change(&self, new_state: ffi::IceGatheringState) { - (*self.observer).on_ice_gathering_change(new_state); + fn on_ice_gathering_change(&self, new_state: ffi::IceGatheringState) { + unsafe { (*self.observer).on_ice_gathering_change(new_state); } } - unsafe fn on_ice_candidate(&self, candidate: UniquePtr) { - (*self.observer).on_ice_candidate(candidate); + fn on_ice_candidate(&self, candidate: UniquePtr) { + unsafe { (*self.observer).on_ice_candidate(candidate); } } - unsafe fn on_ice_candidate_error( + fn on_ice_candidate_error( &self, address: String, port: i32, @@ -350,28 +393,28 @@ impl PeerConnectionObserverWrapper { error_code: i32, error_text: String, ) { - (*self.observer).on_ice_candidate_error(address, port, url, error_code, error_text); + unsafe { (*self.observer).on_ice_candidate_error(address, port, url, error_code, error_text); } } - unsafe fn on_ice_candidates_removed(&self, removed: Vec) { + fn on_ice_candidates_removed(&self, removed: Vec) { let mut vec = Vec::new(); for v in removed { vec.push(v.ptr); } - (*self.observer).on_ice_candidates_removed(vec); + unsafe { (*self.observer).on_ice_candidates_removed(vec); } } - unsafe fn on_ice_connection_receiving_change(&self, receiving: bool) { - (*self.observer).on_ice_connection_receiving_change(receiving); + fn on_ice_connection_receiving_change(&self, receiving: bool) { + unsafe { (*self.observer).on_ice_connection_receiving_change(receiving); } } - unsafe fn on_ice_selected_candidate_pair_changed(&self, event: ffi::CandidatePairChangeEvent) { - (*self.observer).on_ice_selected_candidate_pair_changed(event); + fn on_ice_selected_candidate_pair_changed(&self, event: ffi::CandidatePairChangeEvent) { + unsafe { (*self.observer).on_ice_selected_candidate_pair_changed(event); } } - unsafe fn on_add_track( + fn on_add_track( &self, receiver: UniquePtr, streams: Vec, @@ -382,18 +425,18 @@ impl PeerConnectionObserverWrapper { vec.push(v.ptr); } - (*self.observer).on_add_track(receiver, vec); + unsafe { (*self.observer).on_add_track(receiver, vec); } } - unsafe fn on_track(&self, transceiver: UniquePtr) { - (*self.observer).on_track(transceiver); + fn on_track(&self, transceiver: UniquePtr) { + unsafe { (*self.observer).on_track(transceiver); } } - unsafe fn on_remove_track(&self, receiver: UniquePtr) { - (*self.observer).on_remove_track(receiver); + fn on_remove_track(&self, receiver: UniquePtr) { + unsafe { (*self.observer).on_remove_track(receiver); } } - unsafe fn on_interesting_usage(&self, usage_pattern: i32) { - (*self.observer).on_interesting_usage(usage_pattern); + fn on_interesting_usage(&self, usage_pattern: i32) { + unsafe { (*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 38c0631..0b3f9d6 100644 --- a/crates/livekit-webrtc/libwebrtc-sys/src/peer_connection_factory.cpp +++ b/crates/livekit-webrtc/libwebrtc-sys/src/peer_connection_factory.cpp @@ -70,7 +70,6 @@ namespace livekit{ std::unique_ptr create_rtc_configuration(RTCConfiguration conf){ auto rtc = std::make_unique(); - for (auto &item: conf.ice_servers){ webrtc::PeerConnectionInterface::IceServer ice_server; ice_server.username = item.username.c_str(); diff --git a/crates/livekit-webrtc/libwebrtc-sys/src/rtc_error.cpp b/crates/livekit-webrtc/libwebrtc-sys/src/rtc_error.cpp index 5dd1d05..5ea6217 100644 --- a/crates/livekit-webrtc/libwebrtc-sys/src/rtc_error.cpp +++ b/crates/livekit-webrtc/libwebrtc-sys/src/rtc_error.cpp @@ -7,7 +7,6 @@ #include #include - namespace livekit { RTCError to_error(const webrtc::RTCError &error) { diff --git a/crates/livekit-webrtc/libwebrtc-sys/src/webrtc.cpp b/crates/livekit-webrtc/libwebrtc-sys/src/webrtc.cpp new file mode 100644 index 0000000..b435694 --- /dev/null +++ b/crates/livekit-webrtc/libwebrtc-sys/src/webrtc.cpp @@ -0,0 +1,22 @@ +// +// Created by theom on 18/09/2022. +// + +#include "livekit/webrtc.h" +#include "rtc_base/logging.h" + +namespace livekit { + RTCRuntime::RTCRuntime() { + RTC_LOG(LS_INFO) << "RTCRuntime()"; + RTC_CHECK(rtc::InitializeSSL()) << "Failed to InitializeSSL()"; + } + + RTCRuntime::~RTCRuntime() { + RTC_LOG(LS_INFO) << "~RTCRuntime()"; + RTC_CHECK(rtc::CleanupSSL()) << "Failed to CleanupSSL()"; + } + + std::unique_ptr create_rtc_runtime(){ + return std::make_unique(); + } +} // livekit \ No newline at end of file diff --git a/crates/livekit-webrtc/libwebrtc-sys/src/webrtc.rs b/crates/livekit-webrtc/libwebrtc-sys/src/webrtc.rs new file mode 100644 index 0000000..d890fc9 --- /dev/null +++ b/crates/livekit-webrtc/libwebrtc-sys/src/webrtc.rs @@ -0,0 +1,12 @@ +use cxx::UniquePtr; + +#[cxx::bridge(namespace = "livekit")] +pub mod ffi { + unsafe extern "C++" { + include!("livekit/webrtc.h"); + + type RTCRuntime; + + fn create_rtc_runtime() -> UniquePtr; + } +} diff --git a/crates/livekit-webrtc/src/data_channel.rs b/crates/livekit-webrtc/src/data_channel.rs index 1197138..a1af993 100644 --- a/crates/livekit-webrtc/src/data_channel.rs +++ b/crates/livekit-webrtc/src/data_channel.rs @@ -1,6 +1,8 @@ use cxx::UniquePtr; use libwebrtc_sys::data_channel as sys_dc; +pub use sys_dc::ffi::Priority; + pub struct DataChannel { cxx_handle: UniquePtr, } @@ -10,3 +12,49 @@ impl DataChannel { Self { cxx_handle } } } + +#[derive(Debug)] +pub struct DataChannelInit { + #[deprecated] + reliable: bool, + ordered: bool, + max_retransmit_time: Option, + max_retransmits: Option, + protocol: String, + negotiated: bool, + id: i32, + priority: Option, +} + +impl Default for DataChannelInit { + fn default() -> Self { + Self { + reliable: false, + ordered: true, + max_retransmit_time: None, + max_retransmits: None, + protocol: "".to_string(), + negotiated: false, + id: -1, + priority: None, + } + } +} + +impl From for sys_dc::ffi::DataChannelInit { + fn from(init: DataChannelInit) -> Self { + Self { + reliable: init.reliable, + ordered: init.ordered, + has_max_retransmit_time: init.max_retransmit_time.is_some(), + max_retransmit_time: init.max_retransmit_time.unwrap_or_default(), + has_max_retransmits: init.max_retransmits.is_some(), + max_retransmits: init.max_retransmits.unwrap_or_default(), + protocol: init.protocol, + negotiated: init.negotiated, + id: init.id, + has_priority: init.priority.is_some(), + priority: init.priority.unwrap_or(Priority::Low), + } + } +} diff --git a/crates/livekit-webrtc/src/jsep.rs b/crates/livekit-webrtc/src/jsep.rs index bc76144..3d66eb6 100644 --- a/crates/livekit-webrtc/src/jsep.rs +++ b/crates/livekit-webrtc/src/jsep.rs @@ -2,7 +2,19 @@ use cxx::{SharedPtr, UniquePtr}; use libwebrtc_sys::jsep as sys_jsep; #[derive(Debug)] -pub struct IceCandidate {} +pub struct IceCandidate { + cxx_handle: UniquePtr, +} + +impl IceCandidate { + pub(crate) fn new(cxx_handle: UniquePtr) -> Self { + Self { cxx_handle } + } + + pub(crate) fn release(self) -> UniquePtr { + self.cxx_handle + } +} #[derive(Debug)] pub struct SessionDescription { diff --git a/crates/livekit-webrtc/src/lib.rs b/crates/livekit-webrtc/src/lib.rs index c4403af..4bcbcfc 100644 --- a/crates/livekit-webrtc/src/lib.rs +++ b/crates/livekit-webrtc/src/lib.rs @@ -6,3 +6,4 @@ pub mod peer_connection_factory; pub mod rtc_error; pub mod rtp_receiver; pub mod rtp_transceiver; +pub mod webrtc; diff --git a/crates/livekit-webrtc/src/peer_connection.rs b/crates/livekit-webrtc/src/peer_connection.rs index e855d21..3789b3a 100644 --- a/crates/livekit-webrtc/src/peer_connection.rs +++ b/crates/livekit-webrtc/src/peer_connection.rs @@ -1,12 +1,15 @@ use cxx::UniquePtr; +use libwebrtc_sys::data_channel as sys_dc; use libwebrtc_sys::jsep as sys_jsep; use libwebrtc_sys::peer_connection as sys_pc; use log::trace; +use std::future::Future; +use std::pin::Pin; use std::sync::{Arc, Mutex}; use thiserror::Error; use tokio::sync::{mpsc, oneshot}; -use crate::data_channel::DataChannel; +use crate::data_channel::{DataChannel, DataChannelInit}; use crate::jsep::{IceCandidate, SessionDescription}; use crate::media_stream::MediaStream; use crate::rtc_error::RTCError; @@ -32,19 +35,19 @@ pub struct PeerConnection { observer: Box, // Keep alive for C++ - native_observer: UniquePtr + native_observer: UniquePtr, } impl PeerConnection { pub(crate) fn new( cxx_handle: UniquePtr, observer: Box, - native_observer: UniquePtr + native_observer: UniquePtr, ) -> Self { Self { cxx_handle, observer, - native_observer + native_observer, } } @@ -134,6 +137,38 @@ impl PeerConnection { } } + pub fn create_data_channel( + &mut self, + label: &str, + init: DataChannelInit, + ) -> Result { + let native_init = sys_dc::ffi::create_data_channel_init(init.into()); + let res = self + .cxx_handle + .pin_mut() + .create_data_channel(label.to_string(), native_init); + + match res { + Ok(cxx_handle) => Ok(DataChannel::new(cxx_handle)), + Err(e) => Err(unsafe { RTCError::from(e.what()) }), + } + } + + pub async fn add_ice_candidate(&mut self, candidate: IceCandidate) -> Result<(), SdpError> { + let (tx, mut rx) = mpsc::channel(1); + let observer = sys_pc::AddIceCandidateObserverWrapper::new(Box::new(move |error| { + tx.blocking_send(error).unwrap(); + })); + + let mut native_observer = sys_pc::ffi::create_native_add_ice_candidate_observer(Box::new(observer)); + self.cxx_handle.pin_mut().add_ice_candidate(candidate.release(), native_observer.pin_mut()); + + match rx.recv().await { + Some(value) => Ok(()), + None => Err(SdpError::RecvError("channel closed".to_string())), + } + } + pub fn close(&mut self) { self.cxx_handle.pin_mut().close(); } @@ -467,10 +502,10 @@ impl sys_pc::PeerConnectionObserver for InternalObserver { } fn on_ice_candidate(&self, candidate: UniquePtr) { - trace!("on_ice_candidate"); + trace!("TESTING on_ice_candidate"); let mut handler = self.on_ice_candidate_handler.lock().unwrap(); if let Some(f) = handler.as_mut() { - // TODO(theomonnom) + f(IceCandidate::new(candidate)); } } @@ -567,8 +602,11 @@ impl sys_pc::PeerConnectionObserver for InternalObserver { #[cfg(test)] mod tests { - use crate::peer_connection_factory::PeerConnectionFactory; - use libwebrtc_sys::peer_connection_factory::ffi::RTCConfiguration; + use crate::data_channel::DataChannelInit; + use crate::jsep::IceCandidate; + use crate::peer_connection_factory::{PeerConnectionFactory, ICEServer, RTCConfiguration}; + use tokio::sync::mpsc; + use crate::webrtc::RTCRuntime; fn init_log() { let _ = env_logger::builder().is_test(true).try_init(); @@ -578,14 +616,34 @@ mod tests { async fn create_pc() { init_log(); + let test = RTCRuntime::new(); + let factory = PeerConnectionFactory::new(); let config = RTCConfiguration { - ice_servers: vec![], + ice_servers: vec![ICEServer { + urls: vec!["stun:stun1.l.google.com:19302".to_string()], + username: "".into(), + password: "".into(), + }], }; let mut bob = factory.create_peer_connection(config.clone()).unwrap(); let mut alice = factory.create_peer_connection(config.clone()).unwrap(); + let (bob_ice_tx, mut bob_ice_rx) = mpsc::channel::(1); + let (alice_ice_tx, mut alice_ice_rx) = mpsc::channel::(1); + + bob.on_ice_candidate(Box::new(move |candidate| { + bob_ice_tx.blocking_send(candidate).unwrap(); + })); + + alice.on_ice_candidate(Box::new(move |candidate| { + alice_ice_tx.blocking_send(candidate).unwrap(); + })); + + bob.create_data_channel("test_dc", DataChannelInit::default()) + .unwrap(); + let offer = bob.create_offer().await.unwrap(); bob.set_local_description(offer.clone()).await.unwrap(); alice.set_remote_description(offer).await.unwrap(); @@ -593,6 +651,12 @@ mod tests { alice.set_local_description(answer.clone()).await.unwrap(); bob.set_remote_description(answer).await.unwrap(); + let bob_ice = bob_ice_rx.recv().await.unwrap(); + let alice_ice = alice_ice_rx.recv().await.unwrap(); + + bob.add_ice_candidate(alice_ice).await.unwrap(); + alice.add_ice_candidate(bob_ice).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 2942b63..2724bf5 100644 --- a/crates/livekit-webrtc/src/peer_connection_factory.rs +++ b/crates/livekit-webrtc/src/peer_connection_factory.rs @@ -36,9 +36,7 @@ impl PeerConnectionFactory { match res { Ok(cxx_handle) => Ok(PeerConnection::new(cxx_handle, observer, native_observer)), - Err(e) => { - Err(RTCError::from(e.what())) // TODO - } + Err(e) => Err(RTCError::from(e.what())), } } } diff --git a/crates/livekit-webrtc/src/webrtc.rs b/crates/livekit-webrtc/src/webrtc.rs new file mode 100644 index 0000000..c267100 --- /dev/null +++ b/crates/livekit-webrtc/src/webrtc.rs @@ -0,0 +1,14 @@ +use cxx::UniquePtr; +use libwebrtc_sys::webrtc as sys_rtc; + +pub struct RTCRuntime { + cxx_handle: UniquePtr +} + +impl RTCRuntime { + pub fn new() -> Self { + Self { + cxx_handle: sys_rtc::ffi::create_rtc_runtime() + } + } +} \ No newline at end of file