diff --git a/crates/livekit-core/src/room/mod.rs b/crates/livekit-core/src/room/mod.rs index 97c7a26..75a84f8 100644 --- a/crates/livekit-core/src/room/mod.rs +++ b/crates/livekit-core/src/room/mod.rs @@ -23,8 +23,8 @@ use crate::signal_client::SignalOptions; pub use crate::rtc_engine::SimulateScenario; +mod room_session; pub mod id; -mod internal; pub mod participant; pub mod publication; pub mod track; @@ -48,11 +48,11 @@ pub enum ConnectionState { } #[derive(Clone, Debug)] -pub struct RoomHandle { - session: Arc, +pub struct RoomSession { + internal: Arc, } -impl RoomHandle { +impl RoomSession { pub fn sid(&self) -> String { self.session.sid.lock().clone() } @@ -72,7 +72,7 @@ impl RoomHandle { #[derive(Debug, Default)] pub struct Room { - session: Option, + session: Option>, events: Arc, // Keep the same RoomEvents across sessions } @@ -90,9 +90,73 @@ impl Room { self.events.clone() } - pub fn get_handle(&self) -> Option { + pub fn session(&self) -> Option<> { self.internal.as_ref().map(|internal| RoomHandle { internal: internal.clone(), }) } } + +#[derive(Debug)] +pub struct RoomInternal { + inner: Arc, + session_task: JoinHandle<()>, + close_emitter: oneshot::Sender<()>, +} + +impl RoomInternal { + pub async fn connect(room_events: Arc, url: &str, token: &str) -> RoomResult { + let (rtc_engine, engine_events) = RTCEngine::new(); + let rtc_engine = Arc::new(rtc_engine); + rtc_engine + .connect(url, token, SignalOptions::default()) + .await?; + + let join_response = rtc_engine.join_response().unwrap(); + let pi = join_response.participant.unwrap().clone(); + let local_participant = Arc::new(LocalParticipant::new( + rtc_engine.clone(), + pi.sid.into(), + pi.identity.into(), + pi.name, + pi.metadata, + )); + let room_info = join_response.room.unwrap(); + let inner = Arc::new(SessionInner { + state: AtomicU8::new(ConnectionState::Connecting as u8), + sid: Mutex::new(room_info.sid), + name: Mutex::new(room_info.name), + participants: Default::default(), + rtc_engine, + local_participant, + room_events, + }); + + for pi in join_response.other_participants { + let participant = { + let pi = pi.clone(); + inner.create_participant(pi.sid.into(), pi.identity.into(), pi.name, pi.metadata) + }; + participant.update_info(pi.clone()); + participant + .update_tracks(RoomHandle::from(inner.clone()), pi.tracks) + .await; + } + + let (close_emitter, close_receiver) = oneshot::channel(); + let session_task = tokio::spawn(inner.room_task(engine_events, close_receiver)); + + let session = Self { + inner, + session_task, + close_emitter, + }; + Ok(session) + } + + pub async fn close(self) { + self.inner.close(); + let _ = self.close_emitter.send(()); + self.session_task.await; + } +} diff --git a/crates/livekit-core/src/room/internal.rs b/crates/livekit-core/src/room/room_session.rs similarity index 81% rename from crates/livekit-core/src/room/internal.rs rename to crates/livekit-core/src/room/room_session.rs index e901655..5222f10 100644 --- a/crates/livekit-core/src/room/internal.rs +++ b/crates/livekit-core/src/room/room_session.rs @@ -8,14 +8,14 @@ use tokio::task::JoinHandle; use crate::events::{ParticipantConnectedEvent, ParticipantDisconnectedEvent, RoomEvents}; use crate::proto::{self, participant_info}; use crate::room::ConnectionState; -use crate::rtc_engine::{EngineEvent, EngineEvents, RTCEngine}; +use crate::rtc_engine::{EngineEvent, EngineEvents, EngineResult, RTCEngine}; use crate::signal_client::SignalOptions; use super::id::{ParticipantIdentity, ParticipantSid}; use super::participant::local_participant::LocalParticipant; use super::participant::remote_participant::RemoteParticipant; use super::participant::{ParticipantInternalTrait, ParticipantTrait}; -use super::{RoomError, RoomHandle, RoomResult}; +use super::{RoomError, RoomResult, SimulateScenario}; use tracing::{error, instrument, Level}; #[derive(Debug)] @@ -29,74 +29,40 @@ pub struct SessionInner { pub room_events: Arc, } -#[derive(Debug)] +#[derive(Clone, Debug)] pub struct RoomSession { inner: Arc, - session_task: JoinHandle<()>, - close_emitter: oneshot::Sender<()>, } impl RoomSession { - pub async fn connect( - room_events: Arc, - url: &str, - token: &str, - ) -> RoomResult { - let (rtc_engine, engine_events) = RTCEngine::new(); - let rtc_engine = Arc::new(rtc_engine); - rtc_engine - .connect(url, token, SignalOptions::default()) - .await?; - - let join_response = rtc_engine.join_response().unwrap(); - let pi = join_response.participant.unwrap().clone(); - let local_participant = Arc::new(LocalParticipant::new( - rtc_engine.clone(), - pi.sid.into(), - pi.identity.into(), - pi.name, - pi.metadata, - )); - let room_info = join_response.room.unwrap(); - let inner = Arc::new(SessionInner { - state: AtomicU8::new(ConnectionState::Connecting as u8), - sid: Mutex::new(room_info.sid), - name: Mutex::new(room_info.name), - participants: Default::default(), - rtc_engine, - local_participant, - room_events, - }); - - for pi in join_response.other_participants { - let participant = { - let pi = pi.clone(); - inner.create_participant(pi.sid.into(), pi.identity.into(), pi.name, pi.metadata) - }; - participant.update_info(pi.clone()); - participant - .update_tracks(RoomHandle::from(inner.clone()), pi.tracks) - .await; - } - - let (close_emitter, close_receiver) = oneshot::channel(); - let session_task = tokio::spawn(inner.room_task(engine_events, close_receiver)); - - let session = Self { - inner, - session_task, - close_emitter, - }; - Ok(session) + pub fn sid(&self) -> String { + self.session.sid.lock().clone() } - pub async fn close(self) { - self.inner.close(); - let _ = self.close_emitter.send(()); - self.session_task.await; + pub fn name(&self) -> String { + self.internal.name.lock().clone() + } + + pub fn local_participant(&self) -> Arc { + self.internal.local_participant.clone() + } + + pub async fn simulate_scenario(&self, scenario: SimulateScenario) -> EngineResult<()> { + self.internal.rtc_engine.simulate_scenario(scenario).await } } +impl RoomSession { + + pub(crate) async fn close(&self) -> RoomResult<()> { + self.internal.rtc_engine.close().await?; + self.internal.room_events.close(); + Ok(()) + } +} + +// Connect me to a database + impl SessionInner { async fn room_task( self: Arc,