From 7b0d3d09ea92974a4435d82721c3021ed5b9b73b Mon Sep 17 00:00:00 2001 From: euzu Date: Sat, 11 Oct 2025 16:28:07 +0200 Subject: [PATCH 1/2] fix for hls streaming --- CHANGELOG.md | 3 + backend/src/api/endpoints/hls_api.rs | 2 +- backend/src/api/endpoints/m3u_api.rs | 2 +- backend/src/api/endpoints/xtream_api.rs | 2 +- backend/src/api/model/active_user_manager.rs | 167 +++++++----------- .../src/api/model/streams/provider_stream.rs | 2 +- .../model/streams/provider_stream_factory.rs | 2 +- 7 files changed, 73 insertions(+), 107 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 00f229492..b71c44c76 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,4 +1,7 @@ # Changelog +# 3.1.8 (2025-10-xx) +- Fixed hls streaming caused by session removal and wrong headers. + # 3.1.7 (2025-10-10) - Added Dark/Bright theme switch - Resource proxy retries failed requests up to three times and respects the `Retry-After` header (falls back to 100 ms wait) diff --git a/backend/src/api/endpoints/hls_api.rs b/backend/src/api/endpoints/hls_api.rs index 3b4da92f8..3d5dab20b 100644 --- a/backend/src/api/endpoints/hls_api.rs +++ b/backend/src/api/endpoints/hls_api.rs @@ -162,7 +162,7 @@ async fn hls_api_stream( let user_session_token = format!("{fingerprint}{virtual_id}"); let mut user_session = app_state .active_users - .get_user_session(&user.username, &user_session_token).await; + .get_and_update_user_session(&user.username, &user_session_token).await; if let Some(session) = &mut user_session { if session.permission == UserConnectionPermission::Exhausted { diff --git a/backend/src/api/endpoints/m3u_api.rs b/backend/src/api/endpoints/m3u_api.rs index 519958462..01809be61 100644 --- a/backend/src/api/endpoints/m3u_api.rs +++ b/backend/src/api/endpoints/m3u_api.rs @@ -129,7 +129,7 @@ async fn m3u_api_stream( let session_key = format!("{fingerprint}{virtual_id}"); let user_session = app_state .active_users - .get_user_session(&user.username, &session_key).await; + .get_and_update_user_session(&user.username, &session_key).await; let session_url = if let Some(session) = &user_session { if session.permission == UserConnectionPermission::Exhausted { diff --git a/backend/src/api/endpoints/xtream_api.rs b/backend/src/api/endpoints/xtream_api.rs index aea7450dc..eea038f91 100644 --- a/backend/src/api/endpoints/xtream_api.rs +++ b/backend/src/api/endpoints/xtream_api.rs @@ -277,7 +277,7 @@ async fn xtream_player_api_stream( let session_key = format!("{fingerprint}{virtual_id}"); let user_session = app_state .active_users - .get_user_session(&user.username, &session_key).await; + .get_and_update_user_session(&user.username, &session_key).await; let session_url = if let Some(session) = &user_session { if session.permission == UserConnectionPermission::Exhausted { diff --git a/backend/src/api/model/active_user_manager.rs b/backend/src/api/model/active_user_manager.rs index 194c3de54..ed6bade07 100644 --- a/backend/src/api/model/active_user_manager.rs +++ b/backend/src/api/model/active_user_manager.rs @@ -3,7 +3,7 @@ use crate::api::model::SharedStreamManager; use crate::model::Config; use crate::model::ProxyUserCredentials; use jsonwebtoken::get_current_timestamp; -use log::{debug, info}; +use log::{debug, error, info}; use shared::model::UserConnectionPermission; use shared::utils::{current_time_secs, default_grace_period_millis, default_grace_period_timeout_secs, sanitize_sensitive_info}; use std::collections::HashMap; @@ -11,6 +11,11 @@ use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::sync::Arc; use tokio::sync::RwLock; + +const USER_GC_TTL: u64 = 900; // 15 Min +const USER_CON_TTL: u64 = 10_800; // 3 hours +const USER_SESSION_LIMIT: usize = 50; + type ActiveUserConnectionChangeSender = tokio::sync::mpsc::Sender<(usize, usize)>; pub type ActiveUserConnectionChangeReceiver = tokio::sync::mpsc::Receiver<(usize, usize)>; @@ -18,68 +23,46 @@ macro_rules! active_user_manager_shared_impl { () => { #[inline] async fn get_active_connections(user: &Arc>>) -> usize { - user.read().await.iter().map(|(_, c)| c.connections as usize).sum() + user.read().await.iter().filter(|(_, c)| c.connections > 0).map(|(_, c)| c.connections as usize).sum() } - fn log_active_user(&self) { + async fn log_active_user(&self) { let user = Arc::clone(&self.user); let connection_change_tx = self.connection_change_tx.clone(); let is_log_user_enabled = self.is_log_user_enabled(); - tokio::spawn(async move { - let user_connection_count = Self::get_active_connections(&user).await; - let user_count = user.read().await.len(); - let _= connection_change_tx.try_send((user_count, user_connection_count)); - if is_log_user_enabled { - info!("Active Users: {user_count}, Active User Connections: {user_connection_count}"); - } - }); + let user_connection_count = Self::get_active_connections(&user).await; + let user_count = user.read().await.iter().filter(|(_, c)| c.connections > 0).count(); + let _= connection_change_tx.try_send((user_count, user_connection_count)); + if is_log_user_enabled { + info!("Active Users: {user_count}, Active User Connections: {user_connection_count}"); + } } pub async fn remove_connection(&self, addr: &str) { let username_opt = { - let user_by_addr = self.user_by_addr.read().await; - user_by_addr.get(addr).cloned() + self.user_by_addr.write().await.remove(addr) }; if let Some(username) = username_opt { - { - // Entferne addr aus user_by_addr - let mut user_by_addr = self.user_by_addr.write().await; - user_by_addr.remove(addr); - } - - let mut remove_user = false; - - { - let mut user = self.user.write().await; - if let Some(connection_data) = user.get_mut(&username) { - if connection_data.connections > 0 { - connection_data.connections -= 1; - } - - if connection_data.connections == 0 { - remove_user = true; - } else if connection_data.connections < connection_data.max_connections { - connection_data.granted_grace = false; - connection_data.grace_ts = 0; - } + let mut user = self.user.write().await; + if let Some(connection_data) = user.get_mut(&username) { + if connection_data.connections > 0 { + connection_data.connections -= 1; } - if remove_user { - user.remove(&username); + if connection_data.connections < connection_data.max_connections { + connection_data.granted_grace = false; + connection_data.grace_ts = 0; } } } self.shared_stream_manager.release_connection(addr, true).await; self.provider_manager.release_connection(addr).await; - self.log_active_user(); + self.log_active_user().await; } }; } -const USER_CON_TTL: u64 = 10_800; // 3 hours -const USER_SESSION_LIMIT: usize = 50; - fn get_grace_options(config: &Config) -> (u64, u64) { let (grace_period_millis, grace_period_timeout_secs) = config.reverse_proxy.as_ref() .and_then(|r| r.stream.as_ref()) @@ -112,9 +95,14 @@ impl Drop for UserConnectionGuard { fn drop(&mut self) { let manager = self.manager.clone(); let addr = self.addr.clone(); - tokio::spawn(async move { - manager.remove_connection(&addr).await; - }); + if let Ok(rt) = tokio::runtime::Handle::try_current() { + rt.spawn(async move { + manager.remove_connection(&addr).await; + }); + } else { + // Fallback: no runtime + error!("Runtime not available, cannot cleanly remove connection for {addr}"); + } } } @@ -135,6 +123,7 @@ struct UserConnectionData { granted_grace: bool, grace_ts: u64, sessions: Vec, + ts: u64, } impl UserConnectionData { @@ -145,6 +134,7 @@ impl UserConnectionData { granted_grace: false, grace_ts: 0, sessions: Vec::new(), + ts: current_time_secs(), } } @@ -270,7 +260,7 @@ impl ActiveUserManager { } pub async fn active_users(&self) -> usize { - self.user.read().await.len() + self.user.read().await.iter().filter(|(_, c)| c.connections > 0).count() } pub async fn active_connections(&self) -> usize { @@ -294,7 +284,7 @@ impl ActiveUserManager { user_by_addr.insert(addr.to_owned(), username.to_owned()); } - self.log_active_user(); + self.log_active_user().await; UserConnectionGuard { manager: Arc::new(self.clone_inner()), @@ -325,9 +315,9 @@ impl ActiveUserManager { #[allow(clippy::too_many_arguments)] pub async fn create_user_session(&self, user: &ProxyUserCredentials, session_token: &str, virtual_id: u32, - provider: &str, stream_url: &str, addr: &str, - connection_permission: UserConnectionPermission) -> String { - self.gc().await; + provider: &str, stream_url: &str, addr: &str, + connection_permission: UserConnectionPermission) -> String { + self.gc(); let username = user.username.clone(); let mut user_map = self.user.write().await; @@ -364,25 +354,25 @@ impl ActiveUserManager { } pub async fn update_session_addr(&self, username: &str, token: &str, addr: &str) { - let drop_addr = { - let mut user_map = self.user.write().await; - user_map.get_mut(username).and_then(|connection_data| { - connection_data.sessions.iter_mut().find_map(|session| { - if session.token == token { - let old_addr = session.addr.clone(); - addr.clone_into(&mut session.addr); - Some(old_addr) - } else { - None - } - }) + let drop_addr = { + let mut user_map = self.user.write().await; + user_map.get_mut(username).and_then(|connection_data| { + connection_data.sessions.iter_mut().find_map(|session| { + if session.token == token { + let old_addr = session.addr.clone(); + addr.clone_into(&mut session.addr); + Some(old_addr) + } else { + None + } }) - }; + }) + }; - if let Some(session_addr) = drop_addr { - self.drop_connection(&session_addr); - } + if let Some(session_addr) = drop_addr { + self.drop_connection(&session_addr); } + } fn drop_connection(&self, addr: &str) { let _ = self.close_signal_tx.send(addr.to_string()); @@ -392,47 +382,19 @@ impl ActiveUserManager { self.close_signal_tx.subscribe() } - pub async fn get_user_session(&self, username: &str, token: &str) -> Option { + pub async fn get_and_update_user_session(&self, username: &str, token: &str) -> Option { self.update_user_session(username, token).await } - // fn update_user_session(&self, username: &str, token: &str) -> Option { - // if let Some(mut entry) = self.user.get_mut(username) { - // let connection_data = &mut *entry; - // - // if connection_data.max_connections == 0 { - // return Self::find_user_session(token, &connection_data.sessions).cloned(); - // } - // - // // Separate mutable borrow of the session - // let mut found_session_index = None; - // for (i, session) in connection_data.sessions.iter().enumerate() { - // if session.token == token { - // found_session_index = Some(i); - // break; - // } - // } - // - // if let Some(index) = found_session_index { - // let session_permission = connection_data.sessions[index].permission; - // if session_permission == UserConnectionPermission::GracePeriod { - // let new_permission = self.check_connection_permission(username, connection_data); - // connection_data.sessions[index].permission = new_permission; - // } - // return Some(connection_data.sessions[index].clone()); - // } - // } - // None - // } - async fn update_user_session(&self, username: &str, token: &str) -> Option { let mut users = self.user.write().await; if let Some(connection_data) = users.get_mut(username) { + connection_data.ts = current_time_secs(); if connection_data.max_connections == 0 { return Self::find_user_session(token, &connection_data.sessions).cloned(); } - // Suche nach Index der Session + // Search for index of session if let Some(index) = connection_data .sessions .iter() @@ -448,19 +410,20 @@ impl ActiveUserManager { None } - async fn gc(&self) { + fn gc(&self) { if let Some(gc_ts) = &self.gc_ts { let ts = gc_ts.load(Ordering::Acquire); let now = current_time_secs(); - if now - ts > USER_CON_TTL { - let mut users = self.user.write().await; + if now - ts > USER_GC_TTL { + if let Ok(mut users) = self.user.try_write() { + users.retain(|_k, v| now - v.ts < USER_CON_TTL && v.connections > 0); + for connection_data in users.values_mut() { + connection_data.sessions.retain(|s| now - s.ts < USER_CON_TTL); + } - for connection_data in users.values_mut() { - connection_data.sessions.retain(|s| now - s.ts < USER_CON_TTL); + gc_ts.store(now, Ordering::Release); } - - gc_ts.store(now, Ordering::Release); } } } diff --git a/backend/src/api/model/streams/provider_stream.rs b/backend/src/api/model/streams/provider_stream.rs index a1d506aaa..2af60c4c7 100644 --- a/backend/src/api/model/streams/provider_stream.rs +++ b/backend/src/api/model/streams/provider_stream.rs @@ -72,7 +72,7 @@ pub fn create_custom_video_stream_response(config: &AppConfig, video_response: C } pub fn get_header_filter_for_item_type(item_type: PlaylistItemType) -> HeaderFilter { match item_type { - PlaylistItemType::Live | PlaylistItemType::LiveHls | PlaylistItemType::LiveDash | PlaylistItemType::LiveUnknown => { + PlaylistItemType::Live /*| PlaylistItemType::LiveHls | PlaylistItemType::LiveDash */| PlaylistItemType::LiveUnknown => { Some(Box::new(|key| key != "accept-ranges" && key != "range" && key != "content-range")) } _ => None, diff --git a/backend/src/api/model/streams/provider_stream_factory.rs b/backend/src/api/model/streams/provider_stream_factory.rs index 82c901513..404a3eb01 100644 --- a/backend/src/api/model/streams/provider_stream_factory.rs +++ b/backend/src/api/model/streams/provider_stream_factory.rs @@ -197,7 +197,7 @@ fn prepare_client( let original_headers = stream_options.get_headers(); if log_enabled!(log::Level::Debug) { - let message = format!("original_headers {original_headers:?}"); + let message = format!("original headers {original_headers:?}"); debug!("{}", sanitize_sensitive_info(&message)); } From 50a25582b00e64ca520064d1e438907276d9a3af Mon Sep 17 00:00:00 2001 From: euzu Date: Sat, 11 Oct 2025 17:07:39 +0200 Subject: [PATCH 2/2] fix for hls streaming --- CHANGELOG.md | 2 +- backend/src/api/model/active_user_manager.rs | 58 ++++++++++---------- backend/src/api/serve.rs | 3 +- 3 files changed, 30 insertions(+), 33 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index b71c44c76..3c92dacd5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,6 +1,6 @@ # Changelog # 3.1.8 (2025-10-xx) -- Fixed hls streaming caused by session removal and wrong headers. +- Fixed HLS streaming issues caused by session eviction and incorrect headers. # 3.1.7 (2025-10-10) - Added Dark/Bright theme switch diff --git a/backend/src/api/model/active_user_manager.rs b/backend/src/api/model/active_user_manager.rs index ed6bade07..2ed812b10 100644 --- a/backend/src/api/model/active_user_manager.rs +++ b/backend/src/api/model/active_user_manager.rs @@ -21,11 +21,18 @@ pub type ActiveUserConnectionChangeReceiver = tokio::sync::mpsc::Receiver<(usize macro_rules! active_user_manager_shared_impl { () => { - #[inline] + #[inline] async fn get_active_connections(user: &Arc>>) -> usize { user.read().await.iter().filter(|(_, c)| c.connections > 0).map(|(_, c)| c.connections as usize).sum() } + #[inline] + fn drop_connection(&self, addr: &str) { + if let Err(e) = self.close_signal_tx.send(addr.to_string()) { + debug!("No active receivers for close signal ({addr}): {e:?}"); + } + } + async fn log_active_user(&self) { let user = Arc::clone(&self.user); let connection_change_tx = self.connection_change_tx.clone(); @@ -56,6 +63,7 @@ macro_rules! active_user_manager_shared_impl { } } } + self.drop_connection(&addr); self.shared_stream_manager.release_connection(addr, true).await; self.provider_manager.release_connection(addr).await; self.log_active_user().await; @@ -77,6 +85,7 @@ struct ConnectionGuardUserManager { shared_stream_manager: Arc, provider_manager: Arc, connection_change_tx: ActiveUserConnectionChangeSender, + close_signal_tx: tokio::sync::broadcast::Sender, } impl ConnectionGuardUserManager { @@ -200,6 +209,7 @@ impl ActiveUserManager { shared_stream_manager: Arc::clone(&self.shared_stream_manager), provider_manager: Arc::clone(&self.provider_manager), connection_change_tx: self.connection_change_tx.clone(), + close_signal_tx: self.close_signal_tx.clone(), } } @@ -296,10 +306,6 @@ impl ActiveUserManager { self.log_active_user.load(Ordering::Relaxed) } - fn find_user_session<'a>(token: &'a str, sessions: &'a [UserSession]) -> Option<&'a UserSession> { - sessions.iter().find(|&session| session.token.eq(token)) - } - fn new_user_session(session_token: &str, virtual_id: u32, provider: &str, stream_url: &str, addr: &str, connection_permission: UserConnectionPermission) -> UserSession { UserSession { @@ -354,28 +360,18 @@ impl ActiveUserManager { } pub async fn update_session_addr(&self, username: &str, token: &str, addr: &str) { - let drop_addr = { - let mut user_map = self.user.write().await; - user_map.get_mut(username).and_then(|connection_data| { - connection_data.sessions.iter_mut().find_map(|session| { - if session.token == token { - let old_addr = session.addr.clone(); - addr.clone_into(&mut session.addr); - Some(old_addr) - } else { - None - } - }) + let mut user_map = self.user.write().await; + user_map.get_mut(username).and_then(|connection_data| { + connection_data.sessions.iter_mut().find_map(|session| { + if session.token == token { + let old_addr = session.addr.clone(); + addr.clone_into(&mut session.addr); + Some(old_addr) + } else { + None + } }) - }; - - if let Some(session_addr) = drop_addr { - self.drop_connection(&session_addr); - } - } - - fn drop_connection(&self, addr: &str) { - let _ = self.close_signal_tx.send(addr.to_string()); + }); } pub fn get_close_connection_channel(&self) -> tokio::sync::broadcast::Receiver { @@ -390,9 +386,6 @@ impl ActiveUserManager { let mut users = self.user.write().await; if let Some(connection_data) = users.get_mut(username) { connection_data.ts = current_time_secs(); - if connection_data.max_connections == 0 { - return Self::find_user_session(token, &connection_data.sessions).cloned(); - } // Search for index of session if let Some(index) = connection_data @@ -400,7 +393,12 @@ impl ActiveUserManager { .iter() .position(|s| s.token == token) { - if connection_data.sessions[index].permission == UserConnectionPermission::GracePeriod { + // Refresh session last access + connection_data.sessions[index].ts = current_time_secs(); + // Only re-evaluate permission for limited users during grace + if connection_data.max_connections > 0 + && connection_data.sessions[index].permission == UserConnectionPermission::GracePeriod + { let new_permission = self.check_connection_permission(username, connection_data); connection_data.sessions[index].permission = new_permission; } diff --git a/backend/src/api/serve.rs b/backend/src/api/serve.rs index 1951fb69a..caad8de58 100644 --- a/backend/src/api/serve.rs +++ b/backend/src/api/serve.rs @@ -68,7 +68,6 @@ pub async fn serve(listener: tokio::net::TcpListener, } } - async fn handle_connection( make_service: &mut M, signal_tx: &watch::Sender<()>, @@ -83,7 +82,7 @@ where S::Future: Send, { let Ok(tcp_stream_std) = socket.into_std() else { return; }; - tcp_stream_std.set_nonblocking(true).ok(); // this is not necessary + //tcp_stream_std.set_nonblocking(true).ok(); // this is not necessary // Configure keep alive with socket2 let sock_ref = SockRef::from(&tcp_stream_std);