diff --git a/crates/livekit-core/src/room/room_session.rs b/crates/livekit-core/src/room/room_session.rs index 4f284d1..f74cade 100644 --- a/crates/livekit-core/src/room/room_session.rs +++ b/crates/livekit-core/src/room/room_session.rs @@ -199,7 +199,7 @@ impl SessionInner { error!("failed to handle participant event for {:?}: {:?}", participant.sid(), err); } }, - _ => panic!("engine_events has been closed unexpectedly") + _ => panic!("participant_events has been closed unexpectedly") }; }, _ = &mut close_rx => { diff --git a/crates/livekit-utils/src/observer.rs b/crates/livekit-utils/src/observer.rs index 46cb800..14c26b2 100644 --- a/crates/livekit-utils/src/observer.rs +++ b/crates/livekit-utils/src/observer.rs @@ -34,6 +34,6 @@ where pub fn dispatch(&mut self, msg: &T) { self.senders - .retain(|sender| sender.send(msg.clone()).is_err()); + .retain(|sender| sender.send(msg.clone()).is_ok()); } } diff --git a/examples/Cargo.lock b/examples/Cargo.lock index 5553bfc..4cfaa19 100644 --- a/examples/Cargo.lock +++ b/examples/Cargo.lock @@ -1161,6 +1161,10 @@ dependencies = [ [[package]] name = "livekit-utils" version = "0.1.0" +dependencies = [ + "parking_lot", + "tokio", +] [[package]] name = "livekit-webrtc" diff --git a/examples/simple_room/src/app.rs b/examples/simple_room/src/app.rs index f17cba7..1612c46 100644 --- a/examples/simple_room/src/app.rs +++ b/examples/simple_room/src/app.rs @@ -8,15 +8,17 @@ use std::sync::{ atomic::{AtomicBool, Ordering}, Arc, }; -use tokio::sync::mpsc; +use tokio::sync::{mpsc, oneshot}; +use tracing::error; -use livekit::room::{ConnectionState, Room, RoomError, SimulateScenario}; +use livekit::room::room_session::ConnectionState; +use livekit::room::{Room, RoomEvent, RoomEvents, SimulateScenario}; // Useful default constants for developing const DEFAULT_URL: &str = "ws://localhost:7880"; -const DEFAULT_TOKEN : &str = "eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJleHAiOjIzODQ4MDY0NzMsImlzcyI6IkFQSXpLYkFTaUNWYWtnSiIsIm5hbWUiOiJuYXRpdmUiLCJuYmYiOjE2NjQ4MDY0NzMsInN1YiI6Im5hdGl2ZSIsInZpZGVvIjp7InJvb21DcmVhdGUiOnRydWUsInJvb21Kb2luIjp0cnVlfX0.BgVdBnq3XFD3_BQHoe1azqjifYysubgFl6Qlzu9IQGI"; +const DEFAULT_TOKEN : &str = "eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJleHAiOjE5MDY2MTMyODgsImlzcyI6IkFQSVRzRWZpZFpqclFvWSIsIm5hbWUiOiJuYXRpdmUiLCJuYmYiOjE2NzI2MTMyODgsInN1YiI6Im5hdGl2ZSIsInZpZGVvIjp7InJvb20iOiJ0ZXN0Iiwicm9vbUFkbWluIjp0cnVlLCJyb29tQ3JlYXRlIjp0cnVlLCJyb29tSm9pbiI6dHJ1ZSwicm9vbUxpc3QiOnRydWV9fQ.uSNIangMRu8jZD5mnRYoCHjcsQWCrJXgHCs0aNIgBFY"; -// eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJleHAiOjIzODQ4MDY3MzAsImlzcyI6IkFQSXpLYkFTaUNWYWtnSiIsIm5hbWUiOiJ3ZWIiLCJuYmYiOjE2NjQ4MDY3MzAsInN1YiI6IndlYiIsInZpZGVvIjp7InJvb21DcmVhdGUiOnRydWUsInJvb21Kb2luIjp0cnVlfX0.VbDoULjX1CVGZu2sPy3SvWYlVZUBXxQVPmdB9BnmlN4 +// eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJleHAiOjE5MDY2MTM0MzcsImlzcyI6IkFQSVRzRWZpZFpqclFvWSIsIm5hbWUiOiJ3ZWIiLCJuYmYiOjE2NzI2MTM0MzcsInN1YiI6IndlYiIsInZpZGVvIjp7InJvb20iOiJ0ZXN0Iiwicm9vbUFkbWluIjp0cnVlLCJyb29tQ3JlYXRlIjp0cnVlLCJyb29tSm9pbiI6dHJ1ZSwicm9vbUxpc3QiOnRydWV9fQ.DFTXt60n1kzGq4cSuOhbFBTQW2nd3rlcXKQ54sXsP8s use winit::{ event::*, @@ -25,7 +27,7 @@ use winit::{ }; struct AppState { - room: Mutex, + room: Mutex)>>, connecting: AtomicBool, } @@ -67,7 +69,7 @@ pub fn run(rt: tokio::runtime::Runtime) { let (ui_cmd_tx, ui_cmd_rx) = mpsc::unbounded_channel::(); let state = Arc::new(AppState { - room: Mutex::new(Room::new()), + room: Default::default(), connecting: AtomicBool::new(false), }); @@ -88,35 +90,42 @@ pub fn run(rt: tokio::runtime::Runtime) { // Async event loop tokio::spawn(async move { - { - let events = state.room.lock().events(); - events.on_track_subscribed({ - let ui_cmd_tx = ui_cmd_tx.clone(); - move |event| { - let ui_cmd_tx = ui_cmd_tx.clone(); - async move { - ui_cmd_tx.send(UiCmd::TrackSubscribed { event }).unwrap(); - } - } - }); - } while let Some(event) = async_cmd_rx.recv().await { match event { AsyncCmd::RoomConnect { url, token } => { + if let Some((room, close_tx)) = state.room.lock().take() { + // Close the Room if already connected + let _ = room.close().await; + let _ = close_tx.send(()); + } + state.connecting.store(true, Ordering::SeqCst); - let mut room = state.room.lock(); - ui_cmd_tx - .send(UiCmd::ConnectResult { - result: room.connect(&url, &token).await, - }) - .unwrap(); + let res = Room::connect(&url, &token).await; + match res { + Ok((room, room_events)) => { + let (close_tx, close_rx) = oneshot::channel(); + state.room.lock().replace((room, close_tx)); + + tokio::spawn(room_task( + state.clone(), + room_events, + close_rx, + ui_cmd_tx.clone(), + )); + + let _ = ui_cmd_tx.send(UiCmd::ConnectResult { result: Ok(()) }); + } + Err(err) => { + let _ = ui_cmd_tx.send(UiCmd::ConnectResult { result: Err(err) }); + } + } state.connecting.store(false, Ordering::SeqCst); } AsyncCmd::SimulateScenario { scenario } => { - if let Some(handle) = state.room.lock().get_handle() { - let _ = handle.simulate_scenario(scenario).await; + if let Some((room, _)) = state.room.lock().as_ref() { + let _ = room.session().simulate_scenario(scenario).await; } } } @@ -132,6 +141,24 @@ pub fn run(rt: tokio::runtime::Runtime) { }); } +async fn room_task( + _app_state: Arc, + mut room_events: RoomEvents, + mut close_rx: oneshot::Receiver<()>, + ui_cmd_tx: mpsc::UnboundedSender, +) { + loop { + tokio::select! { + Some(event) = room_events.recv() => { + let _ = ui_cmd_tx.send(UiCmd::RoomEvent{event}); + } + _ = &mut close_rx => { + //break; + } + } + } +} + impl App { fn update(&mut self, event: Event<'_, T>, control_flow: &mut ControlFlow) { if let Ok(cmd) = self.cmd_rx.try_recv() { @@ -143,20 +170,25 @@ impl App { self.connection_failure = None } } - UiCmd::TrackSubscribed { event } => { - match event.track { - RemoteTrackHandle::Video(video_track) => { - // Create a new VideoRenderer - let video_renderer = VideoRenderer::new( - self.egui_painter.render_state().clone().unwrap(), - video_track.rtc_track(), - ); - self.video_renderers.push(video_renderer); + UiCmd::RoomEvent { event } => { + match event { + RoomEvent::TrackSubscribed { track, .. } => { + match track { + RemoteTrackHandle::Video(video_track) => { + // Create a new VideoRenderer + let video_renderer = VideoRenderer::new( + self.egui_painter.render_state().clone().unwrap(), + video_track.rtc_track(), + ); + self.video_renderers.push(video_renderer); + } + RemoteTrackHandle::Audio(_) => { + // The demo doesn't support Audio rendering at the moment. + } + }; } - RemoteTrackHandle::Audio(_) => { - // The demo doesn't support Audio rendering at the moment. - } - }; + _ => {} + } } } } diff --git a/examples/simple_room/src/events.rs b/examples/simple_room/src/events.rs index 8efe803..73df9ea 100644 --- a/examples/simple_room/src/events.rs +++ b/examples/simple_room/src/events.rs @@ -1,17 +1,13 @@ -use livekit::{events::TrackSubscribedEvent, room::SimulateScenario}; +use livekit::room::{RoomEvent, RoomResult, SimulateScenario}; #[derive(Debug)] pub enum AsyncCmd { RoomConnect { url: String, token: String }, - SimulateScenario { scenario: SimulateScenario } + SimulateScenario { scenario: SimulateScenario }, } #[derive(Debug)] pub enum UiCmd { - ConnectResult { - result: livekit::room::RoomResult<()>, - }, - TrackSubscribed { - event: TrackSubscribedEvent, - }, + ConnectResult { result: RoomResult<()> }, + RoomEvent { event: RoomEvent }, } diff --git a/examples/simple_room/src/video_renderer.rs b/examples/simple_room/src/video_renderer.rs index 0f0da9e..60e4a38 100644 --- a/examples/simple_room/src/video_renderer.rs +++ b/examples/simple_room/src/video_renderer.rs @@ -3,13 +3,14 @@ use livekit::webrtc::video_frame_buffer::PlanarYuv8Buffer; use livekit::webrtc::video_frame_buffer::PlanarYuvBuffer; use livekit::webrtc::video_frame_buffer::VideoFrameBufferTrait; use livekit::webrtc::yuv_helper; -use tracing::debug_span; use std::convert::TryInto; use std::num::NonZeroU32; use std::{ ops::DerefMut, sync::{Arc, Mutex}, }; +use tracing::debug_span; +use tracing::{error, warn}; pub struct VideoRenderer { internal: Arc>, @@ -98,10 +99,12 @@ impl VideoRenderer { egui_texture: None, })); + error!("siubsc"); rtc_track.on_frame({ let internal = internal.clone(); Box::new(move |_frame, buffer| { + warn!("got frame"); let span = debug_span!("texture_upload"); let _enter = span.enter();