From 7467364cd64af6c4ae3537c29a9b0e6d9d4b74f1 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Th=C3=A9o=20Monnom?= Date: Mon, 19 Dec 2022 08:09:27 +0100 Subject: [PATCH] switch comp --- .../src/rtc_engine/engine_internal.rs | 85 ------------------ crates/livekit-core/src/rtc_engine/mod.rs | 1 + .../src/rtc_engine/rtc_session.rs | 88 +++++++++++++++++++ 3 files changed, 89 insertions(+), 85 deletions(-) create mode 100644 crates/livekit-core/src/rtc_engine/rtc_session.rs diff --git a/crates/livekit-core/src/rtc_engine/engine_internal.rs b/crates/livekit-core/src/rtc_engine/engine_internal.rs index e83ce00..a67a3ae 100644 --- a/crates/livekit-core/src/rtc_engine/engine_internal.rs +++ b/crates/livekit-core/src/rtc_engine/engine_internal.rs @@ -53,18 +53,10 @@ pub enum PCState { Closed, } -#[derive(Debug, Clone, Default)] -pub struct SessionInfo { - url: String, - token: String, - options: SignalOptions, - join_response: JoinResponse, -} #[derive(Debug)] pub struct EngineInternal { lk_runtime: Arc, - info: Mutex, session: AsyncRwLock, signal_client: SignalClient, reconnecting: AtomicBool, @@ -74,25 +66,6 @@ pub struct EngineInternal { engine_emitter: EngineEmitter, } -/// This struct holds a WebRTC session -/// The session changes at every reconnection -#[derive(Debug)] -pub struct RTCSession { - publisher_pc: AsyncMutex, - subscriber_pc: AsyncMutex, - - // Publisher data channels - // Used to send data to other participants ( The SFU forwards the messages ) - lossy_dc: DataChannel, - reliable_dc: DataChannel, - - // Subscriber data channels - // These fields are never used, we just keep a strong reference to them, - // so we can receive data from other participants - sub_reliable_dc: Mutex>, - sub_lossy_dc: Mutex>, -} - #[derive(Serialize, Deserialize)] #[allow(non_snake_case)] struct IceCandidateJSON { @@ -101,64 +74,6 @@ struct IceCandidateJSON { candidate: String, } -impl RTCSession { - pub fn new( - lk_runtime: Arc, - session_info: SessionInfo, - ) -> EngineResult<(Self, RTCEvents)> { - let (rtc_emitter, events) = mpsc::unbounded_channel(); - let rtc_config = RTCConfiguration::from(session_info.join_response); - - let mut publisher_pc = PCTransport::new( - lk_runtime - .pc_factory - .create_peer_connection(rtc_config.clone())?, - SignalTarget::Publisher, - ); - - let mut subscriber_pc = PCTransport::new( - lk_runtime - .pc_factory - .create_peer_connection(rtc_config.clone())?, - SignalTarget::Subscriber, - ); - - let mut lossy_dc = publisher_pc.peer_connection().create_data_channel( - LOSSY_DC_LABEL, - DataChannelInit { - ordered: true, - max_retransmits: Some(0), - ..DataChannelInit::default() - }, - )?; - - let mut reliable_dc = publisher_pc.peer_connection().create_data_channel( - RELIABLE_DC_LABEL, - DataChannelInit { - ordered: true, - ..DataChannelInit::default() - }, - )?; - - rtc_events::forward_pc_events(&mut publisher_pc, rtc_emitter.clone()); - rtc_events::forward_pc_events(&mut subscriber_pc, rtc_emitter.clone()); - rtc_events::forward_dc_events(&mut lossy_dc, rtc_emitter.clone()); - rtc_events::forward_dc_events(&mut reliable_dc, rtc_emitter.clone()); - - Ok(( - Self { - publisher_pc: AsyncMutex::new(publisher_pc), - subscriber_pc: AsyncMutex::new(subscriber_pc), - sub_lossy_dc: Default::default(), - sub_reliable_dc: Default::default(), - lossy_dc, - reliable_dc, - }, - events, - )) - } -} - impl EngineInternal { #[tracing::instrument] pub async fn connect( diff --git a/crates/livekit-core/src/rtc_engine/mod.rs b/crates/livekit-core/src/rtc_engine/mod.rs index a061d5a..b432d58 100644 --- a/crates/livekit-core/src/rtc_engine/mod.rs +++ b/crates/livekit-core/src/rtc_engine/mod.rs @@ -38,6 +38,7 @@ mod engine_internal; mod lk_runtime; mod pc_transport; mod rtc_events; +mod rtc_session; pub(crate) type EngineEmitter = mpsc::Sender; pub(crate) type EngineEvents = mpsc::Receiver; diff --git a/crates/livekit-core/src/rtc_engine/rtc_session.rs b/crates/livekit-core/src/rtc_engine/rtc_session.rs new file mode 100644 index 0000000..845a1b3 --- /dev/null +++ b/crates/livekit-core/src/rtc_engine/rtc_session.rs @@ -0,0 +1,88 @@ +#[derive(Debug, Clone, Default)] +pub struct SessionInfo { + url: String, + token: String, + options: SignalOptions, + join_response: JoinResponse, +} + +/// This struct holds a WebRTC session +/// The session changes at every reconnection +#[derive(Debug)] +pub struct RTCSession { + info: SessionInfo, + publisher_pc: AsyncMutex, + subscriber_pc: AsyncMutex, + + // Publisher data channels + // Used to send data to other participants ( The SFU forwards the messages ) + lossy_dc: DataChannel, + reliable_dc: DataChannel, + + // Subscriber data channels + // These fields are never used, we just keep a strong reference to them, + // so we can receive data from other participants + sub_reliable_dc: Mutex>, + sub_lossy_dc: Mutex>, +} + +impl RTCSession { + pub fn new( + lk_runtime: Arc, + session_info: SessionInfo, + ) -> EngineResult<(Self, RTCEvents)> { + let (rtc_emitter, events) = mpsc::unbounded_channel(); + let rtc_config = RTCConfiguration::from(session_info.join_response); + + let mut publisher_pc = PCTransport::new( + lk_runtime + .pc_factory + .create_peer_connection(rtc_config.clone())?, + SignalTarget::Publisher, + ); + + let mut subscriber_pc = PCTransport::new( + lk_runtime + .pc_factory + .create_peer_connection(rtc_config.clone())?, + SignalTarget::Subscriber, + ); + + let mut lossy_dc = publisher_pc.peer_connection().create_data_channel( + LOSSY_DC_LABEL, + DataChannelInit { + ordered: true, + max_retransmits: Some(0), + ..DataChannelInit::default() + }, + )?; + + let mut reliable_dc = publisher_pc.peer_connection().create_data_channel( + RELIABLE_DC_LABEL, + DataChannelInit { + ordered: true, + ..DataChannelInit::default() + }, + )?; + + rtc_events::forward_pc_events(&mut publisher_pc, rtc_emitter.clone()); + rtc_events::forward_pc_events(&mut subscriber_pc, rtc_emitter.clone()); + rtc_events::forward_dc_events(&mut lossy_dc, rtc_emitter.clone()); + rtc_events::forward_dc_events(&mut reliable_dc, rtc_emitter.clone()); + + Ok(( + Self { + info: session_info, + publisher_pc: AsyncMutex::new(publisher_pc), + subscriber_pc: AsyncMutex::new(subscriber_pc), + sub_lossy_dc: Default::default(), + sub_reliable_dc: Default::default(), + lossy_dc, + reliable_dc, + }, + events, + )) + } + + +}