diff --git a/CHANGELOG.md b/CHANGELOG.md index 2293a1969..826bd7b0a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -35,7 +35,7 @@ - Added extended debug logging for client requests and ID chain (request/action/virtual) to trace stream resolution. - Fixed xtream series/catchup lookups using the series-info virtual_id so episode requests now keep their own virtual_id/session. - Made cache storage more robust. Incomplete downloads will be deleted from cache. -- `kick_secs` added to config.yaml `web_ui` config. Default 30 seconds, if a user is kicked from the `web_ui`, they can't connect for this duration. +- `kick_secs` added to config.yaml `web_ui` config. Default 90 seconds, if a user is kicked from the `web_ui`, they can't connect for this duration. This setting is also used for sleep-timed streams. # 3.2.0 (2025-11-14) diff --git a/README.md b/README.md index b73339c14..e5d23bf01 100644 --- a/README.md +++ b/README.md @@ -417,7 +417,7 @@ log: - `content_security_policy`: configure Content-Security-Policy headers. When `enabled` is true, the default directives `default-src 'self'`, `script-src 'self' 'wasm-unsafe-eval' 'nonce-{nonce_b64}'`, and `frame-ancestors 'none'` are applied. Additional directives can be added via `custom-attributes`. Enabling CSP may block external images/logos unless allowed via directives like `img-src`. - `path` is for web_ui path like `/ui` for reverse proxy integration if necessary. - `player_server` optional, if set the server setting is used for the web-ui-player. -- `kick_secs` default 30 seconds, if a user is kicked from the `web_ui`, they can't connect for this duration. This setting is also used for sleep-timed streams. +- `kick_secs` default 90 seconds, if a user is kicked from the `web_ui`, they can't connect for this duration. This setting is also used for sleep-timed streams. - `auth` for authentication settings - `enabled` can be deactivated if `enabled` is set to `false`. If not set default is `true`. - `issuer` diff --git a/backend/src/api/endpoints/websocket_api.rs b/backend/src/api/endpoints/websocket_api.rs index 831406ffc..81b31c46e 100644 --- a/backend/src/api/endpoints/websocket_api.rs +++ b/backend/src/api/endpoints/websocket_api.rs @@ -7,7 +7,7 @@ use axum::{extract::ws::{Message, WebSocket, WebSocketUpgrade},response::IntoRes use log::{error, trace}; use shared::model::{ProtocolHandler, ProtocolHandlerMemory, ProtocolMessage, UserCommand, UserRole, WsCloseCode, PROTOCOL_VERSION}; use std::sync::Arc; -use shared::utils::{concat_path_leading_slash}; +use shared::utils::{concat_path_leading_slash, default_kick_secs}; // WebSocket upgrade handler async fn websocket_handler( @@ -316,6 +316,10 @@ async fn handle_socket(mut socket: WebSocket, app_state: Arc, auth_req async fn handle_user_action(app_state: &Arc, cmd: UserCommand) -> bool { match cmd { - UserCommand::Kick(addr, virtual_id, secs) => app_state.connection_manager.kick_connection(&addr, virtual_id, secs).await, + UserCommand::Kick(addr, virtual_id, _secs) => { + // secs could be later used for different kick configurations. Currently, we only have 1. + let kick_secs = app_state.app_config.config.load().web_ui.as_ref().map_or_else(default_kick_secs, |wc| wc.kick_secs); + app_state.connection_manager.kick_connection(&addr, virtual_id, kick_secs).await + } } } diff --git a/backend/src/api/model/active_user_manager.rs b/backend/src/api/model/active_user_manager.rs index 90ca91a3c..db19a20ab 100644 --- a/backend/src/api/model/active_user_manager.rs +++ b/backend/src/api/model/active_user_manager.rs @@ -512,6 +512,10 @@ impl ActiveUserManager { } } + pub async fn get_username_for_addr(&self, addr: &SocketAddr) -> Option{ + self.connections.read().await.key_by_addr.get(addr).cloned() + } + fn gc(&self) { if let Some(gc_ts) = &self.gc_ts { let ts = gc_ts.load(Ordering::Acquire); diff --git a/backend/src/api/model/connection_manager.rs b/backend/src/api/model/connection_manager.rs index 64277d851..f060079e3 100644 --- a/backend/src/api/model/connection_manager.rs +++ b/backend/src/api/model/connection_manager.rs @@ -1,10 +1,12 @@ -use std::borrow::Cow; -use std::net::SocketAddr; use crate::api::model::{ActiveProviderManager, ActiveUserManager, CustomVideoStreamType, EventManager, EventMessage, ProviderHandle, SharedStreamManager}; -use std::sync::Arc; +use crate::auth::Fingerprint; +use crate::utils::debug_if_enabled; use log::{debug, warn}; use shared::model::{ActiveUserConnectionChange, StreamChannel, VirtualId}; -use crate::auth::Fingerprint; +use shared::utils::sanitize_sensitive_info; +use std::borrow::Cow; +use std::net::SocketAddr; +use std::sync::Arc; pub struct ConnectionManager { pub user_manager: Arc, @@ -28,7 +30,7 @@ impl ConnectionManager { provider_manager: Arc::clone(provider_manager), shared_stream_manager: Arc::clone(shared_stream_manager), event_manager: Arc::clone(event_manager), - close_socket_signal_tx + close_socket_signal_tx, } } @@ -37,6 +39,8 @@ impl ConnectionManager { } pub async fn kick_connection(&self, addr: &SocketAddr, virtual_id: VirtualId, block_secs: u64) -> bool { + debug_if_enabled!("User {} kicked for stream with virtual_id {virtual_id} for {block_secs} seconds with addr {}.", + self.user_manager.get_username_for_addr(addr).await.unwrap_or_default(), sanitize_sensitive_info(&addr.to_string())); if block_secs > 0 { self.user_manager.block_user_for_stream(addr, virtual_id, block_secs).await; } @@ -72,12 +76,12 @@ impl ConnectionManager { #[allow(clippy::too_many_arguments)] pub async fn update_connection(&self, username: &str, max_connections: u32, fingerprint: &Fingerprint, - provider: &str, stream_channel: StreamChannel, user_agent: Cow<'_, str>, session_token: Option<&str>) { + provider: &str, stream_channel: StreamChannel, user_agent: Cow<'_, str>, session_token: Option<&str>) { if let Some(stream_info) = self.user_manager.update_connection(username, max_connections, fingerprint, provider, stream_channel, user_agent, session_token).await { self.event_manager.send_event(EventMessage::ActiveUser(ActiveUserConnectionChange::Updated(stream_info))); } else { warn!("Failed to register connection for user {username} at {}; disconnecting client", fingerprint.addr); - let _ = self.kick_connection(&fingerprint.addr,0, 0).await; + let _ = self.kick_connection(&fingerprint.addr, 0, 0).await; } } @@ -86,9 +90,8 @@ impl ConnectionManager { // } pub async fn update_stream_detail(&self, addr: &SocketAddr, video_type: CustomVideoStreamType) { - if let Some(stream_info) = self.user_manager.update_stream_detail(addr, video_type).await { - self.event_manager.send_event(EventMessage::ActiveUser(ActiveUserConnectionChange::Updated(stream_info))); - } + if let Some(stream_info) = self.user_manager.update_stream_detail(addr, video_type).await { + self.event_manager.send_event(EventMessage::ActiveUser(ActiveUserConnectionChange::Updated(stream_info))); + } } - } diff --git a/backend/src/api/model/streams/timed_client_stream.rs b/backend/src/api/model/streams/timed_client_stream.rs index 278a53b15..16fae771d 100644 --- a/backend/src/api/model/streams/timed_client_stream.rs +++ b/backend/src/api/model/streams/timed_client_stream.rs @@ -7,8 +7,9 @@ use std::sync::Arc; use std::task::Poll; use std::time::{Duration, Instant}; use shared::model::VirtualId; -use shared::utils::default_kick_secs; +use shared::utils::{default_kick_secs, sanitize_sensitive_info}; use crate::api::model::{AppState, BoxedProviderStream}; +use crate::utils::debug_if_enabled; pub struct TimedClientStream { inner: BoxedProviderStream, @@ -30,11 +31,13 @@ impl Stream for TimedClientStream { fn poll_next(mut self: Pin<&mut Self>,cx: &mut std::task::Context<'_>,) -> Poll> { if Instant::now() >= self.deadline { let kick_secs = self.app_state.app_config.config.load().web_ui.as_ref().map_or_else(default_kick_secs, |wc| wc.kick_secs); - let user_manager = Arc::clone(&self.app_state.active_users); + let connection_manager = Arc::clone(&self.app_state.connection_manager); let addr = self.addr; let virtual_id = self.virtual_id; + debug_if_enabled!("TimedClient stream exceeds time limit. Closing stream with virtual_id {virtual_id} for addr: {}", + sanitize_sensitive_info(&addr.to_string())); tokio::spawn(async move { - user_manager.block_user_for_stream(&addr, virtual_id, kick_secs).await; + connection_manager.kick_connection(&addr, virtual_id, kick_secs).await; }); return Poll::Ready(None); } diff --git a/backend/src/processing/processor/playlist.rs b/backend/src/processing/processor/playlist.rs index 198a23647..38080b07e 100644 --- a/backend/src/processing/processor/playlist.rs +++ b/backend/src/processing/processor/playlist.rs @@ -340,9 +340,7 @@ async fn process_source(client: &reqwest::Client, cfg: Arc, source_id errors.append(&mut error_list); errors.append(&mut tvguide_errors); let group_count = playlistgroups.len(); - let channel_count = playlistgroups.iter() - .map(|group| group.channels.len()) - .sum(); + let channel_count = playlistgroups.iter().map(|group| group.channels.len()).sum(); let input_name = &input.name; if playlistgroups.is_empty() { info!("Source is empty {input_name}"); diff --git a/shared/src/utils/default_utils.rs b/shared/src/utils/default_utils.rs index 54862ecbf..ae1d11121 100644 --- a/shared/src/utils/default_utils.rs +++ b/shared/src/utils/default_utils.rs @@ -23,4 +23,4 @@ pub fn default_secret() -> String { out.iter().map(|b| format!("{:02X}", b)).collect() } -pub const fn default_kick_secs() -> u64 { 30 } \ No newline at end of file +pub const fn default_kick_secs() -> u64 { 90 } \ No newline at end of file