Merge pull request #366 from euzu/fix/hls_streaming

fix for hls streaming
This commit is contained in:
euzu
2025-10-11 17:08:01 +02:00
committed by GitHub
8 changed files with 87 additions and 124 deletions
+3
View File
@@ -1,4 +1,7 @@
# Changelog
# 3.1.8 (2025-10-xx)
- Fixed HLS streaming issues caused by session eviction and incorrect 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)
+1 -1
View File
@@ -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 {
+1 -1
View File
@@ -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 {
+1 -1
View File
@@ -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 {
+78 -117
View File
@@ -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,75 +11,66 @@ 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)>;
macro_rules! active_user_manager_shared_impl {
() => {
#[inline]
#[inline]
async fn get_active_connections(user: &Arc<RwLock<HashMap<String, UserConnectionData>>>) -> 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) {
#[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();
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.drop_connection(&addr);
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())
@@ -94,6 +85,7 @@ struct ConnectionGuardUserManager {
shared_stream_manager: Arc<SharedStreamManager>,
provider_manager: Arc<ActiveProviderManager>,
connection_change_tx: ActiveUserConnectionChangeSender,
close_signal_tx: tokio::sync::broadcast::Sender<String>,
}
impl ConnectionGuardUserManager {
@@ -112,9 +104,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 +132,7 @@ struct UserConnectionData {
granted_grace: bool,
grace_ts: u64,
sessions: Vec<UserSession>,
ts: u64,
}
impl UserConnectionData {
@@ -145,6 +143,7 @@ impl UserConnectionData {
granted_grace: false,
grace_ts: 0,
sessions: Vec::new(),
ts: current_time_secs(),
}
}
@@ -210,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(),
}
}
@@ -270,7 +270,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 +294,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()),
@@ -306,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 {
@@ -325,9 +321,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,81 +360,45 @@ 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
}
})
})
};
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());
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
}
})
});
}
pub fn get_close_connection_channel(&self) -> tokio::sync::broadcast::Receiver<String> {
self.close_signal_tx.subscribe()
}
pub async fn get_user_session(&self, username: &str, token: &str) -> Option<UserSession> {
pub async fn get_and_update_user_session(&self, username: &str, token: &str) -> Option<UserSession> {
self.update_user_session(username, token).await
}
// fn update_user_session(&self, username: &str, token: &str) -> Option<UserSession> {
// 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<UserSession> {
let mut users = self.user.write().await;
if let Some(connection_data) = users.get_mut(username) {
if connection_data.max_connections == 0 {
return Self::find_user_session(token, &connection_data.sessions).cloned();
}
connection_data.ts = current_time_secs();
// Suche nach Index der Session
// Search for index of session
if let Some(index) = connection_data
.sessions
.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;
}
@@ -448,19 +408,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);
}
}
}
@@ -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,
@@ -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));
}
+1 -2
View File
@@ -68,7 +68,6 @@ pub async fn serve(listener: tokio::net::TcpListener,
}
}
async fn handle_connection<M, S>(
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);