From ac725bb29aa7d5fb65d547a4ab2ae8794de566d4 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Th=C3=A9o=20Monnom?= Date: Tue, 27 Dec 2022 23:06:55 +0100 Subject: [PATCH] connection states and related events --- .../src/room/participant/local_participant.rs | 2 +- .../livekit-core/src/room/participant/mod.rs | 6 +- .../room/participant/remote_participant.rs | 43 ++++----- crates/livekit-core/src/room/room_session.rs | 94 +++++++++++++++---- crates/livekit-core/src/rtc_engine/mod.rs | 2 +- .../src/rtc_engine/rtc_session.rs | 4 +- crates/livekit-core/src/signal_client/mod.rs | 2 +- .../src/signal_client/signal_stream.rs | 2 +- 8 files changed, 105 insertions(+), 50 deletions(-) diff --git a/crates/livekit-core/src/room/participant/local_participant.rs b/crates/livekit-core/src/room/participant/local_participant.rs index 503e87d..2cd42f4 100644 --- a/crates/livekit-core/src/room/participant/local_participant.rs +++ b/crates/livekit-core/src/room/participant/local_participant.rs @@ -49,7 +49,7 @@ impl LocalParticipant { } impl ParticipantInternalTrait for LocalParticipant { - fn update_info(&self, info: ParticipantInfo) { + fn update_info(self: &Arc, info: ParticipantInfo) { self.shared.update_info(info); } } diff --git a/crates/livekit-core/src/room/participant/mod.rs b/crates/livekit-core/src/room/participant/mod.rs index eb7db06..5b473ea 100644 --- a/crates/livekit-core/src/room/participant/mod.rs +++ b/crates/livekit-core/src/room/participant/mod.rs @@ -53,7 +53,7 @@ impl ParticipantShared { } pub(crate) trait ParticipantInternalTrait { - fn update_info(&self, info: ParticipantInfo); + fn update_info(self: &Arc, info: ParticipantInfo); } pub trait ParticipantTrait { @@ -69,7 +69,7 @@ pub enum ParticipantHandle { Remote(Arc), } -impl ParticipantInternalTrait for ParticipantHandle { +impl ParticipantHandle { enum_dispatch!( [Local, Remote] fnc!(update_info, &Self, [info: ParticipantInfo], ()); @@ -114,4 +114,4 @@ macro_rules! impl_participant_trait { pub(super) use impl_participant_trait; -use super::RoomEmitter; + diff --git a/crates/livekit-core/src/room/participant/remote_participant.rs b/crates/livekit-core/src/room/participant/remote_participant.rs index 0ebcbca..98c4fe4 100644 --- a/crates/livekit-core/src/room/participant/remote_participant.rs +++ b/crates/livekit-core/src/room/participant/remote_participant.rs @@ -6,7 +6,6 @@ use crate::room::participant::{ 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; @@ -109,7 +108,8 @@ impl RemoteParticipant { .add_track_publication(TrackPublication::Remote(remote_publication.clone())); track.start(); - self.shared + let _ = self + .shared .internal_tx .send(SessionEvent::Room(RoomEvent::TrackSubscribed { track: track, @@ -119,21 +119,23 @@ impl RemoteParticipant { } else { error!("could not find published track with sid: {:?}", sid); - self.shared - .internal_tx - .send(SessionEvent::Room(RoomEvent::TrackSubscriptionFailed { + let _ = self.shared.internal_tx.send(SessionEvent::Room( + RoomEvent::TrackSubscriptionFailed { sid: sid.clone(), error: TrackError::TrackNotFound(sid.clone().to_string()), participant: self.clone(), - })); + }, + )); } } +} + +impl ParticipantInternalTrait for RemoteParticipant { + fn update_info(self: &Arc, info: ParticipantInfo) { + self.shared.update_info(info.clone()); - #[instrument(level = Level::DEBUG)] - pub(crate) async fn update_tracks(self: Arc, tracks: Vec) { let mut valid_tracks = HashSet::::new(); - - for track in tracks { + for track in info.tracks { if let Some(publication) = self.get_track_publication(&track.sid.clone().into()) { publication.update_info(track.clone()); } else { @@ -141,13 +143,14 @@ impl RemoteParticipant { self.shared .add_track_publication(TrackPublication::Remote(publication.clone())); - // This is a new track, fire publish events - self.shared - .internal_tx - .send(SessionEvent::Room(RoomEvent::TrackPublished { - publication: publication.clone(), - participant: self.clone(), - })); + // This is a new track, fire publish event + let _ = + self.shared + .internal_tx + .send(SessionEvent::Room(RoomEvent::TrackPublished { + publication: publication.clone(), + participant: self.clone(), + })); } valid_tracks.insert(track.sid.into()); @@ -155,10 +158,4 @@ impl RemoteParticipant { } } -impl ParticipantInternalTrait for RemoteParticipant { - fn update_info(&self, info: ParticipantInfo) { - self.shared.update_info(info) - } -} - impl_participant_trait!(RemoteParticipant); diff --git a/crates/livekit-core/src/room/room_session.rs b/crates/livekit-core/src/room/room_session.rs index 3aadbca..9ecf65e 100644 --- a/crates/livekit-core/src/room/room_session.rs +++ b/crates/livekit-core/src/room/room_session.rs @@ -106,7 +106,6 @@ impl SessionHandle { inner.create_participant(pi.sid.into(), pi.identity.into(), pi.name, pi.metadata) }; participant.update_info(pi.clone()); - participant.update_tracks(pi.tracks).await; } let (close_emitter, close_receiver) = oneshot::channel(); @@ -117,9 +116,7 @@ impl SessionHandle { room_emitter, )); - inner - .update_connection_state(ConnectionState::Connected) - .await; + inner.update_connection_state(ConnectionState::Connected); let session = Self { session: RoomSession::from(inner), @@ -233,7 +230,7 @@ impl SessionInner { #[instrument(level = Level::DEBUG)] async fn on_engine_event(self: &Arc, event: EngineEvent) -> RoomResult<()> { match event { - EngineEvent::ParticipantUpdate(update) => self.handle_participant_update(update).await, + EngineEvent::ParticipantUpdate(update) => self.handle_participant_update(update), EngineEvent::MediaTrack { track, stream, @@ -267,11 +264,25 @@ impl SessionInner { )))?; } } - EngineEvent::Resuming => {} - EngineEvent::Resumed => {} - EngineEvent::Restarting => {} - EngineEvent::Restarted => {} - EngineEvent::Disconnected => {} + EngineEvent::Resuming => { + if self.update_connection_state(ConnectionState::Reconnecting) { + let _ = self + .internal_tx + .send(SessionEvent::Room(RoomEvent::Reconnecting)); + } + } + EngineEvent::Resumed => { + self.update_connection_state(ConnectionState::Connected); + let _ = self + .internal_tx + .send(SessionEvent::Room(RoomEvent::Reconnected)); + + // TODO(theomonnom): Update subscriptions settings + // TODO(theomonnom): Send sync state + } + EngineEvent::Restarting => self.handle_restarting(), + EngineEvent::Restarted => self.handle_restarted(), + EngineEvent::Disconnected => self.handle_disconnected(), } Ok(()) @@ -282,30 +293,32 @@ impl SessionInner { self.rtc_engine.close().await; } - fn get_participant(self: &Arc, sid: &ParticipantSid) -> Option> { + fn get_participant(&self, sid: &ParticipantSid) -> Option> { self.participants.read().get(sid).cloned() } /// Change the connection state and emit an event /// Does nothing if the state is already the same #[instrument(level = Level::DEBUG)] - async fn update_connection_state(self: &Arc, state: ConnectionState) { + fn update_connection_state(&self, state: ConnectionState) -> bool { let old_state = self.state.load(Ordering::Acquire); if old_state == state as u8 { - return; + return false; } self.state.store(state as u8, Ordering::Release); let _ = self .internal_tx .send(SessionEvent::Room(RoomEvent::ConnectionStateChanged(state))); + + return true; } /// Update the participants inside a Room. /// It'll create, update or remove a participant /// It also update the participant tracks. #[instrument(level = Level::DEBUG)] - async fn handle_participant_update(self: &Arc, update: proto::ParticipantUpdate) { + fn handle_participant_update(&self, update: proto::ParticipantUpdate) { for pi in update.participants { if pi.sid == self.local_participant.sid() || pi.identity == self.local_participant.identity() @@ -323,7 +336,6 @@ impl SessionInner { } else { // Participant is already connected, update the it remote_participant.update_info(pi.clone()); - remote_participant.update_tracks(pi.tracks).await; } } else { // Create a new participant @@ -339,7 +351,6 @@ impl SessionInner { ))); remote_participant.update_info(pi.clone()); - remote_participant.update_tracks(pi.tracks).await; } } } @@ -347,7 +358,7 @@ impl SessionInner { /// A participant has disconnected /// Cleanup the participant and emit an event #[instrument(level = Level::DEBUG)] - fn handle_participant_disconnect(self: &Arc, remote_participant: Arc) { + fn handle_participant_disconnect(&self, remote_participant: Arc) { self.participants.write().remove(&remote_participant.sid()); // TODO(theomonnom): Unpublish all tracks @@ -358,10 +369,57 @@ impl SessionInner { ))); } + #[instrument(level = Level::DEBUG)] + fn handle_restarting(&self) { + // Remove existing participants/subscriptions on full reconnect + for (_, participant) in self.participants.read().iter() { + self.handle_participant_disconnect(participant.clone()); + } + + if self.update_connection_state(ConnectionState::Reconnecting) { + let _ = self + .internal_tx + .send(SessionEvent::Room(RoomEvent::Reconnecting)); + } + } + + #[instrument(level = Level::DEBUG)] + fn handle_restarted(&self) { + // Full reconnect succeeded! + let join_response = self.rtc_engine.join_response().unwrap(); + + self.update_connection_state(ConnectionState::Connected); + let _ = self + .internal_tx + .send(SessionEvent::Room(RoomEvent::Reconnected)); + + if let Some(pi) = join_response.participant { + self.local_participant.update_info(pi); // The sid may have changed + } + + self.handle_participant_update(proto::ParticipantUpdate { + participants: join_response.other_participants, + }); + + // TODO(theomonnom): unpublish & republish tracks + } + + #[instrument(level = Level::DEBUG)] + fn handle_disconnected(&self) { + if self.state.load(Ordering::Acquire) == ConnectionState::Disconnected as u8 { + return; + } + + self.update_connection_state(ConnectionState::Disconnected); + let _ = self + .internal_tx + .send(SessionEvent::Room(RoomEvent::Disconnected)); + } + /// Create a new participant /// Also add it to the participants list fn create_participant( - self: &Arc, + &self, sid: ParticipantSid, identity: ParticipantIdentity, name: String, diff --git a/crates/livekit-core/src/rtc_engine/mod.rs b/crates/livekit-core/src/rtc_engine/mod.rs index 9f9a7f5..dbe6b63 100644 --- a/crates/livekit-core/src/rtc_engine/mod.rs +++ b/crates/livekit-core/src/rtc_engine/mod.rs @@ -241,7 +241,7 @@ impl EngineInner { self.close().await; } } - SessionEvent::Data { data } => {} + SessionEvent::Data { data: _ } => {} SessionEvent::MediaTrack { track, stream, diff --git a/crates/livekit-core/src/rtc_engine/rtc_session.rs b/crates/livekit-core/src/rtc_engine/rtc_session.rs index a39eef4..7e6d847 100644 --- a/crates/livekit-core/src/rtc_engine/rtc_session.rs +++ b/crates/livekit-core/src/rtc_engine/rtc_session.rs @@ -12,7 +12,7 @@ use tokio::time::sleep; use prost::Message; use serde::{Deserialize, Serialize}; -use tracing::{debug, error, info, trace, warn}; +use tracing::{debug, error, trace, warn}; use crate::{proto, signal_client}; use livekit_webrtc::data_channel::{DataChannel, DataChannelInit, DataState}; @@ -501,7 +501,7 @@ impl SessionInner { let data = DataPacket::decode(&*data)?; match data.value.unwrap() { - Value::User(user) => { + Value::User(_user) => { // TODO(theomonnom) Send event } Value::Speaker(_) => { diff --git a/crates/livekit-core/src/signal_client/mod.rs b/crates/livekit-core/src/signal_client/mod.rs index e840ec5..ad58a88 100644 --- a/crates/livekit-core/src/signal_client/mod.rs +++ b/crates/livekit-core/src/signal_client/mod.rs @@ -1,5 +1,5 @@ use std::fmt::Debug; -use std::sync::RwLockWriteGuard; + use std::time::Duration; use livekit_webrtc::peer_connection_factory::{ diff --git a/crates/livekit-core/src/signal_client/signal_stream.rs b/crates/livekit-core/src/signal_client/signal_stream.rs index c6d4c23..1692397 100644 --- a/crates/livekit-core/src/signal_client/signal_stream.rs +++ b/crates/livekit-core/src/signal_client/signal_stream.rs @@ -4,7 +4,7 @@ use prost::Message as ProstMessage; use tokio::net::TcpStream; use tokio::sync::{mpsc, oneshot}; use tokio::task::JoinHandle; -use tokio_tungstenite::tungstenite::error::ProtocolError; + use tokio_tungstenite::tungstenite::protocol::frame::coding::CloseCode; use tokio_tungstenite::tungstenite::protocol::CloseFrame; use tokio_tungstenite::tungstenite::Message;