From 4281cbc90def1ecf224988aee7f12d967ad7d7f1 Mon Sep 17 00:00:00 2001 From: euzu Date: Sat, 13 Dec 2025 11:32:55 +0100 Subject: [PATCH 1/5] kick debug log added --- CHANGELOG.md | 2 +- README.md | 2 +- backend/src/api/endpoints/websocket_api.rs | 8 ++++-- backend/src/api/model/active_user_manager.rs | 4 +++ backend/src/api/model/connection_manager.rs | 25 +++++++++++-------- .../api/model/streams/timed_client_stream.rs | 9 ++++--- backend/src/processing/processor/playlist.rs | 4 +-- shared/src/utils/default_utils.rs | 2 +- 8 files changed, 34 insertions(+), 22 deletions(-) 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 From 60e74f38a9dc11ce930942503ebfc65c9f8faffe Mon Sep 17 00:00:00 2001 From: euzu Date: Sat, 13 Dec 2025 13:54:01 +0100 Subject: [PATCH 2/5] force reconnect removed --- backend/src/api/api_utils.rs | 8 ++------ backend/src/model/config/stream.rs | 3 --- frontend/public/assets/i18n/en.json | 1 - .../app/components/config/reverse_proxy_config_view.rs | 4 ---- shared/src/model/config/stream.rs | 4 ---- 5 files changed, 2 insertions(+), 18 deletions(-) diff --git a/backend/src/api/api_utils.rs b/backend/src/api/api_utils.rs index d355f2614..254f800a9 100644 --- a/backend/src/api/api_utils.rs +++ b/backend/src/api/api_utils.rs @@ -210,7 +210,6 @@ pub fn get_user_target<'a>( pub struct StreamOptions { pub stream_retry: bool, - pub stream_force_retry_secs: u32, pub buffer_enabled: bool, pub buffer_size: usize, pub pipe_provider_stream: bool, @@ -220,7 +219,6 @@ pub struct StreamOptions { /// /// This function retrieves streaming-related settings from the `AppState`: /// - `stream_retry`: whether retrying the stream is enabled, -/// - `stream_force_retry_secs`: the number of seconds to wait before a forced retry, /// - `buffer_enabled`: whether stream buffering is enabled, /// - `buffer_size`: the size of the stream buffer. /// @@ -236,21 +234,20 @@ pub struct StreamOptions { /// /// Returns a `StreamOptions` instance with the resolved configuration. fn get_stream_options(app_state: &AppState) -> StreamOptions { - let (stream_retry, stream_force_retry_secs, buffer_enabled, buffer_size) = app_state + let (stream_retry, buffer_enabled, buffer_size) = app_state .app_config .config .load() .reverse_proxy .as_ref() .and_then(|reverse_proxy| reverse_proxy.stream.as_ref()) - .map_or((false, 0, false, 0), |stream| { + .map_or((false, false, 0), |stream| { let (buffer_enabled, buffer_size) = stream .buffer .as_ref() .map_or((false, 0), |buffer| (buffer.enabled, buffer.size)); ( stream.retry, - stream.forced_retry_interval_secs, buffer_enabled, buffer_size, ) @@ -258,7 +255,6 @@ fn get_stream_options(app_state: &AppState) -> StreamOptions { let pipe_provider_stream = !stream_retry && !buffer_enabled; StreamOptions { stream_retry, - stream_force_retry_secs, buffer_enabled, buffer_size, pipe_provider_stream, diff --git a/backend/src/model/config/stream.rs b/backend/src/model/config/stream.rs index 25848c223..e3bbaa38c 100644 --- a/backend/src/model/config/stream.rs +++ b/backend/src/model/config/stream.rs @@ -34,7 +34,6 @@ pub struct StreamConfig { pub buffer: Option, pub grace_period_millis: u64, pub grace_period_timeout_secs: u64, - pub forced_retry_interval_secs: u32, pub throttle_str: Option, pub throttle_kbps: u64, pub shared_burst_buffer_mb: u64, @@ -48,7 +47,6 @@ impl From<&StreamConfigDto> for StreamConfig { buffer: dto.buffer.as_ref().map(Into::into), grace_period_millis: dto.grace_period_millis, grace_period_timeout_secs: dto.grace_period_timeout_secs, - forced_retry_interval_secs: dto.forced_retry_interval_secs, throttle_str: dto.throttle.clone(), throttle_kbps: dto.throttle.as_ref().map_or(0u64, |throttle| parse_to_kbps(throttle).unwrap_or(0u64)), shared_burst_buffer_mb: dto.shared_burst_buffer_mb, @@ -63,7 +61,6 @@ impl From<&StreamConfig> for StreamConfigDto { buffer: instance.buffer.as_ref().map(Into::into), grace_period_millis: instance.grace_period_millis, grace_period_timeout_secs: instance.grace_period_timeout_secs, - forced_retry_interval_secs: instance.forced_retry_interval_secs, throttle: instance.throttle_str.clone(), throttle_kbps: instance.throttle_kbps, shared_burst_buffer_mb: instance.shared_burst_buffer_mb, diff --git a/frontend/public/assets/i18n/en.json b/frontend/public/assets/i18n/en.json index 725449cff..9a295eb1c 100644 --- a/frontend/public/assets/i18n/en.json +++ b/frontend/public/assets/i18n/en.json @@ -235,7 +235,6 @@ "RETRY": "Retry", "GRACE_PERIOD_MILLIS": "Grace Period (millis)", "GRACE_PERIOD_TIMEOUT_SECS": "Grace Period Timeout (secs)", - "FORCED_RETRY_INTERVAL_SECS": "Forced Retry Interval (secs)", "THROTTLE_KBPS": "Throttle kbps", "FRIENDLY_NAME": "Friendly Name", "MANUFACTURER": "Manufacturer", diff --git a/frontend/src/app/components/config/reverse_proxy_config_view.rs b/frontend/src/app/components/config/reverse_proxy_config_view.rs index b1c58ea9c..62b51a426 100644 --- a/frontend/src/app/components/config/reverse_proxy_config_view.rs +++ b/frontend/src/app/components/config/reverse_proxy_config_view.rs @@ -22,7 +22,6 @@ const LABEL_RETRY: &str = "LABEL.RETRY"; const LABEL_THROTTLE: &str = "LABEL.THROTTLE"; const LABEL_GRACE_PERIOD_MILLIS: &str = "LABEL.GRACE_PERIOD_MILLIS"; const LABEL_GRACE_PERIOD_TIMEOUT_SECS: &str = "LABEL.GRACE_PERIOD_TIMEOUT_SECS"; -//const LABEL_FORCED_RETRY_INTERVAL_SECS: &str = "LABEL.FORCED_RETRY_INTERVAL_SECS"; const LABEL_THROTTLE_KBPS: &str = "LABEL.THROTTLE_KBPS"; const LABEL_RATE_LIMIT: &str = "LABEL.RATE_LIMIT"; @@ -84,7 +83,6 @@ generate_form_reducer!( Throttle => throttle: Option, GracePeriodMillis => grace_period_millis: u64, GracePeriodTimeoutSecs => grace_period_timeout_secs: u64, - //ForcedRetryIntervalSecs => forced_retry_interval_secs: u32, ThrottleKbps => throttle_kbps: u64, SharedBurstBufferMb => shared_burst_buffer_mb: u64, } @@ -243,7 +241,6 @@ pub fn ReverseProxyConfigView() -> Html { { config_field_optional!(stream_state.form, translate.t(LABEL_THROTTLE), throttle) } { config_field!(stream_state.form, translate.t(LABEL_GRACE_PERIOD_MILLIS), grace_period_millis) } { config_field!(stream_state.form, translate.t(LABEL_GRACE_PERIOD_TIMEOUT_SECS), grace_period_timeout_secs) } - //{ config_field!(stream_state.form, translate.t(LABEL_FORCED_RETRY_INTERVAL_SECS), forced_retry_interval_secs) } { config_field!(stream_state.form, translate.t(LABEL_THROTTLE_KBPS), throttle_kbps) } { config_field!(stream_state.form, translate.t(LABEL_SHARED_BURST_BUFFER_MB), shared_burst_buffer_mb) } @@ -380,7 +377,6 @@ pub fn ReverseProxyConfigView() -> Html { { edit_field_text_option!(stream_state, translate.t(LABEL_THROTTLE), throttle, StreamConfigFormAction::Throttle) } { edit_field_number_u64!(stream_state, translate.t(LABEL_GRACE_PERIOD_MILLIS), grace_period_millis, StreamConfigFormAction::GracePeriodMillis) } { edit_field_number_u64!(stream_state, translate.t(LABEL_GRACE_PERIOD_TIMEOUT_SECS), grace_period_timeout_secs, StreamConfigFormAction::GracePeriodTimeoutSecs) } - //{ edit_field_number!(stream_state, translate.t(LABEL_FORCED_RETRY_INTERVAL_SECS), forced_retry_interval_secs, StreamConfigFormAction::ForcedRetryIntervalSecs) } { edit_field_number_u64!(stream_state, translate.t(LABEL_THROTTLE_KBPS), throttle_kbps, StreamConfigFormAction::ThrottleKbps) } { edit_field_number_u64!(stream_state, translate.t(LABEL_SHARED_BURST_BUFFER_MB), shared_burst_buffer_mb, StreamConfigFormAction::SharedBurstBufferMb) } diff --git a/shared/src/model/config/stream.rs b/shared/src/model/config/stream.rs index bc5288ad7..d27a2da54 100644 --- a/shared/src/model/config/stream.rs +++ b/shared/src/model/config/stream.rs @@ -40,8 +40,6 @@ pub struct StreamConfigDto { pub grace_period_millis: u64, #[serde(default = "default_grace_period_timeout_secs")] pub grace_period_timeout_secs: u64, - #[serde(default)] - pub forced_retry_interval_secs: u32, #[serde(default, skip)] pub throttle_kbps: u64, #[serde(default = "default_shared_burst_buffer_mb")] @@ -56,7 +54,6 @@ impl Default for StreamConfigDto { throttle: None, grace_period_millis: default_grace_period_millis(), grace_period_timeout_secs: default_grace_period_timeout_secs(), - forced_retry_interval_secs: 0, throttle_kbps: 0, shared_burst_buffer_mb: default_shared_burst_buffer_mb(), } @@ -71,7 +68,6 @@ impl StreamConfigDto { && (self.throttle.is_none() || self.throttle.as_ref().is_some_and(|t| t.is_empty())) && self.grace_period_millis == empty.grace_period_millis && self.grace_period_timeout_secs == empty.grace_period_timeout_secs - && self.forced_retry_interval_secs == empty.forced_retry_interval_secs && self.throttle_kbps == empty.throttle_kbps && self.shared_burst_buffer_mb == default_shared_burst_buffer_mb() } From 80b7183d104fa33ed732107f1e17d8e8e073832a Mon Sep 17 00:00:00 2001 From: euzu Date: Sat, 13 Dec 2025 15:07:47 +0100 Subject: [PATCH 3/5] replaced serde_json Value with RawValue to reduce memory usage --- backend/src/model/xtream.rs | 19 ++++-- backend/src/processing/processor/trakt.rs | 15 ++-- .../src/processing/processor/xtream_vod.rs | 17 ++--- backend/src/repository/strm_repository.rs | 5 +- shared/src/model/playlist.rs | 68 +++++++++---------- 5 files changed, 69 insertions(+), 55 deletions(-) diff --git a/backend/src/model/xtream.rs b/backend/src/model/xtream.rs index ef747262e..f162bfaa5 100644 --- a/backend/src/model/xtream.rs +++ b/backend/src/model/xtream.rs @@ -7,6 +7,7 @@ use shared::model::{ClusterFlags, PlaylistEntry, XtreamCluster}; use shared::utils::{deserialize_as_option_string, deserialize_as_string, deserialize_as_string_array, deserialize_number_from_string, get_non_empty_str, opt_string_or_number_u32, string_default_on_null, string_or_number_f64, string_or_number_u32}; use std::iter::FromIterator; +use serde_json::value::RawValue; #[derive(Debug, Default)] pub struct XtreamLoginInfo { @@ -153,7 +154,7 @@ impl XtreamStream { self.stream_id.unwrap_or_else(|| self.series_id.unwrap_or(0)) } - pub fn get_additional_properties(&self) -> Option { + pub fn get_additional_properties(&self) -> Option> { let mut result = Map::new(); if let Some(bdpath) = self.backdrop_path.as_ref() { if !bdpath.is_empty() { @@ -183,7 +184,12 @@ impl XtreamStream { add_opt_i64_property_if_exists!(result, self.tv_archive, "tv_archive"); add_opt_i64_property_if_exists!(result, self.tv_archive_duration, "tv_archive_duration"); add_opt_i64_property_if_exists!(result, self.is_adult, "is_adult"); - if result.is_empty() { None } else { Some(Value::Object(result)) } + if result.is_empty() { + None + } else { + let s = serde_json::to_string(&Value::Object(result)).ok()?; + RawValue::from_string(s).ok() + } } } @@ -403,7 +409,7 @@ pub struct XtreamSeriesInfo { } impl XtreamSeriesInfoEpisode { - pub fn get_additional_properties(&self, series_info: &XtreamSeriesInfo) -> Option { + pub fn get_additional_properties(&self, series_info: &XtreamSeriesInfo) -> Option> { let mut result = Map::new(); let info = series_info.info.as_ref(); let bdpath = info.and_then(|i| i.backdrop_path.as_ref()); @@ -437,7 +443,12 @@ impl XtreamSeriesInfoEpisode { } } - if result.is_empty() { None } else { Some(Value::Object(result)) } + if result.is_empty() { + None + } else { + let s = serde_json::to_string(&Value::Object(result)).ok()?; + RawValue::from_string(s).ok() + } } } diff --git a/backend/src/processing/processor/trakt.rs b/backend/src/processing/processor/trakt.rs index 539f236da..b2083a318 100644 --- a/backend/src/processing/processor/trakt.rs +++ b/backend/src/processing/processor/trakt.rs @@ -40,14 +40,13 @@ fn is_compatible_content_type(cluster: XtreamCluster, content_type: TraktContent /// Extract TMDB ID from playlist item fn extract_tmdb_id_from_playlist_item(item: &PlaylistItem) -> Option { if let Some(additional_props) = &item.header.additional_properties { - if let Some(props_str) = additional_props.as_str() { - if let Ok(props) = serde_json::from_str::>(props_str) { - return props - .get("tmdb_id") - .and_then(get_u32_from_serde_value) - .filter(|&id| id != 0) - .or_else(|| props.get("tmdb").and_then(get_u32_from_serde_value)); - } + let props_str = additional_props.get(); + if let Ok(props) = serde_json::from_str::>(props_str) { + return props + .get("tmdb_id") + .and_then(get_u32_from_serde_value) + .filter(|&id| id != 0) + .or_else(|| props.get("tmdb").and_then(get_u32_from_serde_value)); } } None diff --git a/backend/src/processing/processor/xtream_vod.rs b/backend/src/processing/processor/xtream_vod.rs index 26810540f..badfbc28f 100644 --- a/backend/src/processing/processor/xtream_vod.rs +++ b/backend/src/processing/processor/xtream_vod.rs @@ -8,11 +8,12 @@ use crate::repository::xtream_repository::{write_vod_info_to_wal_file, xtream_up use shared::error::{notify_err}; use crate::processing::processor::{handle_error, handle_error_and_return, create_resolve_options_function_for_xtream_target}; use shared::utils::{get_u32_from_serde_value, get_u64_from_serde_value, get_string_from_serde_value}; -use serde_json::{from_str, Map, Value}; +use serde_json::{Map, Value}; use std::collections::{HashMap, HashSet}; use std::io::{Write}; use std::time::Instant; use log::{info, log_enabled, Level}; +use serde_json::value::RawValue; use crate::utils; use crate::processing::processor::xtream::normalize_json_content; use crate::utils::IO_BUFFER_SIZE; @@ -123,13 +124,13 @@ pub async fn playlist_resolve_vod(app_config: &AppConfig, client: &reqwest::Clie // Update in-memory playlist items with the newly fetched vod info. // This makes the data available for subsequent processing steps like STRM export. - pli.header.additional_properties = from_str::>(normalized_str).ok().and_then(|info_doc| { - info_doc.get("info").cloned().map(|info_content| { - let mut wrapped_info = Map::new(); - wrapped_info.insert("info".to_string(), info_content); - Value::Object(wrapped_info) - }) - }); + if let Ok(value) = serde_json::from_str::(normalized_str) { + if let Some(info_content) = value.get("info") { + if let Ok(raw_info) = serde_json::to_string(info_content) { + pli.header.additional_properties = serde_json::from_str::>(&raw_info).ok(); + } + } + } } } } diff --git a/backend/src/repository/strm_repository.rs b/backend/src/repository/strm_repository.rs index 0cfdeb873..d1f63eb27 100644 --- a/backend/src/repository/strm_repository.rs +++ b/backend/src/repository/strm_repository.rs @@ -749,7 +749,10 @@ async fn prepare_strm_files( let quality_string = if strm_target_output.add_quality_to_filename { pli.header.additional_properties .as_ref() - .and_then(|props| props.get("info")) + .and_then(|raw| { + let v: serde_json::Value = serde_json::from_str(raw.get()).ok()?; + v.get("info").cloned() + }).as_ref() .and_then(MediaQuality::from_ffprobe_info) .map_or_else(String::new, |quality| { let formatted = quality.format_for_filename(separator); diff --git a/shared/src/model/playlist.rs b/shared/src/model/playlist.rs index 1d91befa5..099223bca 100644 --- a/shared/src/model/playlist.rs +++ b/shared/src/model/playlist.rs @@ -5,7 +5,7 @@ use serde_json::{Map, Value}; use std::borrow::Cow; use std::fmt::{Display, Formatter}; use std::str::FromStr; - +use serde_json::value::RawValue; // https://de.wikipedia.org/wiki/M3U // https://siptv.eu/howto/playlist.html @@ -166,7 +166,7 @@ pub struct PlaylistItemHeader { pub url: String, pub epg_channel_id: Option, pub xtream_cluster: XtreamCluster, - pub additional_properties: Option, + pub additional_properties: Option>, #[serde(default)] pub item_type: PlaylistItemType, #[serde(default)] @@ -192,32 +192,33 @@ impl PlaylistItemHeader { } } - pub fn get_additional_property(&self, field: &str) -> Option<&Value> { - self.additional_properties.as_ref().and_then(|v| match v { - Value::Object(map) => { - map.get(field) - } + pub fn get_additional_property(&self, field: &str) -> Option { + let raw = self.additional_properties.as_ref()?; + let value: Value = serde_json::from_str(raw.get()).ok()?; + + match value { + Value::Object(map) => map.get(field).cloned(), _ => None, - }) + } } pub fn get_additional_property_as_u32(&self, field: &str) -> Option { match self.get_additional_property(field) { - Some(value) => get_u32_from_serde_value(value), + Some(value) => get_u32_from_serde_value(&value), None => None } } pub fn get_additional_property_as_u64(&self, field: &str) -> Option { match self.get_additional_property(field) { - Some(value) => get_u64_from_serde_value(value), + Some(value) => get_u64_from_serde_value(&value), None => None } } pub fn get_additional_property_as_str(&self, field: &str) -> Option { match self.get_additional_property(field) { - Some(value) => get_string_from_serde_value(value), + Some(value) => get_string_from_serde_value(&value), None => None } } @@ -607,23 +608,22 @@ impl PlaylistItem { pub fn to_xtream(&self) -> XtreamPlaylistItem { let header = &self.header; let provider_id = header.id.parse::().unwrap_or_default(); - let mut additional_properties = None; + let mut additional_properties: Option = None; + if header.xtream_cluster != XtreamCluster::Live { let add_ext = match header.get_additional_property("container_extension") { None => true, - Some(ext) => ext.as_str().is_none_or(str::is_empty) + Some(ext) => ext.as_str().is_none_or(str::is_empty), }; + if add_ext { if let Some(cont_ext) = extract_extension_from_url(&header.url) { - let ext = if let Some(stripped) = cont_ext.strip_prefix('.') { stripped } else { cont_ext }; - let mut result = match header.additional_properties.as_ref() { + let ext = cont_ext.strip_prefix('.').unwrap_or(cont_ext); + // parse RawValue into Value on-demand + let mut result: Map = match header.additional_properties.as_ref() { None => Map::new(), Some(props) => { - if let Value::Object(map) = props { - map.clone() - } else { - Map::new() - } + serde_json::from_str(props.get()).unwrap_or_else(|_| Map::new()) } }; result.insert("container_extension".to_string(), Value::String(ext.to_string())); @@ -631,14 +631,12 @@ impl PlaylistItem { } } } + if additional_properties.is_none() { additional_properties = header.additional_properties.as_ref().and_then(|props| { - serde_json::to_string(props).ok() + serde_json::to_string(props.get()).ok() }); } - // let additional_properties = header.additional_properties.as_ref().and_then(|props| { - // serde_json::to_string(props).ok() - // }); XtreamPlaylistItem { virtual_id: header.virtual_id, @@ -663,33 +661,35 @@ impl PlaylistItem { pub fn to_common(&self) -> CommonPlaylistItem { let header = &self.header; - let mut additional_properties = None; + let mut additional_properties: Option = None; + if header.xtream_cluster != XtreamCluster::Live { let add_ext = match header.get_additional_property("container_extension") { None => true, - Some(ext) => ext.as_str().is_none_or(str::is_empty) + Some(ext) => ext.as_str().is_none_or(str::is_empty), }; + if add_ext { if let Some(cont_ext) = extract_extension_from_url(&header.url) { - let ext = if let Some(stripped) = cont_ext.strip_prefix('.') { stripped } else { cont_ext }; - let mut result = match header.additional_properties.as_ref() { + let ext = cont_ext.strip_prefix('.').unwrap_or(cont_ext); + + // Parse RawValue on-demand + let mut result: Map = match header.additional_properties.as_ref() { None => Map::new(), Some(props) => { - if let Value::Object(map) = props { - map.clone() - } else { - Map::new() - } + serde_json::from_str(props.get()).unwrap_or_else(|_| Map::new()) } }; + result.insert("container_extension".to_string(), Value::String(ext.to_string())); additional_properties = serde_json::to_string(&Value::Object(result)).ok(); } } } + if additional_properties.is_none() { additional_properties = header.additional_properties.as_ref().and_then(|props| { - serde_json::to_string(props).ok() + serde_json::to_string(props.get()).ok() }); } From 22e6f74bae393a4b176e2f0f06cd6fa72c779963 Mon Sep 17 00:00:00 2001 From: euzu Date: Sat, 13 Dec 2025 15:33:04 +0100 Subject: [PATCH 4/5] replaced serde_json Value with RawValue to reduce memory usage --- CHANGELOG.md | 2 +- backend/src/api/api_utils.rs | 1 - backend/src/api/model/active_user_manager.rs | 17 ++++++---- backend/src/repository/strm_repository.rs | 35 +++++++++++++------- 4 files changed, 34 insertions(+), 21 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 826bd7b0a..5419394b9 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 90 seconds, if a user is kicked from the `web_ui`, they can't connect for this duration. +- `kick_secs` added to config.yml `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/backend/src/api/api_utils.rs b/backend/src/api/api_utils.rs index 254f800a9..8c6c7a2ad 100644 --- a/backend/src/api/api_utils.rs +++ b/backend/src/api/api_utils.rs @@ -224,7 +224,6 @@ pub struct StreamOptions { /// /// If the reverse proxy or stream settings are not defined, default values are used: /// - retry: `false` -/// - forced retry interval: `0` /// - buffering: `false` /// - buffer size: `0` /// diff --git a/backend/src/api/model/active_user_manager.rs b/backend/src/api/model/active_user_manager.rs index db19a20ab..6399e6290 100644 --- a/backend/src/api/model/active_user_manager.rs +++ b/backend/src/api/model/active_user_manager.rs @@ -503,16 +503,19 @@ impl ActiveUserManager { } pub async fn block_user_for_stream(&self, addr: &SocketAddr, virtual_id: VirtualId, blocked_secs: u64) { - let mut connections = self.connections.write().await; - let now = current_time_secs(); - connections.kicked.retain(|_, (expires_at, _)| *expires_at > now); - if let Some(username) = connections.key_by_addr.get(addr).cloned() { - let expires_at = now + blocked_secs.clamp(1,86_400); // max 1 day - connections.kicked.insert(username, (expires_at, virtual_id)); + let block_for_secs = blocked_secs.clamp(0, 86_400); // max 1 day; + if block_for_secs > 0 { + let mut connections = self.connections.write().await; + let now = current_time_secs(); + connections.kicked.retain(|_, (expires_at, _)| *expires_at > now); + if let Some(username) = connections.key_by_addr.get(addr).cloned() { + let expires_at = now + block_for_secs; + connections.kicked.insert(username, (expires_at, virtual_id)); + } } } - pub async fn get_username_for_addr(&self, addr: &SocketAddr) -> Option{ + pub async fn get_username_for_addr(&self, addr: &SocketAddr) -> Option { self.connections.read().await.key_by_addr.get(addr).cloned() } diff --git a/backend/src/repository/strm_repository.rs b/backend/src/repository/strm_repository.rs index d1f63eb27..01ce25600 100644 --- a/backend/src/repository/strm_repository.rs +++ b/backend/src/repository/strm_repository.rs @@ -1,28 +1,37 @@ -use shared::error::{create_tuliprox_error_result, info_err}; -use shared::error::{TuliproxError, TuliproxErrorKind}; +// Import the new MediaQuality struct +use crate::model::MediaQuality; +use crate::model::XtreamSeriesEpisode; use crate::model::{ApiProxyServerInfo, AppConfig, ProxyUserCredentials}; use crate::model::{ConfigTarget, StrmTargetOutput}; -use crate::model::XtreamSeriesEpisode; use crate::repository::bplustree::BPlusTree; use crate::repository::storage::{ensure_target_storage_path, get_input_storage_path}; use crate::repository::storage_const; use crate::repository::xtream_repository::{xtream_get_record_file_path, InputVodInfoRecord}; -use shared::utils::{extract_extension_from_url, hash_bytes, hash_string_as_hex, truncate_string, ExportStyleConfig, CONSTANTS}; +use crate::utils; use crate::utils::{async_file_reader, async_file_writer, normalize_string_path, truncate_filename, FileReadGuard, IO_BUFFER_SIZE}; use chrono::Datelike; use filetime::{set_file_times, FileTime}; use log::{error, trace}; use regex::Regex; use serde::Serialize; +use shared::error::{create_tuliprox_error_result, info_err}; +use shared::error::{TuliproxError, TuliproxErrorKind}; +use shared::model::{ClusterFlags, FieldGetAccessor, PlaylistGroup, PlaylistItem, PlaylistItemType, StrmExportStyle, UUIDType}; +use shared::utils::{extract_extension_from_url, hash_bytes, hash_string_as_hex, truncate_string, ExportStyleConfig, CONSTANTS}; use std::collections::{HashMap, HashSet, VecDeque}; use std::path::{Path, PathBuf}; use std::sync::Arc; use tokio::fs::{create_dir_all, remove_dir, remove_file, File}; use tokio::io::{AsyncBufReadExt, AsyncReadExt, AsyncWriteExt}; -use shared::model::{ClusterFlags, FieldGetAccessor, PlaylistGroup, PlaylistItem, PlaylistItemType, StrmExportStyle, UUIDType}; -use crate::utils; -// Import the new MediaQuality struct -use crate::model::{MediaQuality}; + +use serde::Deserialize; +use serde_json::value::RawValue as JsonRawValue; + +#[derive(Deserialize)] +struct AdditionalProps<'a> { + #[serde(borrow)] + info: Option<&'a JsonRawValue>, +} /// Sanitizes a string to be safe for use as a file or directory name by /// following a strict "allow-list" approach and discarding invalid characters. @@ -750,9 +759,11 @@ async fn prepare_strm_files( pli.header.additional_properties .as_ref() .and_then(|raw| { - let v: serde_json::Value = serde_json::from_str(raw.get()).ok()?; - v.get("info").cloned() - }).as_ref() + let props: AdditionalProps<'_> = serde_json::from_str(raw.get()).ok()?; + let info = props.info?; + serde_json::from_str::(info.get()).ok() + }) + .as_ref() .and_then(MediaQuality::from_ffprobe_info) .map_or_else(String::new, |quality| { let formatted = quality.format_for_filename(separator); @@ -869,7 +880,7 @@ pub async fn write_strm_playlist( for strm_file in strm_files { // file paths let output_path = truncate_filename(&root_path.join(&strm_file.dir_path), 255); - let file_path = output_path.join(format!("{}.strm", truncate_string(&strm_file.file_name, 250))); + let file_path = output_path.join(format!("{}.strm", truncate_string(&strm_file.file_name, 250))); let file_exists = file_path.exists(); let relative_file_path = get_relative_path_str(&file_path, &root_path); From e9766bcdfebd19b4c9ac3874eb232141f0692109 Mon Sep 17 00:00:00 2001 From: euzu Date: Sat, 13 Dec 2025 15:33:35 +0100 Subject: [PATCH 5/5] increased version --- Cargo.lock | 6 +++--- backend/Cargo.toml | 2 +- frontend/Cargo.toml | 4 ++-- shared/Cargo.toml | 2 +- 4 files changed, 7 insertions(+), 7 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index ba2adb119..fde82ca19 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1096,7 +1096,7 @@ dependencies = [ [[package]] name = "frontend" -version = "3.2.19" +version = "3.2.20" dependencies = [ "anyhow", "base64", @@ -3793,7 +3793,7 @@ dependencies = [ [[package]] name = "shared" -version = "3.2.19" +version = "3.2.20" dependencies = [ "base64", "bitflags 2.10.0", @@ -4356,7 +4356,7 @@ checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b" [[package]] name = "tuliprox" -version = "3.2.19" +version = "3.2.20" dependencies = [ "arc-swap", "async-compression", diff --git a/backend/Cargo.toml b/backend/Cargo.toml index 80ab2d728..4f1702942 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "tuliprox" -version = "3.2.19" +version = "3.2.20" edition = "2021" # See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html diff --git a/frontend/Cargo.toml b/frontend/Cargo.toml index 15855ada7..14a93f3e7 100644 --- a/frontend/Cargo.toml +++ b/frontend/Cargo.toml @@ -1,10 +1,10 @@ [package] name = "frontend" -version = "3.2.19" +version = "3.2.20" edition = "2021" [dependencies] -shared = { version = "3.2.19", path = "../shared" } +shared = { version = "3.2.20", path = "../shared" } chrono = "0" yew = "0.21" yew-router = "0.18" diff --git a/shared/Cargo.toml b/shared/Cargo.toml index b658496f0..1ab71dcb2 100644 --- a/shared/Cargo.toml +++ b/shared/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "shared" -version = "3.2.19" +version = "3.2.20" edition = "2021" [dependencies]