From a5caf12902eea1c11d6a10535f21d576ca534fe0 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Th=C3=A9o=20Monnom?= Date: Mon, 24 Apr 2023 20:08:27 +0200 Subject: [PATCH] correctly close rtc_session (#56) --- livekit/src/room/room_session.rs | 30 ++++++++++++--------------- livekit/src/rtc_engine/mod.rs | 9 ++------ livekit/src/rtc_engine/rtc_session.rs | 22 ++++++++------------ 3 files changed, 24 insertions(+), 37 deletions(-) diff --git a/livekit/src/room/room_session.rs b/livekit/src/room/room_session.rs index 6ce16e2..c9689da 100644 --- a/livekit/src/room/room_session.rs +++ b/livekit/src/room/room_session.rs @@ -11,7 +11,7 @@ use std::sync::atomic::{AtomicU8, Ordering}; use std::sync::Arc; use tokio::sync::{mpsc, oneshot}; use tokio::task::JoinHandle; -use tracing::{error, info, instrument, Level}; +use tracing::{error, info, instrument, trace, Level}; #[derive(Debug, Clone, Copy, Eq, PartialEq)] pub enum ConnectionState { @@ -158,16 +158,14 @@ impl SessionInner { loop { tokio::select! { res = engine_events.recv() => { - 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") - }; + if let Some(event) = res { + if let Err(err) = self.on_engine_event(event).await { + error!("failed to handle engine event: {:?}", err); + } + } }, _ = &mut close_receiver => { + trace!("closing room_task"); break; } } @@ -185,16 +183,14 @@ impl SessionInner { loop { tokio::select! { res = participant_events.recv() => { - match res { - Some(event) => { - if let Err(err) = self.on_participant_event(&participant, event).await { - error!("failed to handle participant event for {:?}: {:?}", participant.sid(), err); - } - }, - _ => panic!("participant_events has been closed unexpectedly") - }; + if let Some(event) = res { + if let Err(err) = self.on_participant_event(&participant, event).await { + error!("failed to handle participant event for {:?}: {:?}", participant.sid(), err); + } + } }, _ = &mut close_rx => { + trace!("closing participant_task for {:?}", participant.sid()); break; }, } diff --git a/livekit/src/rtc_engine/mod.rs b/livekit/src/rtc_engine/mod.rs index 4335251..e022135 100644 --- a/livekit/src/rtc_engine/mod.rs +++ b/livekit/src/rtc_engine/mod.rs @@ -15,7 +15,7 @@ use tokio::sync::RwLock as AsyncRwLock; use tokio::sync::{mpsc, oneshot}; use tokio::task::JoinHandle; use tokio::time::{interval, Interval}; -use tracing::{error, info, warn}; +use tracing::{error, info, trace, warn}; pub mod lk_runtime; mod peer_transport; @@ -135,10 +135,6 @@ impl RtcEngine { (Self { inner }, engine_events) } - pub(crate) fn lk_runtime(&self) -> Arc { - self.inner.lk_runtime.clone() - } - #[tracing::instrument] pub async fn connect( &self, @@ -266,11 +262,10 @@ impl EngineInner { if let Err(err) = self.on_session_event(event).await { error!("failed to handle session event: {:?}", err); } - } else { - panic!("rtc_sessions has been closed unexpectedly"); } }, _ = &mut close_receiver => { + trace!("closing engine task"); break; } } diff --git a/livekit/src/rtc_engine/rtc_session.rs b/livekit/src/rtc_engine/rtc_session.rs index a3b6917..3ebed60 100644 --- a/livekit/src/rtc_engine/rtc_session.rs +++ b/livekit/src/rtc_engine/rtc_session.rs @@ -236,7 +236,7 @@ impl RtcSession { // Start session tasks let signal_task = tokio::spawn(inner.clone().signal_task(signal_events, close_rx.clone())); - let rtc_task = tokio::spawn(inner.clone().rtc_task(rtc_events, close_rx.clone())); + let rtc_task = tokio::spawn(inner.clone().rtc_session_task(rtc_events, close_rx.clone())); if !inner.info.join_response.subscriber_primary { inner.negotiate_publisher().await?; @@ -348,10 +348,10 @@ impl RtcSession { } impl SessionInner { - async fn rtc_task( + async fn rtc_session_task( self: Arc, mut rtc_events: RtcEvents, - mut close_receiver: watch::Receiver, + mut close_rx: watch::Receiver, ) { loop { tokio::select! { @@ -360,11 +360,9 @@ impl SessionInner { if let Err(err) = self.on_rtc_event(event).await { error!("failed to handle rtc event: {:?}", err); } - } else { - panic!("rtc_events has been closed unexpectedly"); - } - }, - _ = close_receiver.changed() => { + } }, + _ = close_rx.changed() => { + trace!("closing rtc_session_task"); break; } } @@ -374,7 +372,7 @@ impl SessionInner { async fn signal_task( self: Arc, mut signal_events: SignalEvents, - mut close_receiver: watch::Receiver, + mut close_rx: watch::Receiver, ) { loop { tokio::select! { @@ -397,12 +395,10 @@ impl SessionInner { ); } } - } else { - panic!("signal_events has been closed unexpectedly"); } - }, - _ = close_receiver.changed() => { + _ = close_rx.changed() => { + trace!("closing signal_task"); break; } }