dc progress

This commit is contained in:
Théo Monnom
2022-09-18 14:40:50 +02:00
parent 06587d154f
commit 81b45d9c48
9 changed files with 198 additions and 88 deletions
@@ -13,6 +13,18 @@ namespace livekit {
}
void DataChannel::register_observer(std::unique_ptr<NativeDataChannelObserver> observer) {
data_channel_->RegisterObserver(observer.get());
}
void DataChannel::unregister_observer() {
data_channel_->UnregisterObserver();
}
void DataChannel::close() {
return data_channel_->Close();
}
std::unique_ptr<NativeDataChannelInit> create_data_channel_init(DataChannelInit init) {
auto rtc_init = std::make_unique<webrtc::DataChannelInit>();
rtc_init->id = init.id;
@@ -33,4 +45,27 @@ namespace livekit {
return rtc_init;
}
NativeDataChannelObserver::NativeDataChannelObserver(rust::Box<DataChannelObserverWrapper> observer) : observer_(std::move(observer)){
}
void NativeDataChannelObserver::OnStateChange() {
observer_->on_state_change();
}
void NativeDataChannelObserver::OnMessage(const webrtc::DataBuffer &buffer) {
DataBuffer data{};
data.binary = buffer.data.data();
data.len = buffer.data.size();
data.binary = buffer.binary;
observer_->on_message(data);
}
void NativeDataChannelObserver::OnBufferedAmountChange(uint64_t sent_data_size) {
observer_->on_buffered_amount_change(sent_data_size);
}
std::unique_ptr<NativeDataChannelObserver> create_native_peer_connection_observer(rust::Box<DataChannelObserverWrapper> observer){
return std::make_unique<NativeDataChannelObserver>(std::move(observer));
}
} // livekit
@@ -1,4 +1,5 @@
use cxx::UniquePtr;
use std::slice;
#[cxx::bridge(namespace = "livekit")]
pub mod ffi {
@@ -13,7 +14,7 @@ pub mod ffi {
}
#[derive(Debug)]
//#[allow(deprecated)]
#[allow(deprecated)]
pub struct DataChannelInit {
#[deprecated]
reliable: bool,
@@ -26,7 +27,30 @@ pub mod ffi {
negotiated: bool,
id: i32,
has_priority: bool,
priority: Priority
priority: Priority,
}
#[derive(Debug)]
pub struct DataBuffer {
pub ptr: *const u8,
pub len: usize,
pub binary: bool,
}
#[derive(Debug)]
pub enum DataState {
Connecting,
Open,
Closing,
Closed,
}
extern "Rust" {
type DataChannelObserverWrapper;
fn on_state_change(self: &DataChannelObserverWrapper);
fn on_message(self: &DataChannelObserverWrapper, buffer: DataBuffer);
fn on_buffered_amount_change(self: &DataChannelObserverWrapper, sent_data_size: u64);
}
unsafe extern "C++" {
@@ -34,9 +58,47 @@ pub mod ffi {
type DataChannel;
type NativeDataChannelInit;
type NativeDataChannelObserver;
fn close(self: Pin<&mut DataChannel>);
fn create_data_channel_init(init: DataChannelInit) -> UniquePtr<NativeDataChannelInit>;
fn create_native_data_channel_observer(
observer: Box<DataChannelObserverWrapper>,
) -> UniquePtr<NativeDataChannelObserver>;
fn _unique_data_channel() -> UniquePtr<DataChannel>; // Ignore
}
}
// DataChannelObserver
pub trait DataChannelObserver: Send {
fn on_state_change(&self);
fn on_message(&self, data: &[u8], is_binary: bool);
fn on_buffered_amount_change(&self, sent_data_size: u64);
}
pub struct DataChannelObserverWrapper {
observer: Box<dyn DataChannelObserver>,
}
impl DataChannelObserverWrapper {
pub fn new(observer: Box<dyn DataChannelObserver>) -> Self {
Self { observer }
}
fn on_state_change(&self) {
self.observer.on_state_change();
}
fn on_message(&self, buffer: ffi::DataBuffer) {
let data = unsafe { slice::from_raw_parts(buffer.ptr, buffer.len) };
self.observer.on_message(data, buffer.binary);
}
fn on_buffered_amount_change(&self, sent_data_size: u64) {
self.observer.on_buffered_amount_change(sent_data_size);
}
}
@@ -142,45 +142,45 @@ pub mod ffi {
type PeerConnectionObserverWrapper;
unsafe fn on_signaling_change(
self: &mut PeerConnectionObserverWrapper,
self: &PeerConnectionObserverWrapper,
new_state: SignalingState,
);
unsafe fn on_add_stream(
self: &mut PeerConnectionObserverWrapper,
self: &PeerConnectionObserverWrapper,
stream: UniquePtr<MediaStreamInterface>,
);
unsafe fn on_remove_stream(
self: &mut PeerConnectionObserverWrapper,
self: &PeerConnectionObserverWrapper,
stream: UniquePtr<MediaStreamInterface>,
);
unsafe fn on_data_channel(
self: &mut PeerConnectionObserverWrapper,
self: &PeerConnectionObserverWrapper,
data_channel: UniquePtr<DataChannel>,
);
unsafe fn on_renegotiation_needed(self: &mut PeerConnectionObserverWrapper);
unsafe fn on_negotiation_needed_event(self: &mut PeerConnectionObserverWrapper, event: u32);
unsafe fn on_renegotiation_needed(self: &PeerConnectionObserverWrapper);
unsafe fn on_negotiation_needed_event(self: &PeerConnectionObserverWrapper, event: u32);
unsafe fn on_ice_connection_change(
self: &mut PeerConnectionObserverWrapper,
self: &PeerConnectionObserverWrapper,
new_state: IceConnectionState,
);
unsafe fn on_standardized_ice_connection_change(
self: &mut PeerConnectionObserverWrapper,
self: &PeerConnectionObserverWrapper,
new_state: IceConnectionState,
);
unsafe fn on_connection_change(
self: &mut PeerConnectionObserverWrapper,
self: &PeerConnectionObserverWrapper,
new_state: PeerConnectionState,
);
unsafe fn on_ice_gathering_change(
self: &mut PeerConnectionObserverWrapper,
self: &PeerConnectionObserverWrapper,
new_state: IceGatheringState,
);
unsafe fn on_ice_candidate(
self: &mut PeerConnectionObserverWrapper,
self: &PeerConnectionObserverWrapper,
candidate: UniquePtr<IceCandidate>,
);
unsafe fn on_ice_candidate_error(
self: &mut PeerConnectionObserverWrapper,
self: &PeerConnectionObserverWrapper,
address: String,
port: i32,
url: String,
@@ -188,32 +188,32 @@ pub mod ffi {
error_text: String,
);
unsafe fn on_ice_candidates_removed(
self: &mut PeerConnectionObserverWrapper,
self: &PeerConnectionObserverWrapper,
removed: Vec<CandidatePtr>,
);
unsafe fn on_ice_connection_receiving_change(
self: &mut PeerConnectionObserverWrapper,
self: &PeerConnectionObserverWrapper,
receiving: bool,
);
unsafe fn on_ice_selected_candidate_pair_changed(
self: &mut PeerConnectionObserverWrapper,
self: &PeerConnectionObserverWrapper,
event: CandidatePairChangeEvent,
);
unsafe fn on_add_track(
self: &mut PeerConnectionObserverWrapper,
self: &PeerConnectionObserverWrapper,
receiver: UniquePtr<RtpReceiver>,
streams: Vec<MediaStreamPtr>,
);
unsafe fn on_track(
self: &mut PeerConnectionObserverWrapper,
self: &PeerConnectionObserverWrapper,
transceiver: UniquePtr<RtpTransceiver>,
);
unsafe fn on_remove_track(
self: &mut PeerConnectionObserverWrapper,
self: &PeerConnectionObserverWrapper,
receiver: UniquePtr<RtpReceiver>,
);
unsafe fn on_interesting_usage(
self: &mut PeerConnectionObserverWrapper,
self: &PeerConnectionObserverWrapper,
usage_pattern: i32,
);
}
@@ -245,36 +245,36 @@ impl Default for ffi::RTCOfferAnswerOptions {
}
pub trait PeerConnectionObserver: Send + Sync {
fn on_signaling_change(&mut self, new_state: ffi::SignalingState);
fn on_add_stream(&mut self, stream: UniquePtr<MediaStreamInterface>);
fn on_remove_stream(&mut self, stream: UniquePtr<MediaStreamInterface>);
fn on_data_channel(&mut self, data_channel: UniquePtr<DataChannel>);
fn on_renegotiation_needed(&mut self);
fn on_negotiation_needed_event(&mut self, event: u32);
fn on_ice_connection_change(&mut self, new_state: ffi::IceConnectionState);
fn on_standardized_ice_connection_change(&mut self, new_state: ffi::IceConnectionState);
fn on_connection_change(&mut self, new_state: ffi::PeerConnectionState);
fn on_ice_gathering_change(&mut self, new_state: ffi::IceGatheringState);
fn on_ice_candidate(&mut self, candidate: UniquePtr<IceCandidate>);
fn on_signaling_change(&self, new_state: ffi::SignalingState);
fn on_add_stream(&self, stream: UniquePtr<MediaStreamInterface>);
fn on_remove_stream(&self, stream: UniquePtr<MediaStreamInterface>);
fn on_data_channel(&self, data_channel: UniquePtr<DataChannel>);
fn on_renegotiation_needed(&self);
fn on_negotiation_needed_event(&self, event: u32);
fn on_ice_connection_change(&self, new_state: ffi::IceConnectionState);
fn on_standardized_ice_connection_change(&self, new_state: ffi::IceConnectionState);
fn on_connection_change(&self, new_state: ffi::PeerConnectionState);
fn on_ice_gathering_change(&self, new_state: ffi::IceGatheringState);
fn on_ice_candidate(&self, candidate: UniquePtr<IceCandidate>);
fn on_ice_candidate_error(
&mut self,
&self,
address: String,
port: i32,
url: String,
error_code: i32,
error_text: String,
);
fn on_ice_candidates_removed(&mut self, removed: Vec<UniquePtr<Candidate>>);
fn on_ice_connection_receiving_change(&mut self, receiving: bool);
fn on_ice_selected_candidate_pair_changed(&mut self, event: ffi::CandidatePairChangeEvent);
fn on_ice_candidates_removed(&self, removed: Vec<UniquePtr<Candidate>>);
fn on_ice_connection_receiving_change(&self, receiving: bool);
fn on_ice_selected_candidate_pair_changed(&self, event: ffi::CandidatePairChangeEvent);
fn on_add_track(
&mut self,
&self,
receiver: UniquePtr<RtpReceiver>,
streams: Vec<UniquePtr<MediaStreamInterface>>,
);
fn on_track(&mut self, transceiver: UniquePtr<RtpTransceiver>);
fn on_remove_track(&mut self, receiver: UniquePtr<RtpReceiver>);
fn on_interesting_usage(&mut self, usage_pattern: i32);
fn on_track(&self, transceiver: UniquePtr<RtpTransceiver>);
fn on_remove_track(&self, receiver: UniquePtr<RtpReceiver>);
fn on_interesting_usage(&self, usage_pattern: i32);
}
// Thread safety is handled inside PeerConnectionObserver
@@ -289,52 +289,52 @@ impl PeerConnectionObserverWrapper {
Self { observer }
}
unsafe fn on_signaling_change(&mut self, new_state: ffi::SignalingState) {
unsafe fn on_signaling_change(&self, new_state: ffi::SignalingState) {
(*self.observer).on_signaling_change(new_state);
}
unsafe fn on_add_stream(&mut self, stream: UniquePtr<MediaStreamInterface>) {
unsafe fn on_add_stream(&self, stream: UniquePtr<MediaStreamInterface>) {
(*self.observer).on_add_stream(stream);
}
unsafe fn on_remove_stream(&mut self, stream: UniquePtr<MediaStreamInterface>) {
unsafe fn on_remove_stream(&self, stream: UniquePtr<MediaStreamInterface>) {
(*self.observer).on_remove_stream(stream);
}
unsafe fn on_data_channel(&mut self, data_channel: UniquePtr<DataChannel>) {
unsafe fn on_data_channel(&self, data_channel: UniquePtr<DataChannel>) {
(*self.observer).on_data_channel(data_channel);
}
unsafe fn on_renegotiation_needed(&mut self) {
unsafe fn on_renegotiation_needed(&self) {
(*self.observer).on_renegotiation_needed();
}
unsafe fn on_negotiation_needed_event(&mut self, event: u32) {
unsafe fn on_negotiation_needed_event(&self, event: u32) {
(*self.observer).on_negotiation_needed_event(event);
}
unsafe fn on_ice_connection_change(&mut self, new_state: ffi::IceConnectionState) {
unsafe fn on_ice_connection_change(&self, new_state: ffi::IceConnectionState) {
(*self.observer).on_ice_connection_change(new_state);
}
unsafe fn on_standardized_ice_connection_change(&mut self, new_state: ffi::IceConnectionState) {
unsafe fn on_standardized_ice_connection_change(&self, new_state: ffi::IceConnectionState) {
(*self.observer).on_standardized_ice_connection_change(new_state);
}
unsafe fn on_connection_change(&mut self, new_state: ffi::PeerConnectionState) {
unsafe fn on_connection_change(&self, new_state: ffi::PeerConnectionState) {
(*self.observer).on_connection_change(new_state);
}
unsafe fn on_ice_gathering_change(&mut self, new_state: ffi::IceGatheringState) {
unsafe fn on_ice_gathering_change(&self, new_state: ffi::IceGatheringState) {
(*self.observer).on_ice_gathering_change(new_state);
}
unsafe fn on_ice_candidate(&mut self, candidate: UniquePtr<IceCandidate>) {
unsafe fn on_ice_candidate(&self, candidate: UniquePtr<IceCandidate>) {
(*self.observer).on_ice_candidate(candidate);
}
unsafe fn on_ice_candidate_error(
&mut self,
&self,
address: String,
port: i32,
url: String,
@@ -344,7 +344,7 @@ impl PeerConnectionObserverWrapper {
(*self.observer).on_ice_candidate_error(address, port, url, error_code, error_text);
}
unsafe fn on_ice_candidates_removed(&mut self, removed: Vec<ffi::CandidatePtr>) {
unsafe fn on_ice_candidates_removed(&self, removed: Vec<ffi::CandidatePtr>) {
let mut vec = Vec::new();
for v in removed {
@@ -354,19 +354,19 @@ impl PeerConnectionObserverWrapper {
(*self.observer).on_ice_candidates_removed(vec);
}
unsafe fn on_ice_connection_receiving_change(&mut self, receiving: bool) {
unsafe fn on_ice_connection_receiving_change(&self, receiving: bool) {
(*self.observer).on_ice_connection_receiving_change(receiving);
}
unsafe fn on_ice_selected_candidate_pair_changed(
&mut self,
&self,
event: ffi::CandidatePairChangeEvent,
) {
(*self.observer).on_ice_selected_candidate_pair_changed(event);
}
unsafe fn on_add_track(
&mut self,
&self,
receiver: UniquePtr<RtpReceiver>,
streams: Vec<ffi::MediaStreamPtr>,
) {
@@ -379,15 +379,15 @@ impl PeerConnectionObserverWrapper {
(*self.observer).on_add_track(receiver, vec);
}
unsafe fn on_track(&mut self, transceiver: UniquePtr<RtpTransceiver>) {
unsafe fn on_track(&self, transceiver: UniquePtr<RtpTransceiver>) {
(*self.observer).on_track(transceiver);
}
unsafe fn on_remove_track(&mut self, receiver: UniquePtr<RtpReceiver>) {
unsafe fn on_remove_track(&self, receiver: UniquePtr<RtpReceiver>) {
(*self.observer).on_remove_track(receiver);
}
unsafe fn on_interesting_usage(&mut self, usage_pattern: i32) {
unsafe fn on_interesting_usage(&self, usage_pattern: i32) {
(*self.observer).on_interesting_usage(usage_pattern);
}
}