From 9d05259c5a0fdbe97e63f3faea96df04005bca50 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Th=C3=A9o=20Monnom?= Date: Tue, 27 Dec 2022 20:57:35 +0100 Subject: [PATCH] Internal SessionEvents to avoid confusion with the RoomEvents --- .../src/room/participant/local_participant.rs | 7 +- .../livekit-core/src/room/participant/mod.rs | 7 +- .../room/participant/remote_participant.rs | 53 ++++---- crates/livekit-core/src/room/room_session.rs | 126 ++++++++++++------ 4 files changed, 118 insertions(+), 75 deletions(-) diff --git a/crates/livekit-core/src/room/participant/local_participant.rs b/crates/livekit-core/src/room/participant/local_participant.rs index 754b165..503e87d 100644 --- a/crates/livekit-core/src/room/participant/local_participant.rs +++ b/crates/livekit-core/src/room/participant/local_participant.rs @@ -2,7 +2,8 @@ use crate::proto::{data_packet, DataPacket, UserPacket}; use crate::room::participant::{ impl_participant_trait, ParticipantInternalTrait, ParticipantShared, ParticipantTrait, }; -use crate::room::{RoomError, RoomEmitter}; +use crate::room::room_session::SessionEmitter; +use crate::room::RoomError; use crate::rtc_engine::RTCEngine; #[derive(Debug)] @@ -18,10 +19,10 @@ impl LocalParticipant { identity: ParticipantIdentity, name: String, metadata: String, - room_emitter: RoomEmitter, + internal_tx: SessionEmitter, ) -> Self { Self { - shared: ParticipantShared::new(sid, identity, name, metadata, room_emitter), + shared: ParticipantShared::new(sid, identity, name, metadata, internal_tx), rtc_engine, } } diff --git a/crates/livekit-core/src/room/participant/mod.rs b/crates/livekit-core/src/room/participant/mod.rs index 52781c8..eb7db06 100644 --- a/crates/livekit-core/src/room/participant/mod.rs +++ b/crates/livekit-core/src/room/participant/mod.rs @@ -3,6 +3,7 @@ use crate::room::id::{ParticipantIdentity, ParticipantSid, TrackSid}; use crate::room::participant::local_participant::LocalParticipant; use crate::room::participant::remote_participant::RemoteParticipant; use crate::room::publication::{TrackPublication, TrackPublicationTrait}; +use crate::room::room_session::SessionEmitter; use livekit_utils::enum_dispatch; use parking_lot::{Mutex, RwLock}; use std::collections::HashMap; @@ -18,7 +19,7 @@ pub(super) struct ParticipantShared { pub(super) name: Mutex, pub(super) metadata: Mutex, pub(super) tracks: RwLock>, - pub(super) room_emitter: RoomEmitter, + pub(super) internal_tx: SessionEmitter, } impl ParticipantShared { @@ -27,7 +28,7 @@ impl ParticipantShared { identity: ParticipantIdentity, name: String, metadata: String, - room_emitter: RoomEmitter, + internal_tx: SessionEmitter, ) -> Self { Self { sid: Mutex::new(sid), @@ -35,7 +36,7 @@ impl ParticipantShared { name: Mutex::new(name), metadata: Mutex::new(metadata), tracks: Default::default(), - room_emitter + internal_tx, } } diff --git a/crates/livekit-core/src/room/participant/remote_participant.rs b/crates/livekit-core/src/room/participant/remote_participant.rs index a6caf25..0ebcbca 100644 --- a/crates/livekit-core/src/room/participant/remote_participant.rs +++ b/crates/livekit-core/src/room/participant/remote_participant.rs @@ -7,16 +7,18 @@ use crate::room::publication::{ RemoteTrackPublication, TrackPublication, TrackPublicationInternalTrait, TrackPublicationTrait, }; use crate::room::room_session::RoomSession; +use crate::room::room_session::SessionEmitter; +use crate::room::room_session::SessionEvent; use crate::room::track::remote_audio_track::RemoteAudioTrack; use crate::room::track::remote_track::RemoteTrackHandle; use crate::room::track::remote_video_track::RemoteVideoTrack; use crate::room::track::{TrackKind, TrackTrait}; -use crate::room::{RoomEmitter, RoomEvent, TrackError}; +use crate::room::{RoomEvent, TrackError}; use livekit_webrtc::media_stream::MediaStreamTrackHandle; use std::collections::HashSet; use std::time::Duration; -use tokio::time::{sleep, timeout}; -use tracing::{debug, debug_span, error, instrument, Instrument, Level}; +use tokio::time::timeout; +use tracing::{debug, error, instrument, Level}; use super::ParticipantTrait; @@ -33,10 +35,10 @@ impl RemoteParticipant { identity: ParticipantIdentity, name: String, metadata: String, - room_emitter: RoomEmitter, + internal_tx: SessionEmitter, ) -> Self { Self { - shared: ParticipantShared::new(sid, identity, name, metadata, room_emitter), + shared: ParticipantShared::new(sid, identity, name, metadata, internal_tx), } } @@ -50,10 +52,9 @@ impl RemoteParticipant { }) } - #[instrument(level = Level::DEBUG, skip(room_session))] + #[instrument(level = Level::DEBUG)] pub(crate) async fn add_subscribed_media_track( self: Arc, - room_session: RoomSession, sid: TrackSid, media_track: MediaStreamTrackHandle, ) { @@ -67,7 +68,7 @@ impl RemoteParticipant { return publication; } - tokio::task::yield_now(); + tokio::task::yield_now().await; } } }; @@ -108,30 +109,28 @@ impl RemoteParticipant { .add_track_publication(TrackPublication::Remote(remote_publication.clone())); track.start(); - self.shared.room_emitter.send(RoomEvent::TrackSubscribed { - track: track, - publication: remote_publication, - participant: self.clone(), - }); + self.shared + .internal_tx + .send(SessionEvent::Room(RoomEvent::TrackSubscribed { + track: track, + publication: remote_publication, + participant: self.clone(), + })); } else { error!("could not find published track with sid: {:?}", sid); self.shared - .room_emitter - .send(RoomEvent::TrackSubscriptionFailed { + .internal_tx + .send(SessionEvent::Room(RoomEvent::TrackSubscriptionFailed { sid: sid.clone(), error: TrackError::TrackNotFound(sid.clone().to_string()), participant: self.clone(), - }); + })); } } - #[instrument(level = Level::DEBUG, skip(room_session))] - pub(crate) async fn update_tracks( - self: Arc, - room_session: RoomSession, - tracks: Vec, - ) { + #[instrument(level = Level::DEBUG)] + pub(crate) async fn update_tracks(self: Arc, tracks: Vec) { let mut valid_tracks = HashSet::::new(); for track in tracks { @@ -143,10 +142,12 @@ impl RemoteParticipant { .add_track_publication(TrackPublication::Remote(publication.clone())); // This is a new track, fire publish events - self.shared.room_emitter.send(RoomEvent::TrackPublished { - publication: publication.clone(), - participant: self.clone(), - }); + self.shared + .internal_tx + .send(SessionEvent::Room(RoomEvent::TrackPublished { + publication: publication.clone(), + participant: self.clone(), + })); } valid_tracks.insert(track.sid.into()); diff --git a/crates/livekit-core/src/room/room_session.rs b/crates/livekit-core/src/room/room_session.rs index 5b72afc..3aadbca 100644 --- a/crates/livekit-core/src/room/room_session.rs +++ b/crates/livekit-core/src/room/room_session.rs @@ -10,10 +10,19 @@ use parking_lot::{Mutex, RwLock}; use std::collections::HashMap; use std::sync::atomic::{AtomicU8, Ordering}; use std::sync::Arc; +use tokio::sync::mpsc; use tokio::sync::oneshot; use tokio::task::JoinHandle; use tracing::{error, instrument, Level}; +pub(crate) type SessionEmitter = mpsc::UnboundedSender; +pub(crate) type SessionEvents = mpsc::UnboundedReceiver; + +/// Used internally for participants and tracks +pub(crate) enum SessionEvent { + Room(RoomEvent), // Send a public event +} + #[derive(Debug, Clone, Copy, Eq, PartialEq)] pub enum ConnectionState { Disconnected, @@ -43,7 +52,7 @@ struct SessionInner { participants: RwLock>>, rtc_engine: Arc, local_participant: Arc, - room_emitter: RoomEmitter, + internal_tx: SessionEmitter, } #[derive(Debug)] @@ -68,6 +77,7 @@ impl SessionHandle { .connect(url, token, SignalOptions::default()) .await?; + let (internal_tx, internal_rx) = mpsc::unbounded_channel(); let join_response = rtc_engine.join_response().unwrap(); let pi = join_response.participant.unwrap().clone(); let local_participant = Arc::new(LocalParticipant::new( @@ -76,8 +86,9 @@ impl SessionHandle { pi.identity.into(), pi.name, pi.metadata, - room_emitter.clone(), + internal_tx.clone(), )); + let room_info = join_response.room.unwrap(); let inner = Arc::new(SessionInner { state: AtomicU8::new(ConnectionState::Disconnected as u8), @@ -86,7 +97,7 @@ impl SessionHandle { participants: Default::default(), rtc_engine, local_participant, - room_emitter, + internal_tx, }); for pi in join_response.other_participants { @@ -95,13 +106,16 @@ impl SessionHandle { inner.create_participant(pi.sid.into(), pi.identity.into(), pi.name, pi.metadata) }; participant.update_info(pi.clone()); - participant - .update_tracks(RoomSession::from(inner.clone()), pi.tracks) - .await; + participant.update_tracks(pi.tracks).await; } let (close_emitter, close_receiver) = oneshot::channel(); - let session_task = tokio::spawn(inner.clone().room_task(engine_events, close_receiver)); + let session_task = tokio::spawn(inner.clone().room_task( + engine_events, + internal_rx, + close_receiver, + room_emitter, + )); inner .update_connection_state(ConnectionState::Connected) @@ -161,18 +175,32 @@ impl SessionInner { async fn room_task( self: Arc, mut engine_events: EngineEvents, + mut internal_rx: SessionEvents, mut close_receiver: oneshot::Receiver<()>, + room_emitter: RoomEmitter, ) { loop { tokio::select! { + biased; + res = internal_rx.recv() => { + match res { + Some(event) => { + if let Err(err) = self.on_internal_event(event, &room_emitter).await { + error!("failed to handle internal event: {:?}", err); + } + }, + _ => panic!("internal_rx has been closed unexpectedly") + }; + } res = engine_events.recv() => { - if let Some(event) = res { - if let Err(err) = self.on_engine_event(event).await { - error!("failed to handle engine event: {:?}", err); - } - } else { - panic!("engine_events has been closed unexpectedly"); - } + match res { + Some(event) => { + if let Err(err) = self.on_engine_event(event).await { + error!("failed to handle engine event: {:?}", err); + } + }, + _ => panic!("engine_events has been closed unexpectedly") + }; }, _ = &mut close_receiver => { break; @@ -181,6 +209,27 @@ impl SessionInner { } } + async fn on_internal_event( + &self, + event: SessionEvent, + room_emitter: &RoomEmitter, + ) -> EngineResult<()> { + match event { + SessionEvent::Room(event) => { + if self.state.load(Ordering::Acquire) != ConnectionState::Connected as u8 + && matches!(event, RoomEvent::TrackPublished { .. }) + { + return Ok(()); // Ignore the event + } + + // Forward the event to the public channel + let _ = room_emitter.send(event); + } + } + + Ok(()) + } + #[instrument(level = Level::DEBUG)] async fn on_engine_event(self: &Arc, event: EngineEvent) -> RoomResult<()> { match event { @@ -188,7 +237,7 @@ impl SessionInner { EngineEvent::MediaTrack { track, stream, - receiver, + receiver: _, } => { let stream_id = stream.id(); let lk_stream_id = unpack_stream_id(&stream_id); @@ -200,23 +249,14 @@ impl SessionInner { } let (participant_sid, track_sid) = lk_stream_id.unwrap(); + let track_sid = track_sid.to_owned().into(); let remote_participant = self.get_participant(&participant_sid.to_string().into()); if let Some(remote_participant) = remote_participant { - tokio::spawn({ - let session_inner = self.clone(); - { - let track_sid = track_sid.to_owned().into(); - async move { - remote_participant - .add_subscribed_media_track( - RoomSession::from(session_inner), - track_sid, - track, - ) - .await; - } - } + tokio::spawn(async move { + remote_participant + .add_subscribed_media_track(track_sid, track) + .await; }); } else { // The server should send participant updates before sending a new offer @@ -257,8 +297,8 @@ impl SessionInner { self.state.store(state as u8, Ordering::Release); let _ = self - .room_emitter - .send(RoomEvent::ConnectionStateChanged(state)); + .internal_tx + .send(SessionEvent::Room(RoomEvent::ConnectionStateChanged(state))); } /// Update the participants inside a Room. @@ -283,9 +323,7 @@ impl SessionInner { } else { // Participant is already connected, update the it remote_participant.update_info(pi.clone()); - remote_participant - .update_tracks(RoomSession::from(self.clone()), pi.tracks) - .await; + remote_participant.update_tracks(pi.tracks).await; } } else { // Create a new participant @@ -295,13 +333,13 @@ impl SessionInner { }; let _ = self - .room_emitter - .send(RoomEvent::ParticipantConnected(remote_participant.clone())); + .internal_tx + .send(SessionEvent::Room(RoomEvent::ParticipantConnected( + remote_participant.clone(), + ))); remote_participant.update_info(pi.clone()); - remote_participant - .update_tracks(RoomSession::from(self.clone()), pi.tracks) - .await; + remote_participant.update_tracks(pi.tracks).await; } } } @@ -313,9 +351,11 @@ impl SessionInner { self.participants.write().remove(&remote_participant.sid()); // TODO(theomonnom): Unpublish all tracks - let _ = self.room_emitter.send(RoomEvent::ParticipantDisconnected( - remote_participant.clone(), - )); + let _ = self + .internal_tx + .send(SessionEvent::Room(RoomEvent::ParticipantDisconnected( + remote_participant.clone(), + ))); } /// Create a new participant @@ -332,7 +372,7 @@ impl SessionInner { identity, name, metadata, - self.room_emitter.clone(), + self.internal_tx.clone(), )); self.participants.write().insert(sid, p.clone());