From 4c5d96624842dab0ca63154b9ea1e3697fa1fa9e Mon Sep 17 00:00:00 2001 From: euzu <33094714+euzu@users.noreply.github.com> Date: Thu, 28 May 2026 17:58:41 +0200 Subject: [PATCH] Reduce heap allocation (#764) - Refactored PlaylistSource to avoid Box whre possible - Some other small heap allocation fixes --- backend/src/api/api_utils.rs | 24 +- backend/src/api/endpoints/hls_api.rs | 2 +- backend/src/api/endpoints/m3u_api.rs | 6 +- backend/src/api/endpoints/v1_api.rs | 2 +- backend/src/api/endpoints/v1_api_playlist.rs | 6 +- backend/src/api/endpoints/websocket_api.rs | 170 ++-- backend/src/api/endpoints/xtream_api.rs | 8 +- backend/src/api/model/app_state.rs | 13 +- backend/src/api/model/connection_manager.rs | 11 +- .../src/api/model/provider_lineup_manager.rs | 9 +- .../model/streams/provider_stream_factory.rs | 4 +- backend/src/api/panel_api.rs | 73 +- backend/src/library/processor.rs | 59 +- backend/src/model/config/input.rs | 2 +- backend/src/model/playlist.rs | 14 +- .../processor/filtered_playlist_source.rs | 173 ---- backend/src/processing/processor/library.rs | 31 +- backend/src/processing/processor/mod.rs | 1 - backend/src/processing/processor/playlist.rs | 69 +- .../src/processing/processor/xtream_series.rs | 6 +- backend/src/repository/bplustree.rs | 19 + .../src/repository/m3u_playlist_iterator.rs | 2 +- backend/src/repository/playlist_repository.rs | 53 +- backend/src/repository/playlist_source.rs | 793 ++++++++++++++++-- backend/src/utils/file/file_utils.rs | 3 +- backend/src/utils/network/request.rs | 24 +- .../src/app/components/field_explanation.rs | 4 +- frontend/src/app/components/field_id.rs | 14 +- .../src/app/components/setup/setup_helpers.rs | 9 +- .../source_editor/epg_smart_match_form.rs | 14 +- .../source_editor/provider_item_form.rs | 9 +- frontend/src/utils/mod.rs | 11 + shared/src/error/tuliprox_error.rs | 4 +- shared/src/foundation/filter.rs | 2 +- shared/src/foundation/value_provider.rs | 2 +- shared/src/model/config/input.rs | 2 +- shared/src/model/stream_properties.rs | 15 +- shared/src/utils/mod.rs | 11 +- shared/src/utils/size_utils.rs | 19 +- 39 files changed, 1200 insertions(+), 493 deletions(-) delete mode 100644 backend/src/processing/processor/filtered_playlist_source.rs diff --git a/backend/src/api/api_utils.rs b/backend/src/api/api_utils.rs index 2701faefd..ef90eecf9 100644 --- a/backend/src/api/api_utils.rs +++ b/backend/src/api/api_utils.rs @@ -29,6 +29,7 @@ use crate::{ BUILD_TIMESTAMP, }; use arc_swap::ArcSwapOption; +use tokio::sync::RwLock; use axum::{ body::Body, http::{header, Extensions, HeaderMap, HeaderName, HeaderValue, Response, StatusCode}, @@ -62,10 +63,7 @@ use std::{ path::{Path, PathBuf}, sync::Arc, }; -use tokio::{ - io::{AsyncReadExt, AsyncSeekExt}, - sync::Mutex, -}; +use tokio::io::{AsyncReadExt, AsyncSeekExt}; use tokio_util::io::ReaderStream; use url::Url; @@ -1979,7 +1977,7 @@ where let redirect_url = match resolve_redirect_location(Some(params.input), &redirect_url) { Ok(url) => url, Err(err) => { - error!("Failed to resolve redirect url: {}", sanitize_sensitive_info(err.to_string().as_str())); + error!("Failed to resolve redirect url: {}", sanitize_sensitive_info(&err.to_string())); return Some(StatusCode::BAD_REQUEST.into_response()); } }; @@ -2018,7 +2016,7 @@ where let stream_url = match resolve_redirect_location(Some(params.input), &stream_url) { Ok(url) => url, Err(err) => { - error!("Failed to resolve redirect url: {}", sanitize_sensitive_info(err.to_string().as_str())); + error!("Failed to resolve redirect url: {}", sanitize_sensitive_info(&err.to_string())); return Some(StatusCode::BAD_REQUEST.into_response()); } }; @@ -3179,7 +3177,7 @@ pub fn get_headers_from_request(req_headers: &HeaderMap, filter: &HeaderFilter) fn get_add_cache_content( res_url: &str, mime_type: Option, - cache: &Arc>>, + cache: &Arc>>, ) -> Arc { let resource_url = String::from(res_url); let cache = Arc::clone(cache); @@ -3190,7 +3188,7 @@ fn get_add_cache_content( let cache = Arc::clone(&cache); tokio::spawn(async move { if let Some(cache) = cache.load().as_ref() { - let _ = cache.lock().await.add_content(&res_url, mime_type, size); + let _ = cache.write().await.add_content(&res_url, mime_type, size); } }); }); @@ -3246,7 +3244,7 @@ async fn build_resource_stream_response( if can_cache { debug!("Caching eligible resource stream {sanitized_resource_url}"); let cache_resource_path = if let Some(cache) = app_state.cache.load().as_ref() { - Some(cache.lock().await.store_path(resource_url, mime_type.as_deref())) + Some(cache.write().await.store_path(resource_url, mime_type.as_deref())) } else { None }; @@ -3351,8 +3349,12 @@ pub async fn resource_response( let filter: HeaderFilter = Some(Box::new(|key| key != "if-none-match" && key != "if-modified-since")); let req_headers = get_headers_from_request(req_headers, &filter); if let Some(cache) = app_state.cache.load().as_ref() { - let mut guard = cache.lock().await; - if let Some((resource_path, mime_type)) = guard.get_content(resource_url) { + let cache_hit = { + let mut guard = cache.write().await; + guard.get_content(resource_url) + }; + + if let Some((resource_path, mime_type)) = cache_hit { trace_if_enabled!("Responding resource from cache {}", sanitize_sensitive_info(resource_url)); return serve_file( &resource_path, diff --git a/backend/src/api/endpoints/hls_api.rs b/backend/src/api/endpoints/hls_api.rs index 4770faf82..aa3cb4cee 100644 --- a/backend/src/api/endpoints/hls_api.rs +++ b/backend/src/api/endpoints/hls_api.rs @@ -310,7 +310,7 @@ pub(in crate::api) async fn handle_hls_stream_request( hls_response(hls_content).into_response() } Err(err) => { - error!("Failed to download m3u8 {}", sanitize_sensitive_info(err.to_string().as_str())); + error!("Failed to download m3u8 {}", sanitize_sensitive_info(&err.to_string())); if let Some(session_token) = session_token.as_deref() { terminate_failed_hls_manifest_session(app_state, &user.username, session_token).await; } diff --git a/backend/src/api/endpoints/m3u_api.rs b/backend/src/api/endpoints/m3u_api.rs index 907fa1984..907dfb508 100644 --- a/backend/src/api/endpoints/m3u_api.rs +++ b/backend/src/api/endpoints/m3u_api.rs @@ -57,7 +57,7 @@ async fn m3u_api( try_unwrap_body!(builder.body(axum::body::Body::from_stream(content_stream))) } Err(err) => { - error!("{}", sanitize_sensitive_info(err.to_string().as_str())); + error!("{}", sanitize_sensitive_info(&err.to_string())); axum::http::StatusCode::NO_CONTENT.into_response() } } @@ -424,7 +424,7 @@ async fn m3u_api_resource( let m3u_item = match m3u_get_item_for_stream_id(m3u_stream_id, &app_state, &target).await { Ok(item) => item, Err(err) => { - error!("Failed to get m3u url: {}", sanitize_sensitive_info(err.to_string().as_str())); + error!("Failed to get m3u url: {}", sanitize_sensitive_info(&err.to_string())); return axum::http::StatusCode::NOT_FOUND.into_response(); } }; @@ -446,7 +446,7 @@ async fn m3u_api_resource( redirect(redirect_url.as_ref()).into_response() } Err(err) => { - error!("Failed to resolve redirect url: {}", sanitize_sensitive_info(err.to_string().as_str())); + error!("Failed to resolve redirect url: {}", sanitize_sensitive_info(&err.to_string())); axum::http::StatusCode::BAD_REQUEST.into_response() } } diff --git a/backend/src/api/endpoints/v1_api.rs b/backend/src/api/endpoints/v1_api.rs index 162af5686..76dc4c6f8 100644 --- a/backend/src/api/endpoints/v1_api.rs +++ b/backend/src/api/endpoints/v1_api.rs @@ -42,7 +42,7 @@ async fn create_ipinfo_check(app_state: &Arc) -> Option<(Option) -> StatusCheck { let cache = match app_state.cache.load().as_ref().as_ref() { None => None, - Some(lock) => Some(lock.lock().await.get_size_text()), + Some(lock) => Some(lock.read().await.get_size_text()), }; let (active_users, active_user_connections, active_user_streams) = { let active_user = &app_state.active_users; diff --git a/backend/src/api/endpoints/v1_api_playlist.rs b/backend/src/api/endpoints/v1_api_playlist.rs index dc1cae350..9ac5d2436 100644 --- a/backend/src/api/endpoints/v1_api_playlist.rs +++ b/backend/src/api/endpoints/v1_api_playlist.rs @@ -225,7 +225,7 @@ async fn playlist_update( } } Err(err) => { - error!("Failed playlist update {}", sanitize_sensitive_info(err.to_string().as_str())); + error!("Failed playlist update {}", sanitize_sensitive_info(&err.to_string())); (axum::http::StatusCode::BAD_REQUEST, axum::Json(json!({"error": err.to_string()}))).into_response() } } @@ -404,7 +404,7 @@ async fn playlist_epg( } Ok(None) => return axum::http::StatusCode::NO_CONTENT.into_response(), Err(err) => { - error!("Failed to load input EPG for '{}': {}", input.name, sanitize_sensitive_info(err.to_string().as_str())); + error!("Failed to load input EPG for '{}': {}", input.name, sanitize_sensitive_info(&err.to_string())); return ( axum::http::StatusCode::INTERNAL_SERVER_ERROR, axum::Json(serde_json::json!({"error": "Failed to load EPG"})), @@ -436,7 +436,7 @@ async fn playlist_epg( return json_or_bin_response(accept.as_deref(), &epg).into_response(); } Err(err) => { - error!("Failed to load custom EPG: {}", sanitize_sensitive_info(err.to_string().as_str())); + error!("Failed to load custom EPG: {}", sanitize_sensitive_info(&err.to_string())); return ( axum::http::StatusCode::INTERNAL_SERVER_ERROR, axum::Json(serde_json::json!({"error": "Failed to load EPG"})), diff --git a/backend/src/api/endpoints/websocket_api.rs b/backend/src/api/endpoints/websocket_api.rs index c88a922bc..32e892dac 100644 --- a/backend/src/api/endpoints/websocket_api.rs +++ b/backend/src/api/endpoints/websocket_api.rs @@ -17,7 +17,45 @@ use shared::{ }, utils::{concat_path_leading_slash, default_kick_secs}, }; -use std::sync::Arc; +use std::{fmt, io, sync::Arc}; + +#[derive(Debug)] +enum WebSocketApiError { + Transport(axum::Error), + Protocol(io::Error), + ProtocolVersionMismatch, + EventSend { context: &'static str, source: axum::Error }, +} + +impl fmt::Display for WebSocketApiError { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::Transport(err) => write!(f, "{err}"), + Self::Protocol(err) => write!(f, "{err}"), + Self::ProtocolVersionMismatch => f.write_str("Protocol version mismatch"), + Self::EventSend { context, source } => write!(f, "{context}: {source}"), + } + } +} + +impl std::error::Error for WebSocketApiError { + fn source(&self) -> Option<&(dyn std::error::Error + 'static)> { + match self { + Self::Transport(err) => Some(err), + Self::Protocol(err) => Some(err), + Self::ProtocolVersionMismatch => None, + Self::EventSend { source, .. } => Some(source), + } + } +} + +impl From for WebSocketApiError { + fn from(value: axum::Error) -> Self { Self::Transport(value) } +} + +impl From for WebSocketApiError { + fn from(value: io::Error) -> Self { Self::Protocol(value) } +} // WebSocket upgrade handler async fn websocket_handler( @@ -87,12 +125,12 @@ fn get_secret_key(app_state: &AppState, auth: bool) -> Option> { }) } -async fn handle_handshake(msg: Message, socket: &mut WebSocket, version: u8) -> Result<(), String> { +async fn handle_handshake(msg: Message, socket: &mut WebSocket, version: u8) -> Result<(), WebSocketApiError> { if let Message::Binary(bytes) = msg { if bytes.len() == 1 { let client_version = bytes[0]; if client_version == version { - socket.send(Message::binary(bytes)).await.map_err(|e| e.to_string())?; + socket.send(Message::binary(bytes)).await?; return Ok(()); } error!("Protocol Version mismatch: server={version}, client={client_version}"); @@ -106,7 +144,7 @@ async fn handle_handshake(msg: Message, socket: &mut WebSocket, version: u8) -> }))) .await; - Err("Protocol version mismatch".into()) + Err(WebSocketApiError::ProtocolVersionMismatch) } async fn handle_protocol_message( @@ -251,8 +289,8 @@ async fn handle_incoming_message( app_state: &Arc, auth_required: bool, secret_key: Option<&Vec>, -) -> Result<(), String> { - let msg = result.map_err(|e| e.to_string())?; +) -> Result<(), WebSocketApiError> { + let msg = result?; match handler { ProtocolHandler::Version(version) => { @@ -272,9 +310,9 @@ async fn handle_incoming_message( Some(protocol_msg) => { let bytes = match protocol_msg.to_bytes() { Ok(bytes) => bytes, - Err(err) => ProtocolMessage::Error(err.to_string()).to_bytes().map_err(|e| e.to_string())?, + Err(err) => ProtocolMessage::Error(err.to_string()).to_bytes()?, }; - Ok(socket.send(Message::Binary(bytes)).await.map_err(|e| e.to_string())?) + Ok(socket.send(Message::Binary(bytes)).await?) } } } @@ -285,78 +323,82 @@ async fn handle_event_message( socket: &mut WebSocket, event: EventMessage, handler: &ProtocolHandler, -) -> Result<(), String> { +) -> Result<(), WebSocketApiError> { match handler { ProtocolHandler::Version(_) => {} ProtocolHandler::Default(mem) => { if websocket_can_receive_runtime_events(mem, &event) { match event { EventMessage::ServerError(error) => { - let msg = ProtocolMessage::ServerError(error).to_bytes().map_err(|e| e.to_string())?; - socket.send(Message::Binary(msg)).await.map_err(|e| format!("Server Error event: {e} "))?; + send_event_response(socket, ProtocolMessage::ServerError(error), "Server Error event").await?; } EventMessage::ActiveUser(event) => { - let msg = ProtocolMessage::ActiveUserResponse(event).to_bytes().map_err(|e| e.to_string())?; - socket - .send(Message::Binary(msg)) - .await - .map_err(|e| format!("Active user connection change event: {e} "))?; + send_event_response( + socket, + ProtocolMessage::ActiveUserResponse(event), + "Active user connection change event", + ) + .await?; } EventMessage::ActiveProvider(provider, connections) => { - let msg = ProtocolMessage::ActiveProviderResponse(provider, connections) - .to_bytes() - .map_err(|e| e.to_string())?; - socket - .send(Message::Binary(msg)) - .await - .map_err(|e| format!("Provider connection change event: {e} "))?; + send_event_response( + socket, + ProtocolMessage::ActiveProviderResponse(provider, connections), + "Provider connection change event", + ) + .await?; } EventMessage::ConfigChange(config) => { - let msg = - ProtocolMessage::ConfigChangeResponse(config).to_bytes().map_err(|e| e.to_string())?; - socket - .send(Message::Binary(msg)) - .await - .map_err(|e| format!("Configuration files change event: {e} "))?; + send_event_response( + socket, + ProtocolMessage::ConfigChangeResponse(config), + "Configuration files change event", + ) + .await?; } EventMessage::PlaylistUpdate(state) => { - let msg = - ProtocolMessage::PlaylistUpdateResponse(state).to_bytes().map_err(|e| e.to_string())?; - socket.send(Message::Binary(msg)).await.map_err(|e| format!("Playlist update event: {e} "))?; + send_event_response( + socket, + ProtocolMessage::PlaylistUpdateResponse(state), + "Playlist update event", + ) + .await?; } EventMessage::PlaylistUpdateProgress(target, msg) => { - let msg = ProtocolMessage::PlaylistUpdateProgressResponse(target, msg) - .to_bytes() - .map_err(|e| e.to_string())?; - socket - .send(Message::Binary(msg)) - .await - .map_err(|e| format!("Playlist update progress event: {e} "))?; + send_event_response( + socket, + ProtocolMessage::PlaylistUpdateProgressResponse(target, msg), + "Playlist update progress event", + ) + .await?; } EventMessage::SystemInfoUpdate(system_info) => { - let msg = - ProtocolMessage::SystemInfoResponse(system_info).to_bytes().map_err(|e| e.to_string())?; - socket.send(Message::Binary(msg)).await.map_err(|e| format!("System info event: {e} "))?; + send_event_response( + socket, + ProtocolMessage::SystemInfoResponse(system_info), + "System info event", + ) + .await?; } EventMessage::LibraryScanProgress(summary) => { - let msg = ProtocolMessage::LibraryScanProgressResponse(summary) - .to_bytes() - .map_err(|e| e.to_string())?; - socket - .send(Message::Binary(msg)) - .await - .map_err(|e| format!("Library scan progress event: {e} "))?; + send_event_response( + socket, + ProtocolMessage::LibraryScanProgressResponse(summary), + "Library scan progress event", + ) + .await?; } EventMessage::DownloadsUpdate(downloads) => { - let msg = ProtocolMessage::DownloadsResponse(downloads).to_bytes().map_err(|e| e.to_string())?; - socket.send(Message::Binary(msg)).await.map_err(|e| format!("Downloads event: {e} "))?; + send_event_response(socket, ProtocolMessage::DownloadsResponse(downloads), "Downloads event") + .await?; } EventMessage::DownloadsDeltaUpdate(delta) => { - let msg = ProtocolMessage::DownloadsDeltaResponse(delta).to_bytes().map_err(|e| e.to_string())?; - socket - .send(Message::Binary(msg)) - .await - .map_err(|e| format!("Downloads delta event: {e} "))?; + send_event_response( + socket, + ProtocolMessage::DownloadsDeltaResponse(delta), + "Downloads delta event", + ) + .await?; } EventMessage::InputMetadataUpdatesCompleted(_) | EventMessage::InputMetadataUpdatesStarted(_) => { @@ -369,6 +411,18 @@ async fn handle_event_message( Ok(()) } +async fn send_event_response( + socket: &mut WebSocket, + message: ProtocolMessage, + context: &'static str, +) -> Result<(), WebSocketApiError> { + let msg = message.to_bytes()?; + socket + .send(Message::Binary(msg)) + .await + .map_err(|source| WebSocketApiError::EventSend { context, source }) +} + // WebSocket communication logic async fn handle_socket(mut socket: WebSocket, app_state: Arc, auth_required: bool) { let secret_key = get_secret_key(&app_state, auth_required); @@ -410,9 +464,7 @@ async fn handle_socket(mut socket: WebSocket, app_state: Arc, auth_req Ok(entries) => { if let ProtocolHandler::Default(mem) = &handler { if mem.stream_meter_subscribed { - let msg = ProtocolMessage::StreamMeterBatchResponse(entries) - .to_bytes() - .map_err(|e| e.to_string()); + let msg = ProtocolMessage::StreamMeterBatchResponse(entries).to_bytes(); match msg { Ok(msg) => { if let Err(e) = socket.send(Message::Binary(msg)).await { diff --git a/backend/src/api/endpoints/xtream_api.rs b/backend/src/api/endpoints/xtream_api.rs index 50dfc9c0d..9fbe1b398 100644 --- a/backend/src/api/endpoints/xtream_api.rs +++ b/backend/src/api/endpoints/xtream_api.rs @@ -777,7 +777,7 @@ async fn xtream_player_api_resource( redirect(redirect_url.as_ref()).into_response() } Err(err) => { - error!("Failed to resolve redirect url: {}", sanitize_sensitive_info(err.to_string().as_str())); + error!("Failed to resolve redirect url: {}", sanitize_sensitive_info(&err.to_string())); axum::http::StatusCode::BAD_REQUEST.into_response() } } @@ -1017,7 +1017,7 @@ pub async fn xtream_get_stream_info_response( return match api_utils::resolve_redirect_location(Some(&input), &info_url) { Ok(redirect_url) => redirect(redirect_url.as_ref()).into_response(), Err(err) => { - error!("Failed to resolve redirect url: {}", sanitize_sensitive_info(err.to_string().as_str())); + error!("Failed to resolve redirect url: {}", sanitize_sensitive_info(&err.to_string())); axum::http::StatusCode::BAD_REQUEST.into_response() } }; @@ -1108,7 +1108,7 @@ async fn xtream_get_short_epg( return match api_utils::resolve_redirect_location(Some(&input), &info_url) { Ok(redirect_url) => redirect(redirect_url.as_ref()).into_response(), Err(err) => { - error!("Failed to resolve redirect url: {}", sanitize_sensitive_info(err.to_string().as_str())); + error!("Failed to resolve redirect url: {}", sanitize_sensitive_info(&err.to_string())); axum::http::StatusCode::BAD_REQUEST.into_response() } }; @@ -1132,7 +1132,7 @@ async fn xtream_get_short_epg( ) .into_response(), Err(err) => { - error!("Failed to download epg {}", sanitize_sensitive_info(err.to_string().as_str())); + error!("Failed to download epg {}", sanitize_sensitive_info(&err.to_string())); axum::Json(json!(ShortEpgResultDto::default())).into_response() } }; diff --git a/backend/src/api/model/app_state.rs b/backend/src/api/model/app_state.rs index 2b5004785..7f7e00b38 100644 --- a/backend/src/api/model/app_state.rs +++ b/backend/src/api/model/app_state.rs @@ -37,7 +37,8 @@ use std::{ sync::{atomic::AtomicI8, Arc}, time::Duration, }; -use tokio::sync::{mpsc, Mutex}; +use tokio::sync::mpsc; +use tokio::sync::RwLock; use tokio::task; use tokio_util::sync::CancellationToken; use url::Url; @@ -323,7 +324,7 @@ fn build_http_client_with_fallback( Ok(fallback_client()) } -pub fn create_cache(config: &Config) -> Option>> { +pub fn create_cache(config: &Config) -> Option>> { let lru_cache = config.reverse_proxy.as_ref().and_then(|r| r.cache.as_ref()).and_then(|c| { if c.enabled { Some(LRUResourceCache::new(c.size, c.directory.as_str())) @@ -335,11 +336,11 @@ pub fn create_cache(config: &Config) -> Option>> { if cache_enabled { info!("Scanning cache"); if let Some(res_cache) = lru_cache { - let cache = Arc::new(Mutex::new(res_cache)); + let cache = Arc::new(RwLock::new(res_cache)); let cache_scanner = Arc::clone(&cache); tokio::spawn(async move { let scan_result = { - let mut cache = cache_scanner.lock().await; + let mut cache = cache_scanner.write().await; task::block_in_place(|| cache.scan()) }; if let Err(err) = scan_result { @@ -396,7 +397,7 @@ pub struct AppState { pub http_client: Arc>, pub http_client_no_redirect: Arc>, pub downloads: Arc, - pub cache: Arc>>, + pub cache: Arc>>, pub shared_stream_manager: Arc, pub active_users: Arc, pub active_provider: Arc, @@ -465,7 +466,7 @@ impl AppState { if let Some(cache) = self.cache.load().as_ref() { if enabled { - cache.lock().await.update_config(size, cache_dir); + cache.write().await.update_config(size, cache_dir); } else { self.cache.store(None); } diff --git a/backend/src/api/model/connection_manager.rs b/backend/src/api/model/connection_manager.rs index f71ce1098..5628d7ae4 100644 --- a/backend/src/api/model/connection_manager.rs +++ b/backend/src/api/model/connection_manager.rs @@ -223,13 +223,10 @@ impl SocketActivityTracker { fn lock_socket_activity_pending( pending: &Mutex>, ) -> MutexGuard<'_, HashMap> { - match pending.lock() { - Ok(guard) => guard, - Err(poisoned) => { - warn!("Socket activity state was poisoned, continuing with recovered state"); - poisoned.into_inner() - } - } + pending.lock().unwrap_or_else(|poisoned| { + warn!("Socket activity state was poisoned, continuing with recovered state"); + poisoned.into_inner() + }) } struct CleanupWorkerDeps { diff --git a/backend/src/api/model/provider_lineup_manager.rs b/backend/src/api/model/provider_lineup_manager.rs index a402a26dc..ed9dab5a2 100644 --- a/backend/src/api/model/provider_lineup_manager.rs +++ b/backend/src/api/model/provider_lineup_manager.rs @@ -983,11 +983,12 @@ mod tests { model::{InputFetchMethod, InputType}, utils::Internable, }; - use std::{sync::atomic::AtomicU16, thread}; + use std::sync::atomic::AtomicU16; + use tokio::time::{sleep, Duration}; macro_rules! should_available { ($lineup:expr, $provider_id:expr, $grace_period_timeout_secs: expr) => { - thread::sleep(std::time::Duration::from_millis(200)); + sleep(Duration::from_millis(200)).await; match $lineup.acquire(true, $grace_period_timeout_secs).await { ProviderAllocation::Exhausted => assert!(false, "Should available and not exhausted"), ProviderAllocation::Available(provider) => assert_eq!(provider.id, $provider_id), @@ -999,7 +1000,7 @@ mod tests { } macro_rules! should_grace_period { ($lineup:expr, $provider_id:expr, $grace_period_timeout_secs: expr) => { - thread::sleep(std::time::Duration::from_millis(200)); + sleep(Duration::from_millis(200)).await; match $lineup.acquire(true, $grace_period_timeout_secs).await { ProviderAllocation::Exhausted => assert!(false, "Should grace period and not exhausted"), ProviderAllocation::Available(provider) => { @@ -1012,7 +1013,7 @@ mod tests { macro_rules! should_exhausted { ($lineup:expr, $grace_period_timeout_secs: expr) => { - thread::sleep(std::time::Duration::from_millis(200)); + sleep(Duration::from_millis(200)).await; match $lineup.acquire(true, $grace_period_timeout_secs).await { ProviderAllocation::Exhausted => {} ProviderAllocation::Available(provider) => { diff --git a/backend/src/api/model/streams/provider_stream_factory.rs b/backend/src/api/model/streams/provider_stream_factory.rs index 37cee60ed..f0142dbcd 100644 --- a/backend/src/api/model/streams/provider_stream_factory.rs +++ b/backend/src/api/model/streams/provider_stream_factory.rs @@ -509,7 +509,7 @@ async fn send_with_manual_redirects( // send_with_retry_and_provider already applies provider failover policy. // Do not rotate again here, otherwise non-failover errors (e.g. auth) may // incorrectly switch provider URLs. - debug!("Manual redirect failed: {}", sanitize_sensitive_info(e.to_string().as_str())); + debug!("Manual redirect failed: {}", sanitize_sensitive_info(&e.to_string())); return Err(e); } }; @@ -662,7 +662,7 @@ async fn provider_stream_request( let diagnostics = preview_request_diagnostics_for_logging(stream_options.get_url(), stream_options.get_provider()); debug!( "Provider request failed: {}, {}", - sanitize_sensitive_info(err.to_string().as_str()), + sanitize_sensitive_info(&err.to_string()), sanitize_sensitive_info(&diagnostics) ); Err(ProviderStreamRequestFailure::Status { diff --git a/backend/src/api/panel_api.rs b/backend/src/api/panel_api.rs index 8d37a4e75..53197d03f 100644 --- a/backend/src/api/panel_api.rs +++ b/backend/src/api/panel_api.rs @@ -662,7 +662,7 @@ async fn fetch_root_user_api_info( debug_if_enabled!( "panel_api user_api account_info failed for {}: {}", sanitize_sensitive_info(&input.name), - sanitize_sensitive_info(err.to_string().as_str()) + sanitize_sensitive_info(&err.to_string()) ); return Ok(None); } @@ -1094,7 +1094,7 @@ pub(crate) fn can_provision_on_exhausted(app_state: &AppState, input: &ConfigInp return false; } if let Err(err) = validate_panel_api_config(panel_cfg) { - debug_if_enabled!("panel_api config invalid: {}", sanitize_sensitive_info(err.to_string().as_str())); + debug_if_enabled!("panel_api config invalid: {}", sanitize_sensitive_info(&err.to_string())); return false; } if is_alias_pool_max_reached(app_state, input) { @@ -1604,7 +1604,7 @@ async fn try_renew_expired_account( debug_if_enabled!( "panel_api client_adult_content failed for {}: {}", sanitize_sensitive_info(&acct.name), - sanitize_sensitive_info(err.to_string().as_str()) + sanitize_sensitive_info(&err.to_string()) ); } } @@ -1657,7 +1657,7 @@ async fn try_renew_expired_account( debug_if_enabled!( "panel_api client_renew failed for {}: {}", sanitize_sensitive_info(&acct.name), - sanitize_sensitive_info(err.to_string().as_str()) + sanitize_sensitive_info(&err.to_string()) ); } } @@ -1698,7 +1698,7 @@ async fn try_create_new_account( debug_if_enabled!( "panel_api client_adult_content failed for {}: {}", sanitize_sensitive_info(&alias_name), - sanitize_sensitive_info(err.to_string().as_str()) + sanitize_sensitive_info(&err.to_string()) ); } } @@ -1767,7 +1767,7 @@ async fn try_create_new_account( Some(PanelApiProvisionOutcome::Created { username, password }) } Err(err) => { - debug_if_enabled!("panel_api client_new failed: {}", sanitize_sensitive_info(err.to_string().as_str())); + debug_if_enabled!("panel_api client_new failed: {}", sanitize_sensitive_info(&err.to_string())); None } } @@ -1803,7 +1803,7 @@ pub async fn try_provision_account_on_exhausted( app_state.app_config.file_locks.write_lock_str(format!("panel_api:{}", input.name).as_str()).await; if let Err(err) = validate_panel_api_config(panel_cfg) { - debug_if_enabled!("panel_api config invalid: {}", sanitize_sensitive_info(err.to_string().as_str())); + debug_if_enabled!("panel_api config invalid: {}", sanitize_sensitive_info(&err.to_string())); return None; } if is_alias_pool_max_reached(app_state, input) { @@ -1908,7 +1908,7 @@ async fn ensure_alias_pool_min( debug_if_enabled!( "panel_api client_adult_content failed for {}: {}", sanitize_sensitive_info(&acct.name), - sanitize_sensitive_info(err.to_string().as_str()) + sanitize_sensitive_info(&err.to_string()) ); } } @@ -1960,7 +1960,7 @@ async fn ensure_alias_pool_min( debug_if_enabled!( "panel_api client_renew failed for {}: {}", sanitize_sensitive_info(&acct.name), - sanitize_sensitive_info(err.to_string().as_str()) + sanitize_sensitive_info(&err.to_string()) ); } } @@ -1989,7 +1989,7 @@ async fn ensure_alias_pool_min( debug_if_enabled!( "panel_api client_adult_content failed for {}: {}", sanitize_sensitive_info(&alias_name), - sanitize_sensitive_info(err.to_string().as_str()) + sanitize_sensitive_info(&err.to_string()) ); } } @@ -2044,7 +2044,7 @@ async fn ensure_alias_pool_min( } } Err(err) => { - debug_if_enabled!("panel_api client_new failed: {}", sanitize_sensitive_info(err.to_string().as_str())); + debug_if_enabled!("panel_api client_new failed: {}", sanitize_sensitive_info(&err.to_string())); break; } } @@ -2070,7 +2070,7 @@ async fn sync_panel_api_for_input_on_boot( debug_if_enabled!( "panel_api boot sync skipped for {}: {}", sanitize_sensitive_info(&input.name), - sanitize_sensitive_info(err.to_string().as_str()) + sanitize_sensitive_info(&err.to_string()) ); return false; } @@ -2126,7 +2126,7 @@ async fn sync_panel_api_for_input_on_boot( debug_if_enabled!( "panel_api root client_info (raw) failed for {}: {}", sanitize_sensitive_info(&input.name), - sanitize_sensitive_info(err.to_string().as_str()) + sanitize_sensitive_info(&err.to_string()) ); None } @@ -2200,7 +2200,7 @@ async fn sync_panel_api_for_input_on_boot( debug_if_enabled!( "panel_api client_info failed for {}: {}", sanitize_sensitive_info(&acct.name), - sanitize_sensitive_info(err.to_string().as_str()) + sanitize_sensitive_info(&err.to_string()) ); None } @@ -2311,7 +2311,7 @@ async fn sync_panel_api_for_input_on_boot( debug_if_enabled!( "panel_api client_renew failed for root {}: {}", sanitize_sensitive_info(&input.name), - sanitize_sensitive_info(err.to_string().as_str()) + sanitize_sensitive_info(&err.to_string()) ); if new_enabled { match panel_client_new(app_state.as_ref(), panel_cfg).await { @@ -2431,7 +2431,7 @@ async fn sync_panel_api_for_input_on_boot( debug_if_enabled!( "panel_api client_new failed for root {}: {}", sanitize_sensitive_info(&input.name), - sanitize_sensitive_info(err.to_string().as_str()) + sanitize_sensitive_info(&err.to_string()) ); (old_username.clone(), old_password.clone(), false) } @@ -2552,7 +2552,7 @@ async fn sync_panel_api_for_input_on_boot( debug_if_enabled!( "panel_api client_new failed for root {}: {}", sanitize_sensitive_info(&input.name), - sanitize_sensitive_info(err.to_string().as_str()) + sanitize_sensitive_info(&err.to_string()) ); (old_username.clone(), old_password.clone(), false) } @@ -2572,7 +2572,7 @@ async fn sync_panel_api_for_input_on_boot( debug_if_enabled!( "panel_api client_adult_content failed for root {}: {}", sanitize_sensitive_info(&input.name), - sanitize_sensitive_info(err.to_string().as_str()) + sanitize_sensitive_info(&err.to_string()) ); } } @@ -2800,7 +2800,7 @@ async fn sync_panel_api_for_input_on_boot( debug_if_enabled!( "panel_api client_renew failed for alias {}: {}", sanitize_sensitive_info(&account_name), - sanitize_sensitive_info(err.to_string().as_str()) + sanitize_sensitive_info(&err.to_string()) ); if new_enabled { match panel_client_new(app_state.as_ref(), panel_cfg).await { @@ -2827,7 +2827,7 @@ async fn sync_panel_api_for_input_on_boot( debug_if_enabled!( "panel_api client_adult_content failed for {}: {}", sanitize_sensitive_info(&alias_name), - sanitize_sensitive_info(err.to_string().as_str()) + sanitize_sensitive_info(&err.to_string()) ); } } @@ -2896,7 +2896,7 @@ async fn sync_panel_api_for_input_on_boot( debug_if_enabled!( "panel_api client_new failed for input {}: {}", sanitize_sensitive_info(&input.name), - sanitize_sensitive_info(err.to_string().as_str()) + sanitize_sensitive_info(&err.to_string()) ); continue; } @@ -2929,7 +2929,7 @@ async fn sync_panel_api_for_input_on_boot( debug_if_enabled!( "panel_api client_adult_content failed for {}: {}", sanitize_sensitive_info(&alias_name), - sanitize_sensitive_info(err.to_string().as_str()) + sanitize_sensitive_info(&err.to_string()) ); } } @@ -2998,7 +2998,7 @@ async fn sync_panel_api_for_input_on_boot( debug_if_enabled!( "panel_api client_new failed for input {}: {}", sanitize_sensitive_info(&input.name), - sanitize_sensitive_info(err.to_string().as_str()) + sanitize_sensitive_info(&err.to_string()) ); continue; } @@ -3018,7 +3018,7 @@ async fn sync_panel_api_for_input_on_boot( debug_if_enabled!( "panel_api client_adult_content failed for {}: {}", sanitize_sensitive_info(&account_name), - sanitize_sensitive_info(err.to_string().as_str()) + sanitize_sensitive_info(&err.to_string()) ); } } @@ -3113,7 +3113,7 @@ async fn sync_panel_api_for_input_on_boot( debug_if_enabled!( "panel_api client_adult_content failed for {}: {}", sanitize_sensitive_info(&acct.name), - sanitize_sensitive_info(err.to_string().as_str()) + sanitize_sensitive_info(&err.to_string()) ); } } @@ -3193,7 +3193,7 @@ async fn sync_panel_api_for_input_on_boot( debug_if_enabled!( "panel_api account_info failed for {}: {}", sanitize_sensitive_info(&input.name), - sanitize_sensitive_info(err.to_string().as_str()) + sanitize_sensitive_info(&err.to_string()) ); } } @@ -3277,7 +3277,7 @@ pub(crate) async fn sync_panel_api_alias_pool_for_target(app_state: &Arc String { + let mut result = String::new(); + for (idx, target) in targets.iter().enumerate() { + if idx > 0 { + result.push(','); + } + result.push_str(target.action()); + } + result +} + fn build_panel_api_probe_targets(input: &ConfigInput, username: &str, password: &str) -> Vec { let mut targets = Vec::new(); for action in ["client_info", "get_live_categories", "get_series_categories", "get_vod_categories"] { @@ -3445,7 +3456,7 @@ async fn wait_for_panel_api_account_ready( return false; } - let targets_list = probe_targets.iter().map(PanelApiProbeTarget::action).collect::>().join(","); + let targets_list = format_probe_target_actions(&probe_targets); debug_if_enabled!( "panel_api probe start for {} (input={} timeout={}s interval={}s method={}) targets={}", sanitize_sensitive_info(account_name), @@ -3593,7 +3604,7 @@ pub(crate) async fn run_panel_api_provisioning_probe( debug_if_enabled!( "panel_api provisioning closing client connection for input {} addr={}", sanitize_sensitive_info(&input.name), - sanitize_sensitive_info(addr.to_string().as_str()) + sanitize_sensitive_info(&addr.to_string()) ); let _ = app_state .connection_manager @@ -3702,7 +3713,7 @@ pub(crate) async fn run_panel_api_provisioning_probe( debug_if_enabled!( "panel_api provisioning closing client connection for input {} addr={}", sanitize_sensitive_info(&input.name), - sanitize_sensitive_info(addr.to_string().as_str()) + sanitize_sensitive_info(&addr.to_string()) ); let _ = app_state .connection_manager diff --git a/backend/src/library/processor.rs b/backend/src/library/processor.rs index 77d364250..7b36dd6e2 100644 --- a/backend/src/library/processor.rs +++ b/backend/src/library/processor.rs @@ -10,7 +10,9 @@ use log::{debug, error, info, warn}; use path_clean::PathClean; use shared::model::{LibraryMetadataFormat, LibraryScanResult}; use std::collections::HashMap; +use std::fmt; use std::future::Future; +use std::io; use std::path::PathBuf; use std::pin::Pin; use std::sync::Arc; @@ -26,6 +28,34 @@ enum ProcessAction { Unchanged, } +#[derive(Debug)] +enum LibraryProcessError { + Resolve(String), + Io(io::Error), +} + +impl fmt::Display for LibraryProcessError { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::Resolve(message) => f.write_str(message), + Self::Io(err) => write!(f, "{err}"), + } + } +} + +impl std::error::Error for LibraryProcessError { + fn source(&self) -> Option<&(dyn std::error::Error + 'static)> { + match self { + Self::Resolve(_) => None, + Self::Io(err) => Some(err), + } + } +} + +impl From for LibraryProcessError { + fn from(value: io::Error) -> Self { Self::Io(value) } +} + // VOD processor that orchestrates scanning, classification, metadata resolution, and storage pub struct LibraryProcessor { config: LibraryConfig, @@ -225,7 +255,7 @@ impl LibraryProcessor { force_rescan: bool, can_probe: bool, can_extract_thumbnails: bool, - ) -> Result { + ) -> Result { match group { MediaGroup::Movie { .. } => { self.process_movie(group, existing_map, force_rescan, can_probe, can_extract_thumbnails).await @@ -264,8 +294,10 @@ impl LibraryProcessor { force_rescan: bool, can_probe: bool, can_extract_thumbnails: bool, - ) -> Result { - let MediaGroup::Movie { file, .. } = group else { return Err(format!("Expected movie to resolve but got {group}")) }; + ) -> Result { + let MediaGroup::Movie { file, .. } = group else { + return Err(LibraryProcessError::Resolve(format!("Expected movie to resolve but got {group}"))); + }; // Check if file already exists in cache let (mut cache_entry, status) = if let Some(existing_entry) = existing_map.get(&file.file_path) { // Check if file has been modified @@ -311,8 +343,8 @@ impl LibraryProcessor { file.modified_timestamp, can_extract_thumbnails, ).await; - self.storage.store(&cache_entry).await.map_err(|e| e.to_string())?; - self.write_metadata_files(&cache_entry).await.map_err(|e| e.to_string())?; + self.storage.store(&cache_entry).await?; + self.write_metadata_files(&cache_entry).await?; Ok(status) } @@ -325,8 +357,10 @@ impl LibraryProcessor { force_rescan: bool, can_probe: bool, can_extract_thumbnails: bool, - ) -> Result { - let MediaGroup::Series { show_key, episodes } = group else { return Err(format!("Expected series to resolve but got {group}")) }; + ) -> Result { + let MediaGroup::Series { show_key, episodes } = group else { + return Err(LibraryProcessError::Resolve(format!("Expected series to resolve but got {group}"))); + }; let series_file_path = episodes .iter() .find_map(|episode| { @@ -508,14 +542,17 @@ impl LibraryProcessor { } } - self.storage.store(&chache_entry).await.map_err(|e| e.to_string())?; - self.write_metadata_files(&chache_entry).await.map_err(|e| e.to_string())?; + self.storage.store(&chache_entry).await?; + self.write_metadata_files(&chache_entry).await?; Ok(status) } // Resolves metadata for a video file - async fn resolve_metadata(&self, file: &MediaGroup) -> Result { - self.resolver.resolve(file).await.ok_or_else(|| format!("Could not resolve metadata for {file}")) + async fn resolve_metadata(&self, file: &MediaGroup) -> Result { + self.resolver + .resolve(file) + .await + .ok_or_else(|| LibraryProcessError::Resolve(format!("Could not resolve metadata for {file}"))) } /// Extracts and caches a thumbnail if no TMDB poster is available. diff --git a/backend/src/model/config/input.rs b/backend/src/model/config/input.rs index b2d8167f5..5823f102f 100644 --- a/backend/src/model/config/input.rs +++ b/backend/src/model/config/input.rs @@ -327,7 +327,7 @@ impl ConfigInput { TuliproxError::ConfigInput(format!( "Malformed provider URL {}: {}", sanitize_sensitive_info(url), - sanitize_sensitive_info(err.to_string().as_str()) + sanitize_sensitive_info(&err.to_string()) )) })?; diff --git a/backend/src/model/playlist.rs b/backend/src/model/playlist.rs index 1bb6f5fa4..8ad60e927 100644 --- a/backend/src/model/playlist.rs +++ b/backend/src/model/playlist.rs @@ -1,12 +1,13 @@ -use std::collections::HashSet; use crate::model::{ConfigInput, TVGuide}; -use shared::model::{PlaylistGroup, PlaylistItem}; -use shared::model::UUIDType; use crate::repository::PlaylistSource; +use shared::error::TuliproxError; +use shared::model::UUIDType; +use shared::model::{PlaylistGroup, PlaylistItem}; +use std::collections::HashSet; pub struct FetchedPlaylist<'a> { pub input: &'a ConfigInput, - pub source: Box, + pub source: PlaylistSource, pub epg: Option, } @@ -54,8 +55,5 @@ impl FetchedPlaylist<'_> { self.source.deduplicate(duplicates); } - pub fn clone_source(&self) -> Box { - self.source.clone_box() - } + pub fn clone_source(&self) -> Result { self.source.clone_source() } } - diff --git a/backend/src/processing/processor/filtered_playlist_source.rs b/backend/src/processing/processor/filtered_playlist_source.rs deleted file mode 100644 index 2dcfb2b5e..000000000 --- a/backend/src/processing/processor/filtered_playlist_source.rs +++ /dev/null @@ -1,173 +0,0 @@ -use crate::repository::{MemoryPlaylistSource, PlaylistSource}; -use futures::future::BoxFuture; -use shared::model::{PlaylistGroup, PlaylistItem, PlaylistItemType, UUIDType, XtreamCluster}; -use std::collections::HashSet; -use std::sync::Arc; - -pub(crate) struct FilteredPlaylistSource { - inner: Box, - skip_set: Arc>, -} - -impl FilteredPlaylistSource { - pub(crate) fn new(inner: Box, skip_set: HashSet) -> Self { - Self { - inner, - skip_set: Arc::new(skip_set), - } - } - - fn filter_group(&self, mut group: PlaylistGroup) -> Option { - if self.skip_set.contains(&group.xtream_cluster) { - return None; - } - group.channels.retain(|item| !self.skip_set.contains(&item.header.xtream_cluster)); - if group.channels.is_empty() { - None - } else { - Some(group) - } - } -} - -impl PlaylistSource for FilteredPlaylistSource { - fn is_memory(&self) -> bool { - self.inner.is_memory() - } - - fn get_channel_count(&mut self) -> usize { - let skip_set = Arc::clone(&self.skip_set); - self.inner - .items() - .filter(move |item| !skip_set.contains(&item.as_ref().header.xtream_cluster)) - .count() - } - - fn get_group_count(&mut self) -> usize { - let skip_set = Arc::clone(&self.skip_set); - let mut groups = HashSet::<(XtreamCluster, Arc)>::new(); - for item in self - .inner - .items() - .filter(move |item| !skip_set.contains(&item.as_ref().header.xtream_cluster)) - { - let pli = item.as_ref(); - groups.insert((pli.header.xtream_cluster, Arc::clone(&pli.header.group))); - } - groups.len() - } - - fn is_empty(&mut self) -> bool { - self.get_channel_count() == 0 - } - - #[allow(clippy::wrong_self_convention)] - fn into_items(&mut self) -> Box + Send + '_> { - let skip_set = Arc::clone(&self.skip_set); - Box::new( - self.inner - .into_items() - .filter(move |item| !skip_set.contains(&item.header.xtream_cluster)), - ) - } - - fn items_mut(&mut self) -> Box + Send + '_> { - let skip_set = Arc::clone(&self.skip_set); - Box::new(self.inner.items_mut().filter_map(move |item| { - if skip_set.contains(&item.header.xtream_cluster) { - None - } else { - Some(item) - } - })) - } - - fn items<'a>(&'a mut self) -> Box> + Send + 'a> { - let skip_set = Arc::clone(&self.skip_set); - Box::new( - self.inner - .items() - .filter(move |item| !skip_set.contains(&item.as_ref().header.xtream_cluster)), - ) - } - - fn update_playlist<'a>(&'a mut self, plg: &'a PlaylistGroup) -> BoxFuture<'a, ()> { - if self.skip_set.contains(&plg.xtream_cluster) { - return Box::pin(async move {}); - } - self.inner.update_playlist(plg) - } - - fn get_missing_vod_info_count(&mut self) -> usize { - let skip_set = Arc::clone(&self.skip_set); - self.inner - .items() - .filter(move |item| !skip_set.contains(&item.as_ref().header.xtream_cluster)) - .filter(|item| { - let pli = item.as_ref(); - pli.header.xtream_cluster == XtreamCluster::Video - && pli.header.item_type == PlaylistItemType::Video - && !pli.has_details() - }) - .count() - } - - fn get_missing_series_info_count(&mut self) -> usize { - let skip_set = Arc::clone(&self.skip_set); - self.inner - .items() - .filter(move |item| !skip_set.contains(&item.as_ref().header.xtream_cluster)) - .filter(|item| { - let pli = item.as_ref(); - pli.header.xtream_cluster == XtreamCluster::Series - && pli.header.item_type == PlaylistItemType::SeriesInfo - && pli.header.id.parse::().is_ok_and(|id| id > 0) - && !pli.has_details() - }) - .count() - } - - fn deduplicate(&mut self, duplicates: &mut HashSet) { - if self.inner.is_memory() { - let filtered_groups = self - .inner - .take_groups() - .into_iter() - .filter_map(|group| self.filter_group(group)) - .collect::>(); - let mut memory = MemoryPlaylistSource::new(filtered_groups).boxed(); - memory.deduplicate(duplicates); - self.inner = memory; - return; - } - - self.inner.deduplicate(duplicates); - } - - fn take_groups(&mut self) -> Vec { - self.inner - .take_groups() - .into_iter() - .filter_map(|group| self.filter_group(group)) - .collect() - } - - fn clone_box(&self) -> Box { - Box::new(Self { - inner: self.inner.clone_box(), - skip_set: Arc::clone(&self.skip_set), - }) - } - - fn release_resources(&mut self, cluster: XtreamCluster) { - self.inner.release_resources(cluster); - } - - fn obtain_resources(&mut self) -> BoxFuture<'_, ()> { - self.inner.obtain_resources() - } - - fn sort_by_provider_ordinal(&mut self) { - self.inner.sort_by_provider_ordinal(); - } -} diff --git a/backend/src/processing/processor/library.rs b/backend/src/processing/processor/library.rs index ad32139ae..12e12b257 100644 --- a/backend/src/processing/processor/library.rs +++ b/backend/src/processing/processor/library.rs @@ -1,4 +1,4 @@ -use crate::library::{EpisodeMetadata, MediaMetadata, MetadataAsyncIter, MetadataCacheEntry, TechnicalMetadata}; +use crate::library::{Actor, EpisodeMetadata, MediaMetadata, MetadataAsyncIter, MetadataCacheEntry, TechnicalMetadata, VideoClipMetadata}; use crate::library::resolve_metadata_storage_path; use crate::model::{AppConfig, ConfigInput}; use shared::concat_string; @@ -41,6 +41,24 @@ fn technical_bitrate(technical: Option<&TechnicalMetadata>) -> u32 { .unwrap_or_default() } +fn join_actor_names(actors: &[Actor]) -> Arc { + let mut names = String::new(); + for (idx, actor) in actors.iter().enumerate() { + if idx > 0 { + names.push_str(", "); + } + names.push_str(&actor.name); + } + names.intern() +} + +fn youtube_trailer_key(videos: Option<&[VideoClipMetadata]>) -> Option<&str> { + videos? + .iter() + .find(|video| video.site.eq_ignore_ascii_case("youtube")) + .map(|video| video.key.as_str()) +} + pub async fn download_library_playlist(_client: &reqwest::Client, app_config: &Arc, input: &ConfigInput) -> (Vec, Vec) { let config = &*app_config.config.load(); let Some(library_config) = config.library.as_ref() else { return (vec![], vec![]) }; @@ -280,10 +298,11 @@ pub fn metadata_cache_entry_to_xtream_movie_info( .and_then(|s| s.to_str()) .map(ToString::to_string).unwrap_or_default(); - let actor_names = movie.actors.as_ref().map(|a| a.iter().map(|a| a.name.clone()).collect::>().join(", ").intern()); + let actor_names = movie.actors.as_deref().map(join_actor_names); let technical = movie.technical.as_ref(); let duration_secs = technical_duration_secs(technical).or_else(|| movie.runtime.map(|runtime| runtime * 60)); let duration = duration_secs.map(duration_secs_to_xtream_duration); + let youtube_trailer = youtube_trailer_key(movie.videos.as_deref()); let properties = VideoStreamProperties { name: movie.title.clone().into(), @@ -297,7 +316,7 @@ pub fn metadata_cache_entry_to_xtream_movie_info( rating: movie.rating, rating_5based: None, stream_type: Some("movie".intern()), - trailer: movie.videos.as_ref().and_then(|v| v.iter().find(|video| video.site.eq_ignore_ascii_case("youtube")).map(|video| video.key.clone().into())), + trailer: youtube_trailer.map(Into::into), tmdb: movie.tmdb_id, is_adult: 0, details: Some(VideoStreamDetailProperties { @@ -308,7 +327,7 @@ pub fn metadata_cache_entry_to_xtream_movie_info( release_date: movie.year.map(|y| format!("{y}-01-01").into()), episode_run_time: movie.runtime, director: movie.directors.as_ref().map(|d| d.join(", ").into()), - youtube_trailer: movie.videos.as_ref().and_then(|v| v.iter().find(|video| video.site.eq_ignore_ascii_case("youtube")).map(|video| video.key.clone().into())), + youtube_trailer: youtube_trailer.map(Into::into), actors: actor_names.clone(), cast: actor_names.clone(), genre: movie.genres.as_ref().map(|g| g.join(", ").into()), @@ -347,9 +366,9 @@ pub fn metadata_cache_entry_to_xtream_series_info( MediaMetadata::Series(m) => m, }; - let actor_names: Arc = series.actors.as_ref().map(|a| a.iter().map(|a| a.name.clone()).collect::>().join(", ")).unwrap_or_default().into(); + let actor_names: Arc = series.actors.as_deref().map(join_actor_names).unwrap_or_default(); let release_date = series.year.map(|y| format!("{y}-01-01")); - let youtube_trailer = series.videos.as_ref().and_then(|v| v.iter().find(|video| video.site.eq_ignore_ascii_case("youtube")).map(|video| video.key.clone())).unwrap_or_default(); + let youtube_trailer = youtube_trailer_key(series.videos.as_deref()).unwrap_or_default(); let series_thumbnail = thumbnail_url(entry, api_base_path); let series_art = series.poster.clone().or(series_thumbnail.clone()); diff --git a/backend/src/processing/processor/mod.rs b/backend/src/processing/processor/mod.rs index 0c4c039ad..7034c75d0 100644 --- a/backend/src/processing/processor/mod.rs +++ b/backend/src/processing/processor/mod.rs @@ -8,7 +8,6 @@ mod sort; mod trakt; mod library; mod stream_probe; -mod filtered_playlist_source; mod probe_handle_guard; mod resolve_options; pub use self::playlist::*; diff --git a/backend/src/processing/processor/playlist.rs b/backend/src/processing/processor/playlist.rs index 68f329316..7f27dbf55 100644 --- a/backend/src/processing/processor/playlist.rs +++ b/backend/src/processing/processor/playlist.rs @@ -1,4 +1,3 @@ -use super::filtered_playlist_source::FilteredPlaylistSource; use crate::{ api::{ model::{ @@ -66,12 +65,23 @@ use tokio::{ const PLAYLIST_UPDATE_MAX_DURATION_SECS: u64 = 3600; +fn join_arc_strs(values: &[Arc], separator: &str) -> String { + let mut result = String::new(); + for value in values { + if !result.is_empty() { + result.push_str(separator); + } + result.push_str(value.as_ref()); + } + result +} + fn is_valid(pli: &PlaylistItem, filter: &Filter, match_as_ascii: bool) -> bool { let provider = ValueProvider { pli, match_as_ascii }; filter.filter(&provider) } -pub fn apply_filter_to_source(source: &mut dyn PlaylistSource, filter: &Filter) -> Option> { +pub fn apply_filter_to_source(source: &mut PlaylistSource, filter: &Filter) -> Option> { let mut groups: IndexMap = IndexMap::new(); for pli in source.into_items() { if is_valid(&pli, filter, false) { @@ -100,7 +110,7 @@ pub fn apply_filter_to_source(source: &mut dyn PlaylistSource, filter: &Filter) } } -fn filter_playlist(source: &mut dyn PlaylistSource, target: &ConfigTarget) -> Option> { +fn filter_playlist(source: &mut PlaylistSource, target: &ConfigTarget) -> Option> { apply_filter_to_source(source, &target.filter) } @@ -152,14 +162,13 @@ fn exec_rename(pli: &mut PlaylistItem, rename: Option<&Vec>) { if log_enabled!(log::Level::Debug) && *value != *cap { trace_if_enabled!("Renamed {}={value} to {cap}", &r.field); } - let value = cap.into_owned(); - set_field_value(result, r.field, value); + set_field_value(result, r.field, cap.as_ref()); } } } } -fn rename_playlist(source: &mut dyn PlaylistSource, target: &ConfigTarget) -> Option> { +fn rename_playlist(source: &mut PlaylistSource, target: &ConfigTarget) -> Option> { match &target.rename { Some(renames) if !renames.is_empty() => { let mut groups: IndexMap<(XtreamCluster, Arc), PlaylistGroup> = IndexMap::new(); @@ -233,7 +242,7 @@ fn map_channel_and_flatten(channel: PlaylistItem, mapping: &Mapping) -> Vec Option> { +fn map_playlist(source: &mut PlaylistSource, target: &ConfigTarget) -> Option> { let mapping_binding = target.mapping.load(); let mappings = mapping_binding.as_ref()?; let valid_mappings = mappings.iter().filter(|m| m.mapper.as_ref().is_some_and(|v| !v.is_empty())); @@ -396,17 +405,14 @@ async fn download_plex_media_server_playlist( } } -fn filter_skipped_clusters_from_source( - source: Box, - input: &ConfigInput, -) -> Box { +fn filter_skipped_clusters_from_source(source: PlaylistSource, input: &ConfigInput) -> PlaylistSource { let skip_clusters = collect_effective_skip_clusters(input); if skip_clusters.is_empty() { return source; } let skip_set: HashSet = skip_clusters.into_iter().collect(); - Box::new(FilteredPlaylistSource::new(source, skip_set)) + PlaylistSource::filtered(source, skip_set) } #[allow(clippy::too_many_lines)] @@ -679,7 +685,7 @@ async fn process_source( if !disabled_inputs.is_empty() && !source_downloaded { warn!( "Source at index {source_idx} has no enabled inputs for the given targets. Disabled: {}", - disabled_inputs.iter().map(std::convert::AsRef::as_ref).collect::>().join(", ") + join_arc_strs(&disabled_inputs, ", ") ); } if source_downloaded { @@ -687,7 +693,7 @@ async fn process_source( debug!("Source at index {source_idx} is empty"); errors.push(TuliproxError::RepositoryPlaylist(format!( "Source at index {source_idx} is empty: {}", - source.inputs.iter().map(Clone::clone).collect::>>().join(", ") + join_arc_strs(&source.inputs, ", ") ))); } else { debug_if_enabled!( @@ -768,17 +774,17 @@ async fn invalidate_input_cache_status(ctx: &PlaylistProcessingContext, input: & async fn load_cached_input_playlist( ctx: &PlaylistProcessingContext, input: &Arc, -) -> (Box, Option) { +) -> (PlaylistSource, Option) { match load_input_playlist(ctx, input, None).await { Ok(pl_source) => (pl_source, None), - Err(err) => (MemoryPlaylistSource::default().boxed(), Some(err)), + Err(err) => (MemoryPlaylistSource::default().into_source(), Some(err)), } } async fn download_input( ctx: &PlaylistProcessingContext, input: &Arc, -) -> (Vec, Box, Option) { +) -> (Vec, PlaylistSource, Option) { // Coordination Logic let need_download = !ctx.is_input_downloaded(&input.name).await; // Keep this lock for the whole critical section (download + persist/load + mark processed) @@ -806,7 +812,7 @@ async fn download_input( PlaylistDownloadResult::new(vec![], vec![], true, false) }; - let mut preloaded_playlist: Option<(Box, Option)> = None; + let mut preloaded_playlist: Option<(PlaylistSource, Option)> = None; if playlist_download_result.was_cached { let (cached_playlist, cached_error) = load_cached_input_playlist(ctx, input).await; // Defensive fallback: if cache metadata says "valid" but persisted data is unreadable, @@ -852,12 +858,12 @@ async fn download_input( } else if playlist_download_result.was_cached || playlist_download_result.persisted { match load_input_playlist(ctx, input, None).await { Ok(pl_source) => (pl_source, None), - Err(e) => (MemoryPlaylistSource::default().boxed(), Some(e)), + Err(e) => (MemoryPlaylistSource::default().into_source(), Some(e)), } } else { debug!("Persisting input '{}' playlist", input.name); let (pl, err) = persist_input_playlist(&ctx.config, input, playlist_download_result.downloaded_playlist).await; - (MemoryPlaylistSource::new(pl).boxed(), err) + (MemoryPlaylistSource::new(pl).into_source(), err) }; let playlist = filter_skipped_clusters_from_source(playlist, input); @@ -1018,7 +1024,7 @@ async fn process_sources(processing_ctx: &PlaylistProcessingContext) -> (Vec Option>>; +pub type ProcessingPipe = Vec Option>>; fn get_processing_pipe(target: &ConfigTarget) -> ProcessingPipe { match &target.processing_order { @@ -1037,15 +1043,15 @@ fn execute_pipe<'a>( fpl: &mut FetchedPlaylist<'a>, duplicates: &mut HashSet, consume_source: bool, -) -> FetchedPlaylist<'a> { +) -> Result, TuliproxError> { let source = if consume_source { if fpl.is_memory() { - MemoryPlaylistSource::new(fpl.source.take_groups()).boxed() + MemoryPlaylistSource::new(fpl.source.take_groups()).into_source() } else { - std::mem::replace(&mut fpl.source, MemoryPlaylistSource::default().boxed()) + std::mem::replace(&mut fpl.source, MemoryPlaylistSource::default().into_source()) } } else { - fpl.clone_source() + fpl.clone_source()? }; let mut new_fpl = FetchedPlaylist { input: fpl.input, source, epg: fpl.epg.clone() }; @@ -1054,15 +1060,15 @@ fn execute_pipe<'a>( } for f in pipe { - if let Some(groups) = f(new_fpl.source.as_mut(), target) { - new_fpl.source = MemoryPlaylistSource::new(groups).boxed(); + if let Some(groups) = f(&mut new_fpl.source, target) { + new_fpl.source = MemoryPlaylistSource::new(groups).into_source(); } } // Ensure source is memory-based for downstream mutable processing (VOD/series resolution) if !new_fpl.is_memory() { - new_fpl.source = MemoryPlaylistSource::new(new_fpl.source.take_groups()).boxed(); + new_fpl.source = MemoryPlaylistSource::new(new_fpl.source.take_groups()).into_source(); } - new_fpl + Ok(new_fpl) } // This method is needed, because of duplicate group names in different inputs. @@ -1116,7 +1122,8 @@ async fn process_playlist_for_target( format!("target '{}' input '{}' before_pipe", target.name, provider_fpl.input.name).as_str(), ); step.broadcast("Executing transformations on '{}' playlist", &target.name); - let mut processed_fpl = execute_pipe(target, &pipe, provider_fpl, &mut duplicates, consume_input_source); + let mut processed_fpl = + execute_pipe(target, &pipe, provider_fpl, &mut duplicates, consume_input_source).map_err(|err| vec![err])?; log_memory_snapshot( format!("target '{}' input '{}' after_pipe", target.name, provider_fpl.input.name).as_str(), ); @@ -1830,7 +1837,7 @@ mod tests { }, ]; - let source = MemoryPlaylistSource::new(groups).boxed(); + let source = MemoryPlaylistSource::new(groups).into_source(); let input = ConfigInput { name: "skip_live".intern(), input_type: InputType::Xtream, diff --git a/backend/src/processing/processor/xtream_series.rs b/backend/src/processing/processor/xtream_series.rs index c0409c9a4..f02e65b22 100644 --- a/backend/src/processing/processor/xtream_series.rs +++ b/backend/src/processing/processor/xtream_series.rs @@ -18,9 +18,7 @@ use crate::processing::processor::{ FOREGROUND_RETRY_BATCH_MAX_SIZE as RETRY_BATCH_MAX_SIZE, }; use crate::repository::persists_input_series_info; -use crate::repository::{ - get_input_storage_path, persist_input_series_info_batch, MemoryPlaylistSource, PlaylistSource, -}; +use crate::repository::{get_input_storage_path, persist_input_series_info_batch, MemoryPlaylistSource}; use crate::repository::{xtream_get_file_path, BPlusTreeQuery}; use crate::utils::ffmpeg::{is_supported_probe_url, FfmpegExecutor, ProbeFailureKind, ProbeUrlOutcome}; use crate::utils::{debug_if_enabled, xtream}; @@ -124,7 +122,7 @@ async fn playlist_resolve_series_info( // Apply pipe transformations to new groups let mut new_playlist = groups_to_add; for f in pipe { - let mut source = MemoryPlaylistSource::new(new_playlist); + let mut source = MemoryPlaylistSource::new(new_playlist).into_source(); if let Some(v) = f(&mut source, target) { new_playlist = v; } else { diff --git a/backend/src/repository/bplustree.rs b/backend/src/repository/bplustree.rs index 83893d512..7a12026f5 100644 --- a/backend/src/repository/bplustree.rs +++ b/backend/src/repository/bplustree.rs @@ -3346,6 +3346,25 @@ where /// Returns the filepath this query was opened from. pub fn filepath(&self) -> &Path { &self.filepath } + #[cfg(test)] + pub(crate) fn clone_error_fixture() -> Self { + Self { + file: None, + mmap: None, + filepath: PathBuf::new(), + file_identity: None, + has_tombstones: false, + buffer: vec![0u8; PAGE_SIZE_USIZE], + cache: SplitNodeCache::new(), + node_cache: SplitNodeCache::new(), + last_refresh_at: Instant::now(), + refresh_interval: QUERY_REFRESH_INTERVAL, + root_offset: 0, + _marker_k: PhantomData, + _marker_v: PhantomData, + } + } + /// Force a header/root refresh, bypassing the automatic refresh throttle. pub fn refresh(&mut self) -> io::Result<()> { self.refresh_root_offset() } diff --git a/backend/src/repository/m3u_playlist_iterator.rs b/backend/src/repository/m3u_playlist_iterator.rs index 399d47b0b..063cb91bc 100644 --- a/backend/src/repository/m3u_playlist_iterator.rs +++ b/backend/src/repository/m3u_playlist_iterator.rs @@ -109,7 +109,7 @@ fn resolve_effective_source_url<'a>( "Failed to resolve provider URL '{}' for input '{}': {}", sanitize_sensitive_info(&m3u_pli.url), m3u_pli.input_name, - sanitize_sensitive_info(err.to_string().as_str()) + sanitize_sensitive_info(&err.to_string()) ); Cow::Borrowed(m3u_pli.url.as_ref()) } diff --git a/backend/src/repository/playlist_repository.rs b/backend/src/repository/playlist_repository.rs index 2028959c5..65042c70c 100644 --- a/backend/src/repository/playlist_repository.rs +++ b/backend/src/repository/playlist_repository.rs @@ -21,7 +21,7 @@ use shared::error::{ TuliproxError}; use shared::model::xtream_const::XTREAM_CLUSTER; use shared::model::{InputType, M3uPlaylistItem, PlaylistEntry, PlaylistGroup, PlaylistItem, PlaylistItemHeader, PlaylistItemType, SeriesStreamDetailEpisodeProperties, SeriesStreamDetailProperties, StreamProperties, VirtualId, XtreamCluster, XtreamPlaylistItem}; use shared::utils::{is_dash_url, is_hls_url, Internable}; -use std::collections::HashMap; +use std::collections::{HashMap, HashSet}; use std::path::{Path, PathBuf}; use std::sync::Arc; @@ -568,7 +568,7 @@ pub async fn persist_input_playlist(app_config: &Arc, input: &ConfigI } } -pub async fn load_input_playlist(ctx: &PlaylistProcessingContext, input: &ConfigInput, clusters: Option<&[XtreamCluster]>) -> Result, TuliproxError> { +pub async fn load_input_playlist(ctx: &PlaylistProcessingContext, input: &ConfigInput, clusters: Option<&[XtreamCluster]>) -> Result { let app_config = &ctx.config; let cfg = app_config.config.load(); let storage_path = get_input_storage_path(&input.name, &cfg.storage_dir).await @@ -577,49 +577,62 @@ pub async fn load_input_playlist(ctx: &PlaylistProcessingContext, input: &Config match input.get_download_input_type() { InputType::Xtream | InputType::XtreamBatch => { + let clusters_to_load = clusters.unwrap_or(&XTREAM_CLUSTER); if disk_based_processing { - Ok(Box::new(XtreamDiskPlaylistSource::new(app_config, &storage_path).await)) + let source = PlaylistSource::xtream_disk( + XtreamDiskPlaylistSource::new(app_config, &storage_path).await, + ); + Ok(PlaylistSource::filtered(source, skipped_clusters(clusters_to_load))) } else { - let clusters_to_load = if let Some(c) = clusters { - c - } else { - &XTREAM_CLUSTER - }; let groups = load_input_xtream_playlist(app_config, &storage_path, clusters_to_load).await?; - Ok(Box::new(MemoryPlaylistSource::new(groups))) + Ok(MemoryPlaylistSource::new(groups).into_source()) } } InputType::M3u | InputType::M3uBatch => { // Load M3U let file_path = get_input_m3u_playlist_file_path(&storage_path, &input.name); if disk_based_processing && file_path.exists() { - Ok(Box::new(M3uDiskPlaylistSource::new(app_config, &file_path).await)) + Ok(PlaylistSource::m3u_disk( + M3uDiskPlaylistSource::new(app_config, &file_path).await, + )) } else { let groups = load_input_m3u_playlist(app_config, &file_path).await?; - Ok(Box::new(MemoryPlaylistSource::new(groups))) + Ok(MemoryPlaylistSource::new(groups).into_source()) } } InputType::Library => { let file_path = get_input_local_library_playlist_file_path(&storage_path, &input.name); if disk_based_processing && file_path.exists() { - Ok(Box::new(LocalLibraryDiskPlaylistSource::new(app_config, &file_path).await)) + Ok(PlaylistSource::local_library_disk( + LocalLibraryDiskPlaylistSource::new(app_config, &file_path).await, + )) } else { let groups = load_input_local_library_playlist(app_config, &file_path).await?; - Ok(Box::new(MemoryPlaylistSource::new(groups))) + Ok(MemoryPlaylistSource::new(groups).into_source()) } } InputType::Emby | InputType::Jellyfin | InputType::Plex => { let file_path = get_input_media_server_playlist_file_path(&storage_path, &input.name); if disk_based_processing && file_path.exists() { - Ok(Box::new(MediaServerDiskPlaylistSource::new(app_config, &file_path).await)) + Ok(PlaylistSource::media_server_disk( + MediaServerDiskPlaylistSource::new(app_config, &file_path).await, + )) } else { let groups = load_input_media_server_playlist(app_config, &file_path).await?; - Ok(Box::new(MemoryPlaylistSource::new(groups))) + Ok(MemoryPlaylistSource::new(groups).into_source()) } } } } +fn skipped_clusters(clusters_to_load: &[XtreamCluster]) -> HashSet { + XTREAM_CLUSTER + .iter() + .copied() + .filter(|cluster| !clusters_to_load.contains(cluster)) + .collect() +} + pub fn get_input_m3u_playlist_file_path(storage_path: &Path, input_name: &Arc) -> PathBuf { let sanitized_input_name: String = input_name.chars() .map(|c| if c.is_alphanumeric() { c } else { '_' }) @@ -662,7 +675,7 @@ mod tests { assign_local_series_info_episode_key, assign_media_server_series_info_episode, get_input_media_server_playlist_file_path, materialize_media_server_series_info_episodes, rewrite_local_series_info_episode_virtual_id, rewrite_series_episode_parent_virtual_ids, - rewrite_series_info_episode_virtual_id, LocalEpisodeKey, ProviderEpisodeKey, + rewrite_series_info_episode_virtual_id, skipped_clusters, LocalEpisodeKey, ProviderEpisodeKey, }; use crate::repository::{BPlusTreeQuery, TargetIdMapping, VirtualIdRecord}; use shared::model::{ @@ -682,6 +695,14 @@ mod tests { assert!(path.ends_with("media_server_Media_Server_Input.db")); } + #[test] + fn skipped_clusters_converts_loaded_clusters_to_exclusions() { + let skipped = skipped_clusters(&[XtreamCluster::Live, XtreamCluster::Series]); + + assert_eq!(skipped.len(), 1); + assert!(skipped.contains(&XtreamCluster::Video)); + } + fn make_local_series_info(series_uuid: &str, episodes: Vec<(u32, &str, &str)>) -> PlaylistItem { let episode_props = episodes .into_iter() diff --git a/backend/src/repository/playlist_source.rs b/backend/src/repository/playlist_source.rs index b27364eba..7a44e8c0f 100644 --- a/backend/src/repository/playlist_source.rs +++ b/backend/src/repository/playlist_source.rs @@ -6,41 +6,415 @@ use futures::future::BoxFuture; use indexmap::IndexMap; use log::{error, warn}; use serde::{Deserialize, Serialize}; +use shared::error::TuliproxError; use shared::model::{M3uPlaylistItem, PlaylistEntry, PlaylistGroup, PlaylistItem, PlaylistItemType, XtreamCluster, XtreamPlaylistItem}; +use shared::model::UUIDType; +use shared::utils::Internable; use std::borrow::Cow; use std::collections::HashSet; use std::path::{Path, PathBuf}; use std::sync::Arc; -use shared::model::UUIDType; -use shared::utils::Internable; -pub trait PlaylistSource: Send + Sync { +trait PlaylistSourceOps: Send + Sync { fn is_memory(&self) -> bool; fn get_channel_count(&mut self) -> usize; fn get_group_count(&mut self) -> usize; + fn get_channel_count_excluding_clusters(&mut self, skip_set: &HashSet) -> usize; + fn get_group_count_excluding_clusters(&mut self, skip_set: &HashSet) -> usize; fn is_empty(&mut self) -> bool; #[allow(clippy::wrong_self_convention)] - fn into_items(&mut self) -> Box + Send + '_>; - fn items_mut(&mut self) -> Box + Send + '_>; - fn items<'a>(&'a mut self) -> Box> + Send + 'a>; + fn into_items(&mut self) -> Box + Send + '_>; + fn items_mut(&mut self) -> Box + Send + '_>; + fn items<'a>(&'a mut self) -> Box> + Send + 'a>; fn update_playlist<'a>(&'a mut self, plg: &'a PlaylistGroup) -> BoxFuture<'a, ()>; fn get_missing_vod_info_count(&mut self) -> usize; fn get_missing_series_info_count(&mut self) -> usize; + fn get_missing_vod_info_count_excluding_clusters(&mut self, skip_set: &HashSet) -> usize; + fn get_missing_series_info_count_excluding_clusters(&mut self, skip_set: &HashSet) -> usize; fn deduplicate(&mut self, duplicates: &mut HashSet); fn take_groups(&mut self) -> Vec; - fn clone_box(&self) -> Box; + fn clone_source(&self) -> Result; fn release_resources(&mut self, cluster: XtreamCluster); fn obtain_resources(&mut self) -> BoxFuture<'_, ()>; fn sort_by_provider_ordinal(&mut self); } +pub struct PlaylistSource { + kind: PlaylistSourceKind, + skip_set: Option>>, +} + +enum PlaylistSourceKind { + Empty(EmptyPlaylistSource), + XtreamDisk(Box), + M3uDisk(Box), + LocalLibraryDisk(Box), + MediaServerDisk(Box), + Memory(MemoryPlaylistSource), +} + +type XtreamQueryHandle = (BPlusTreeQuery, Arc); + +impl Default for PlaylistSource { + fn default() -> Self { Self::new(PlaylistSourceKind::Empty(EmptyPlaylistSource::default())) } +} + +impl PlaylistSource { + fn new(kind: PlaylistSourceKind) -> Self { Self { kind, skip_set: None } } + + pub fn xtream_disk(source: XtreamDiskPlaylistSource) -> Self { Self::new(PlaylistSourceKind::XtreamDisk(Box::new(source))) } + + pub fn m3u_disk(source: M3uDiskPlaylistSource) -> Self { Self::new(PlaylistSourceKind::M3uDisk(Box::new(source))) } + + pub fn local_library_disk(source: LocalLibraryDiskPlaylistSource) -> Self { + Self::new(PlaylistSourceKind::LocalLibraryDisk(Box::new(source))) + } + + pub fn media_server_disk(source: MediaServerDiskPlaylistSource) -> Self { + Self::new(PlaylistSourceKind::MediaServerDisk(Box::new(source))) + } + + pub fn memory(source: MemoryPlaylistSource) -> Self { Self::new(PlaylistSourceKind::Memory(source)) } + + pub fn filtered(mut inner: Self, skip_set: HashSet) -> Self { + if skip_set.is_empty() { + return inner; + } + + if let Some(existing) = &inner.skip_set { + let mut merged = HashSet::with_capacity(existing.len() + skip_set.len()); + merged.extend(existing.iter().copied()); + merged.extend(skip_set); + inner.skip_set = Some(Arc::new(merged)); + } else { + inner.skip_set = Some(Arc::new(skip_set)); + } + + inner + } + + pub fn is_memory(&self) -> bool { + match &self.kind { + PlaylistSourceKind::Empty(source) => source.is_memory(), + PlaylistSourceKind::XtreamDisk(source) => source.is_memory(), + PlaylistSourceKind::M3uDisk(source) => source.is_memory(), + PlaylistSourceKind::LocalLibraryDisk(source) => source.is_memory(), + PlaylistSourceKind::MediaServerDisk(source) => source.is_memory(), + PlaylistSourceKind::Memory(source) => source.is_memory(), + } + } + + pub fn get_channel_count(&mut self) -> usize { + if let Some(skip_set) = self.skip_set.as_ref() { + return match &mut self.kind { + PlaylistSourceKind::Empty(source) => source.get_channel_count_excluding_clusters(skip_set), + PlaylistSourceKind::XtreamDisk(source) => source.get_channel_count_excluding_clusters(skip_set), + PlaylistSourceKind::M3uDisk(source) => source.get_channel_count_excluding_clusters(skip_set), + PlaylistSourceKind::LocalLibraryDisk(source) => source.get_channel_count_excluding_clusters(skip_set), + PlaylistSourceKind::MediaServerDisk(source) => source.get_channel_count_excluding_clusters(skip_set), + PlaylistSourceKind::Memory(source) => source.get_channel_count_excluding_clusters(skip_set), + }; + } + + match &mut self.kind { + PlaylistSourceKind::Empty(source) => source.get_channel_count(), + PlaylistSourceKind::XtreamDisk(source) => source.get_channel_count(), + PlaylistSourceKind::M3uDisk(source) => source.get_channel_count(), + PlaylistSourceKind::LocalLibraryDisk(source) => source.get_channel_count(), + PlaylistSourceKind::MediaServerDisk(source) => source.get_channel_count(), + PlaylistSourceKind::Memory(source) => source.get_channel_count(), + } + } + + pub fn get_group_count(&mut self) -> usize { + if let Some(skip_set) = self.skip_set.as_ref() { + return match &mut self.kind { + PlaylistSourceKind::Empty(source) => source.get_group_count_excluding_clusters(skip_set), + PlaylistSourceKind::XtreamDisk(source) => source.get_group_count_excluding_clusters(skip_set), + PlaylistSourceKind::M3uDisk(source) => source.get_group_count_excluding_clusters(skip_set), + PlaylistSourceKind::LocalLibraryDisk(source) => source.get_group_count_excluding_clusters(skip_set), + PlaylistSourceKind::MediaServerDisk(source) => source.get_group_count_excluding_clusters(skip_set), + PlaylistSourceKind::Memory(source) => source.get_group_count_excluding_clusters(skip_set), + }; + } + + match &mut self.kind { + PlaylistSourceKind::Empty(source) => source.get_group_count(), + PlaylistSourceKind::XtreamDisk(source) => source.get_group_count(), + PlaylistSourceKind::M3uDisk(source) => source.get_group_count(), + PlaylistSourceKind::LocalLibraryDisk(source) => source.get_group_count(), + PlaylistSourceKind::MediaServerDisk(source) => source.get_group_count(), + PlaylistSourceKind::Memory(source) => source.get_group_count(), + } + } + + pub fn is_empty(&mut self) -> bool { + if self.skip_set.is_some() { + return self.get_channel_count() == 0; + } + + match &mut self.kind { + PlaylistSourceKind::Empty(source) => source.is_empty(), + PlaylistSourceKind::XtreamDisk(source) => source.is_empty(), + PlaylistSourceKind::M3uDisk(source) => source.is_empty(), + PlaylistSourceKind::LocalLibraryDisk(source) => source.is_empty(), + PlaylistSourceKind::MediaServerDisk(source) => source.is_empty(), + PlaylistSourceKind::Memory(source) => source.is_empty(), + } + } + + #[allow(clippy::wrong_self_convention)] + pub fn into_items(&mut self) -> Box + Send + '_> { + let iter = match &mut self.kind { + PlaylistSourceKind::Empty(source) => source.into_items(), + PlaylistSourceKind::XtreamDisk(source) => source.into_items(), + PlaylistSourceKind::M3uDisk(source) => source.into_items(), + PlaylistSourceKind::LocalLibraryDisk(source) => source.into_items(), + PlaylistSourceKind::MediaServerDisk(source) => source.into_items(), + PlaylistSourceKind::Memory(source) => source.into_items(), + }; + + if let Some(skip_set) = self.skip_set.clone() { + Box::new(iter.filter(move |item| !skip_set.contains(&item.header.xtream_cluster))) + } else { + iter + } + } + + pub fn items_mut(&mut self) -> Box + Send + '_> { + let iter = match &mut self.kind { + PlaylistSourceKind::Empty(source) => source.items_mut(), + PlaylistSourceKind::XtreamDisk(source) => source.items_mut(), + PlaylistSourceKind::M3uDisk(source) => source.items_mut(), + PlaylistSourceKind::LocalLibraryDisk(source) => source.items_mut(), + PlaylistSourceKind::MediaServerDisk(source) => source.items_mut(), + PlaylistSourceKind::Memory(source) => source.items_mut(), + }; + + if let Some(skip_set) = self.skip_set.clone() { + Box::new(iter.filter_map(move |item| { + if skip_set.contains(&item.header.xtream_cluster) { + None + } else { + Some(item) + } + })) + } else { + iter + } + } + + pub fn items<'a>(&'a mut self) -> Box> + Send + 'a> { + let iter = match &mut self.kind { + PlaylistSourceKind::Empty(source) => source.items(), + PlaylistSourceKind::XtreamDisk(source) => source.items(), + PlaylistSourceKind::M3uDisk(source) => source.items(), + PlaylistSourceKind::LocalLibraryDisk(source) => source.items(), + PlaylistSourceKind::MediaServerDisk(source) => source.items(), + PlaylistSourceKind::Memory(source) => source.items(), + }; + + if let Some(skip_set) = self.skip_set.clone() { + Box::new(iter.filter(move |item| !skip_set.contains(&item.as_ref().header.xtream_cluster))) + } else { + iter + } + } + + pub fn update_playlist<'a>(&'a mut self, plg: &'a PlaylistGroup) -> BoxFuture<'a, ()> { + if self + .skip_set + .as_ref() + .is_some_and(|skip_set| skip_set.contains(&plg.xtream_cluster)) + { + return Box::pin(async move {}); + } + + match &mut self.kind { + PlaylistSourceKind::Empty(source) => source.update_playlist(plg), + PlaylistSourceKind::XtreamDisk(source) => source.update_playlist(plg), + PlaylistSourceKind::M3uDisk(source) => source.update_playlist(plg), + PlaylistSourceKind::LocalLibraryDisk(source) => source.update_playlist(plg), + PlaylistSourceKind::MediaServerDisk(source) => source.update_playlist(plg), + PlaylistSourceKind::Memory(source) => source.update_playlist(plg), + } + } + + pub fn get_missing_vod_info_count(&mut self) -> usize { + if let Some(skip_set) = self.skip_set.as_ref() { + return match &mut self.kind { + PlaylistSourceKind::Empty(source) => source.get_missing_vod_info_count_excluding_clusters(skip_set), + PlaylistSourceKind::XtreamDisk(source) => source.get_missing_vod_info_count_excluding_clusters(skip_set), + PlaylistSourceKind::M3uDisk(source) => source.get_missing_vod_info_count_excluding_clusters(skip_set), + PlaylistSourceKind::LocalLibraryDisk(source) => { + source.get_missing_vod_info_count_excluding_clusters(skip_set) + } + PlaylistSourceKind::MediaServerDisk(source) => source.get_missing_vod_info_count_excluding_clusters(skip_set), + PlaylistSourceKind::Memory(source) => source.get_missing_vod_info_count_excluding_clusters(skip_set), + }; + } + + match &mut self.kind { + PlaylistSourceKind::Empty(source) => source.get_missing_vod_info_count(), + PlaylistSourceKind::XtreamDisk(source) => source.get_missing_vod_info_count(), + PlaylistSourceKind::M3uDisk(source) => source.get_missing_vod_info_count(), + PlaylistSourceKind::LocalLibraryDisk(source) => source.get_missing_vod_info_count(), + PlaylistSourceKind::MediaServerDisk(source) => source.get_missing_vod_info_count(), + PlaylistSourceKind::Memory(source) => source.get_missing_vod_info_count(), + } + } + + pub fn get_missing_series_info_count(&mut self) -> usize { + if let Some(skip_set) = self.skip_set.as_ref() { + return match &mut self.kind { + PlaylistSourceKind::Empty(source) => source.get_missing_series_info_count_excluding_clusters(skip_set), + PlaylistSourceKind::XtreamDisk(source) => { + source.get_missing_series_info_count_excluding_clusters(skip_set) + } + PlaylistSourceKind::M3uDisk(source) => source.get_missing_series_info_count_excluding_clusters(skip_set), + PlaylistSourceKind::LocalLibraryDisk(source) => { + source.get_missing_series_info_count_excluding_clusters(skip_set) + } + PlaylistSourceKind::MediaServerDisk(source) => { + source.get_missing_series_info_count_excluding_clusters(skip_set) + } + PlaylistSourceKind::Memory(source) => source.get_missing_series_info_count_excluding_clusters(skip_set), + }; + } + + match &mut self.kind { + PlaylistSourceKind::Empty(source) => source.get_missing_series_info_count(), + PlaylistSourceKind::XtreamDisk(source) => source.get_missing_series_info_count(), + PlaylistSourceKind::M3uDisk(source) => source.get_missing_series_info_count(), + PlaylistSourceKind::LocalLibraryDisk(source) => source.get_missing_series_info_count(), + PlaylistSourceKind::MediaServerDisk(source) => source.get_missing_series_info_count(), + PlaylistSourceKind::Memory(source) => source.get_missing_series_info_count(), + } + } + + pub fn deduplicate(&mut self, duplicates: &mut HashSet) { + if self.skip_set.is_some() && self.is_memory() { + let mut memory = MemoryPlaylistSource::new(self.take_groups()).into_source(); + memory.deduplicate(duplicates); + *self = memory; + return; + } + + match &mut self.kind { + PlaylistSourceKind::Empty(source) => source.deduplicate(duplicates), + PlaylistSourceKind::XtreamDisk(source) => source.deduplicate(duplicates), + PlaylistSourceKind::M3uDisk(source) => source.deduplicate(duplicates), + PlaylistSourceKind::LocalLibraryDisk(source) => source.deduplicate(duplicates), + PlaylistSourceKind::MediaServerDisk(source) => source.deduplicate(duplicates), + PlaylistSourceKind::Memory(source) => source.deduplicate(duplicates), + } + } + + pub fn take_groups(&mut self) -> Vec { + let groups = match &mut self.kind { + PlaylistSourceKind::Empty(source) => source.take_groups(), + PlaylistSourceKind::XtreamDisk(source) => source.take_groups(), + PlaylistSourceKind::M3uDisk(source) => source.take_groups(), + PlaylistSourceKind::LocalLibraryDisk(source) => source.take_groups(), + PlaylistSourceKind::MediaServerDisk(source) => source.take_groups(), + PlaylistSourceKind::Memory(source) => source.take_groups(), + }; + + if let Some(skip_set) = self.skip_set.clone() { + groups.into_iter().filter_map(|group| filter_group(&skip_set, group)).collect() + } else { + groups + } + } + + pub fn clone_source(&self) -> Result { + let mut cloned = match &self.kind { + PlaylistSourceKind::Empty(source) => source.clone_source(), + PlaylistSourceKind::XtreamDisk(source) => source.clone_source(), + PlaylistSourceKind::M3uDisk(source) => source.clone_source(), + PlaylistSourceKind::LocalLibraryDisk(source) => source.clone_source(), + PlaylistSourceKind::MediaServerDisk(source) => source.clone_source(), + PlaylistSourceKind::Memory(source) => source.clone_source(), + }?; + cloned.skip_set.clone_from(&self.skip_set); + Ok(cloned) + } + + pub fn release_resources(&mut self, cluster: XtreamCluster) { + match &mut self.kind { + PlaylistSourceKind::Empty(source) => source.release_resources(cluster), + PlaylistSourceKind::XtreamDisk(source) => source.release_resources(cluster), + PlaylistSourceKind::M3uDisk(source) => source.release_resources(cluster), + PlaylistSourceKind::LocalLibraryDisk(source) => source.release_resources(cluster), + PlaylistSourceKind::MediaServerDisk(source) => source.release_resources(cluster), + PlaylistSourceKind::Memory(source) => source.release_resources(cluster), + } + } + + pub fn obtain_resources(&mut self) -> BoxFuture<'_, ()> { + match &mut self.kind { + PlaylistSourceKind::Empty(source) => source.obtain_resources(), + PlaylistSourceKind::XtreamDisk(source) => source.obtain_resources(), + PlaylistSourceKind::M3uDisk(source) => source.obtain_resources(), + PlaylistSourceKind::LocalLibraryDisk(source) => source.obtain_resources(), + PlaylistSourceKind::MediaServerDisk(source) => source.obtain_resources(), + PlaylistSourceKind::Memory(source) => source.obtain_resources(), + } + } + + pub fn sort_by_provider_ordinal(&mut self) { + match &mut self.kind { + PlaylistSourceKind::Empty(source) => source.sort_by_provider_ordinal(), + PlaylistSourceKind::XtreamDisk(source) => source.sort_by_provider_ordinal(), + PlaylistSourceKind::M3uDisk(source) => source.sort_by_provider_ordinal(), + PlaylistSourceKind::LocalLibraryDisk(source) => source.sort_by_provider_ordinal(), + PlaylistSourceKind::MediaServerDisk(source) => source.sort_by_provider_ordinal(), + PlaylistSourceKind::Memory(source) => source.sort_by_provider_ordinal(), + } + } +} + +fn filter_group(skip_set: &HashSet, mut group: PlaylistGroup) -> Option { + if skip_set.contains(&group.xtream_cluster) { + return None; + } + group + .channels + .retain(|item| !skip_set.contains(&item.header.xtream_cluster)); + if group.channels.is_empty() { + None + } else { + Some(group) + } +} + +fn cluster_from_item_type(item_type: PlaylistItemType) -> XtreamCluster { + XtreamCluster::try_from(item_type).unwrap_or(XtreamCluster::Live) +} + +fn clone_xtream_query( + label: &str, + source: Option<&XtreamQueryHandle>, +) -> Result, TuliproxError> { + source + .map(|(query, guard)| { + query + .try_clone() + .map(|cloned_query| (cloned_query, Arc::clone(guard))) + .map_err(|err| TuliproxError::RepositoryPlaylist(format!("Failed to clone {label} disk playlist query: {err}"))) + }) + .transpose() +} + #[derive(Default)] pub struct EmptyPlaylistSource {} -impl PlaylistSource for EmptyPlaylistSource { +impl PlaylistSourceOps for EmptyPlaylistSource { fn is_memory(&self) -> bool { true } fn get_channel_count(&mut self) -> usize { 0 } fn get_group_count(&mut self) -> usize { 0 } + fn get_channel_count_excluding_clusters(&mut self, _skip_set: &HashSet) -> usize { 0 } + fn get_group_count_excluding_clusters(&mut self, _skip_set: &HashSet) -> usize { 0 } fn is_empty(&mut self) -> bool { true } fn into_items(&mut self) -> Box + Send + '_> { Box::new(std::iter::empty()) } fn items_mut(&mut self) -> Box + Send + '_> { Box::new(std::iter::empty()) } @@ -48,9 +422,13 @@ impl PlaylistSource for EmptyPlaylistSource { fn update_playlist<'a>(&'a mut self, _plg: &'a PlaylistGroup) -> BoxFuture<'a, ()> { Box::pin(async move {}) } fn get_missing_vod_info_count(&mut self) -> usize { 0 } fn get_missing_series_info_count(&mut self) -> usize { 0 } + fn get_missing_vod_info_count_excluding_clusters(&mut self, _skip_set: &HashSet) -> usize { 0 } + fn get_missing_series_info_count_excluding_clusters(&mut self, _skip_set: &HashSet) -> usize { 0 } fn deduplicate(&mut self, _duplicates: &mut HashSet) { /* noop */ } fn take_groups(&mut self) -> Vec { vec![] } - fn clone_box(&self) -> Box { Box::new(EmptyPlaylistSource::default()) } + fn clone_source(&self) -> Result { + Ok(PlaylistSource::new(PlaylistSourceKind::Empty(EmptyPlaylistSource::default()))) + } fn release_resources(&mut self, _cluster: XtreamCluster) { /* noop */ } fn obtain_resources(&mut self) -> BoxFuture<'_, ()> { Box::pin(async move {}) } fn sort_by_provider_ordinal(&mut self) { /* noop */ } @@ -59,9 +437,9 @@ impl PlaylistSource for EmptyPlaylistSource { pub struct XtreamDiskPlaylistSource { app_config: Arc, storage_path: PathBuf, - live: Option<(BPlusTreeQuery, Arc)>, - vod: Option<(BPlusTreeQuery, Arc)>, - series: Option<(BPlusTreeQuery, Arc)>, + live: Option, + vod: Option, + series: Option, } impl XtreamDiskPlaylistSource { @@ -97,7 +475,7 @@ impl XtreamDiskPlaylistSource { } } -impl PlaylistSource for XtreamDiskPlaylistSource { +impl PlaylistSourceOps for XtreamDiskPlaylistSource { fn is_memory(&self) -> bool { false } fn get_channel_count(&mut self) -> usize { @@ -118,6 +496,49 @@ impl PlaylistSource for XtreamDiskPlaylistSource { groups.len() } + fn get_channel_count_excluding_clusters(&mut self, skip_set: &HashSet) -> usize { + let live_count = if skip_set.contains(&XtreamCluster::Live) { + 0 + } else { + self.live.as_mut().map_or(0usize, |(query, _)| query.len().unwrap_or(0usize)) + }; + let vod_count = if skip_set.contains(&XtreamCluster::Video) { + 0 + } else { + self.vod.as_mut().map_or(0usize, |(query, _)| query.len().unwrap_or(0usize)) + }; + let series_count = if skip_set.contains(&XtreamCluster::Series) { + 0 + } else { + self.series.as_mut().map_or(0usize, |(query, _)| query.len().unwrap_or(0usize)) + }; + live_count + vod_count + series_count + } + + fn get_group_count_excluding_clusters(&mut self, skip_set: &HashSet) -> usize { + fn collect_groups( + cluster: XtreamCluster, + query: &mut Option<(BPlusTreeQuery, Q)>, + groups: &mut HashSet<(XtreamCluster, Arc)>, + skip_set: &HashSet, + ) { + if skip_set.contains(&cluster) { + return; + } + if let Some((query, _)) = query { + for (_, item) in query.iter() { + groups.insert((cluster, Arc::clone(&item.group))); + } + } + } + + let mut groups = HashSet::new(); + collect_groups(XtreamCluster::Live, &mut self.live, &mut groups, skip_set); + collect_groups(XtreamCluster::Video, &mut self.vod, &mut groups, skip_set); + collect_groups(XtreamCluster::Series, &mut self.series, &mut groups, skip_set); + groups.len() + } + fn is_empty(&mut self) -> bool { self.live.as_mut().is_none_or(|(q, _)| q.is_empty().unwrap_or(true)) && self.vod.as_mut().is_none_or(|(q, _)| q.is_empty().unwrap_or(true)) @@ -180,6 +601,22 @@ impl PlaylistSource for XtreamDiskPlaylistSource { }) } + fn get_missing_vod_info_count_excluding_clusters(&mut self, skip_set: &HashSet) -> usize { + if skip_set.contains(&XtreamCluster::Video) { + 0 + } else { + self.get_missing_vod_info_count() + } + } + + fn get_missing_series_info_count_excluding_clusters(&mut self, skip_set: &HashSet) -> usize { + if skip_set.contains(&XtreamCluster::Series) { + 0 + } else { + self.get_missing_series_info_count() + } + } + fn deduplicate(&mut self, _duplicates: &mut HashSet) { warn!("Deduplication is not supported for disk based playlist updates"); } @@ -224,42 +661,18 @@ impl PlaylistSource for XtreamDiskPlaylistSource { groups } - fn clone_box(&self) -> Box { - let live = self.live.as_ref().and_then(|(query, guard)| { - match query.try_clone() { - Ok(q) => Some((q, Arc::clone(guard))), - Err(err) => { - warn!("PlaylistSource::clone_box failed to clone live query: {err}"); - None - } - } - }); - let vod = self.vod.as_ref().and_then(|(query, guard)| { - match query.try_clone() { - Ok(q) => Some((q, Arc::clone(guard))), - Err(err) => { - warn!("PlaylistSource::clone_box failed to clone vod query: {err}"); - None - } - } - }); - let series = self.series.as_ref().and_then(|(query, guard)| { - match query.try_clone() { - Ok(q) => Some((q, Arc::clone(guard))), - Err(err) => { - warn!("PlaylistSource::clone_box failed to clone series query: {err}"); - None - } - } - }); + fn clone_source(&self) -> Result { + let live = clone_xtream_query("live", self.live.as_ref())?; + let vod = clone_xtream_query("vod", self.vod.as_ref())?; + let series = clone_xtream_query("series", self.series.as_ref())?; - Box::new(Self { + Ok(PlaylistSource::xtream_disk(Self { app_config: Arc::clone(&self.app_config), storage_path: self.storage_path.clone(), live, vod, series, - }) + })) } fn release_resources(&mut self, cluster: XtreamCluster) { @@ -312,7 +725,7 @@ macro_rules! impl_single_file_disk_source { } } - impl PlaylistSource for [<$name DiskPlaylistSource>] { + impl PlaylistSourceOps for [<$name DiskPlaylistSource>] { fn get_channel_count(&mut self) -> usize { self.playlist.as_mut().map_or(0usize, |t: &mut BPlusTreeQuery<$key_type, $entry_type>| t.len().unwrap_or(0usize)) } @@ -324,6 +737,28 @@ macro_rules! impl_single_file_disk_source { groups.len() } + fn get_channel_count_excluding_clusters(&mut self, skip_set: &HashSet) -> usize { + self.playlist.as_mut().map_or(0, |query| { + query + .iter() + .filter(|(_, item)| !skip_set.contains(&cluster_from_item_type(item.item_type))) + .count() + }) + } + + fn get_group_count_excluding_clusters(&mut self, skip_set: &HashSet) -> usize { + let mut groups = HashSet::<(XtreamCluster, Arc)>::new(); + if let Some(query) = self.playlist.as_mut() { + for (_, item) in query.iter() { + let cluster = cluster_from_item_type(item.item_type); + if !skip_set.contains(&cluster) { + groups.insert((cluster, Arc::clone(&item.group))); + } + } + } + groups.len() + } + fn is_empty(&mut self) -> bool { self.playlist.as_mut().map_or(true, |t| t.is_empty().unwrap_or(true)) } fn into_items(&mut self) -> Box + Send + '_> { @@ -356,6 +791,8 @@ macro_rules! impl_single_file_disk_source { fn get_missing_vod_info_count(&mut self) -> usize { 0 } fn get_missing_series_info_count(&mut self) -> usize { 0 } + fn get_missing_vod_info_count_excluding_clusters(&mut self, _skip_set: &HashSet) -> usize { 0 } + fn get_missing_series_info_count_excluding_clusters(&mut self, _skip_set: &HashSet) -> usize { 0 } fn deduplicate(&mut self, _duplicates: &mut HashSet) { warn!("Deduplication is not supported for disk based playlist updates"); } @@ -391,27 +828,27 @@ macro_rules! impl_single_file_disk_source { vec![] } } - fn clone_box(&self) -> Box { - let playlist = self.playlist.as_ref().and_then(|query| { - match query.try_clone() { - Ok(cloned_query) => Some(cloned_query), - Err(err) => { - warn!( - "PlaylistSource::clone_box failed to clone {} disk query {}: {err}", + fn clone_source(&self) -> Result { + let playlist = self + .playlist + .as_ref() + .map(|query| { + query.try_clone().map_err(|err| { + TuliproxError::RepositoryPlaylist(format!( + "Failed to clone {} disk playlist query {}: {err}", stringify!($name), self.file_path.display() - ); - None - } - } - }); + )) + }) + }) + .transpose()?; - Box::new(Self { + Ok(PlaylistSource::[<$name:snake _disk>](Self { app_config: Arc::clone(&self.app_config), file_path: self.file_path.clone(), playlist, guard: self.guard.clone(), - }) + })) } fn release_resources(&mut self, _cluster: XtreamCluster) { @@ -447,9 +884,7 @@ impl MemoryPlaylistSource { Self { playlist: Arc::new(groups) } } - pub fn boxed(self) -> Box { - Box::new(self) - } + pub fn into_source(self) -> PlaylistSource { PlaylistSource::memory(self) } } impl Default for MemoryPlaylistSource { @@ -458,10 +893,30 @@ impl Default for MemoryPlaylistSource { } } -impl PlaylistSource for MemoryPlaylistSource { +impl PlaylistSourceOps for MemoryPlaylistSource { fn is_memory(&self) -> bool { true } fn get_channel_count(&mut self) -> usize { self.playlist.iter().map(|group| group.channels.len()).sum() } fn get_group_count(&mut self) -> usize { self.playlist.len() } + fn get_channel_count_excluding_clusters(&mut self, skip_set: &HashSet) -> usize { + self.playlist + .iter() + .filter(|group| !skip_set.contains(&group.xtream_cluster)) + .map(|group| { + group + .channels + .iter() + .filter(|item| !skip_set.contains(&item.header.xtream_cluster)) + .count() + }) + .sum() + } + fn get_group_count_excluding_clusters(&mut self, skip_set: &HashSet) -> usize { + self.playlist + .iter() + .filter(|group| !skip_set.contains(&group.xtream_cluster)) + .filter(|group| group.channels.iter().any(|item| !skip_set.contains(&item.header.xtream_cluster))) + .count() + } fn is_empty(&mut self) -> bool { self.playlist.is_empty() } fn into_items(&mut self) -> Box + Send + '_> { let playlist = Arc::make_mut(&mut self.playlist); @@ -514,6 +969,20 @@ impl PlaylistSource for MemoryPlaylistSource { && pli.get_provider_id().is_some_and(|id| id > 0) && !pli.has_details()).count() } + fn get_missing_vod_info_count_excluding_clusters(&mut self, skip_set: &HashSet) -> usize { + if skip_set.contains(&XtreamCluster::Video) { + 0 + } else { + self.get_missing_vod_info_count() + } + } + fn get_missing_series_info_count_excluding_clusters(&mut self, skip_set: &HashSet) -> usize { + if skip_set.contains(&XtreamCluster::Series) { + 0 + } else { + self.get_missing_series_info_count() + } + } fn deduplicate(&mut self, duplicates: &mut HashSet) { let playlist = Arc::make_mut(&mut self.playlist); for group in playlist { @@ -523,8 +992,8 @@ impl PlaylistSource for MemoryPlaylistSource { fn take_groups(&mut self) -> Vec { std::mem::take(Arc::make_mut(&mut self.playlist)) } - fn clone_box(&self) -> Box { - Box::new(MemoryPlaylistSource { playlist: Arc::clone(&self.playlist) }) + fn clone_source(&self) -> Result { + Ok(PlaylistSource::memory(MemoryPlaylistSource { playlist: Arc::clone(&self.playlist) })) } fn release_resources(&mut self, _cluster: XtreamCluster) { /* noop */ } fn obtain_resources(&mut self) -> BoxFuture<'_, ()> { Box::pin(async move {}) } @@ -577,17 +1046,62 @@ where #[cfg(test)] mod tests { - use super::{MemoryPlaylistSource, PlaylistGroup, PlaylistItem, PlaylistSource, XtreamCluster}; - use shared::model::PlaylistItemHeader; + use super::{MemoryPlaylistSource, PlaylistGroup, PlaylistItem, PlaylistSource, XtreamCluster, XtreamDiskPlaylistSource}; + use crate::model::{AppConfig, Config, MediaToolCapabilities, SourcesConfig}; + use crate::repository::BPlusTreeQuery; + use crate::utils::FileLockManager; + use arc_swap::{ArcSwap, ArcSwapOption}; + use shared::error::TuliproxError; + use shared::model::{ConfigPaths, PlaylistItemHeader, PlaylistItemType, XtreamPlaylistItem}; use shared::utils::Internable; + use std::collections::HashSet; + use std::path::PathBuf; + use std::sync::Arc; + + fn test_app_config() -> Arc { + Arc::new(AppConfig { + config: Arc::new(ArcSwap::from_pointee(Config::default())), + sources: Arc::new(ArcSwap::from_pointee(SourcesConfig { + batch_files: Vec::new(), + templates: None, + provider: Vec::new(), + inputs: Vec::new(), + sources: Vec::new(), + })), + hdhomerun: Arc::new(ArcSwapOption::empty()), + api_proxy: Arc::new(ArcSwapOption::empty()), + file_locks: Arc::new(FileLockManager::default()), + paths: Arc::new(ArcSwap::from_pointee(ConfigPaths { + home_path: String::new(), + config_path: String::new(), + storage_path: String::new(), + config_file_path: String::new(), + sources_file_path: String::new(), + mapping_file_path: None, + mapping_files_used: None, + template_file_path: None, + template_files_used: None, + api_proxy_file_path: String::new(), + custom_stream_response_path: None, + })), + custom_stream_response: Arc::new(ArcSwapOption::empty()), + access_token_secret: [0u8; 32], + encrypt_secret: [0u8; 16], + media_tools: Arc::new(MediaToolCapabilities::new()), + }) + } fn make_item(title: &str, group: &str, category_id: u32) -> PlaylistItem { PlaylistItem { header: PlaylistItemHeader { + id: title.intern(), title: title.intern(), group: group.intern(), category_id, xtream_cluster: XtreamCluster::Series, + item_type: PlaylistItemType::SeriesInfo, + url: format!("http://example.test/{title}").intern(), + input_name: "test-input".intern(), ..Default::default() }, } @@ -607,7 +1121,7 @@ mod tests { channels: vec![make_item("target-item", "B-Series", 22)], xtream_cluster: XtreamCluster::Series, }; - let mut source = MemoryPlaylistSource::new(vec![first_group, target_group]); + let mut source = MemoryPlaylistSource::new(vec![first_group, target_group]).into_source(); // Simulates a mapped delta group whose local pipeline id restarted at 1. let incoming = PlaylistGroup { @@ -625,4 +1139,149 @@ mod tests { assert_eq!(groups[1].title.as_ref(), "B-Series"); assert_eq!(groups[1].channels.len(), 2); } + + #[tokio::test] + async fn cloned_memory_source_is_copy_on_write() { + let group = PlaylistGroup { + id: 1, + title: "Series".intern(), + channels: vec![make_item("original", "Series", 1)], + xtream_cluster: XtreamCluster::Series, + }; + let source = MemoryPlaylistSource::new(vec![group]).into_source(); + let mut original = source.clone_source().expect("memory source clone should succeed"); + let mut cloned = source.clone_source().expect("memory source clone should succeed"); + + let incoming = PlaylistGroup { + id: 1, + title: "Series".intern(), + channels: vec![make_item("new", "Series", 1)], + xtream_cluster: XtreamCluster::Series, + }; + cloned.update_playlist(&incoming).await; + + assert_eq!(original.get_channel_count(), 1); + assert_eq!(cloned.get_channel_count(), 2); + } + + #[tokio::test] + async fn xtream_disk_clone_source_returns_error_when_query_clone_fails() { + let app_config = test_app_config(); + let storage_path = PathBuf::from("clone-error-fixture"); + let guard = Arc::new(app_config.file_locks.read_lock(&storage_path).await); + let source = PlaylistSource::xtream_disk(XtreamDiskPlaylistSource { + app_config, + storage_path, + live: Some((BPlusTreeQuery::::clone_error_fixture(), guard)), + vod: None, + series: None, + }); + + let result = source.clone_source(); + + assert!(matches!( + result, + Err(TuliproxError::RepositoryPlaylist(message)) + if message.contains("Failed to clone live disk playlist query") + && message.contains("No data source available to clone") + )); + } + + fn make_cluster_item(title: &str, cluster: XtreamCluster) -> PlaylistItem { + PlaylistItem { + header: PlaylistItemHeader { + id: title.intern(), + title: title.intern(), + group: format!("{cluster:?}").intern(), + xtream_cluster: cluster, + item_type: PlaylistItemType::from(cluster), + url: format!("http://example.test/{title}").intern(), + input_name: "test-input".intern(), + ..Default::default() + }, + } + } + + fn make_cluster_group(title: &str, cluster: XtreamCluster, channels: Vec) -> PlaylistGroup { + PlaylistGroup { + id: cluster as u32, + title: title.intern(), + channels, + xtream_cluster: cluster, + } + } + + fn filtered_source(skip_cluster: XtreamCluster) -> PlaylistSource { + let groups = vec![ + make_cluster_group("Live", XtreamCluster::Live, vec![make_cluster_item("live-1", XtreamCluster::Live)]), + make_cluster_group( + "Video", + XtreamCluster::Video, + vec![make_cluster_item("video-1", XtreamCluster::Video)], + ), + make_cluster_group( + "Series", + XtreamCluster::Series, + vec![make_cluster_item("series-1", XtreamCluster::Series)], + ), + ]; + PlaylistSource::filtered( + MemoryPlaylistSource::new(groups).into_source(), + HashSet::from([skip_cluster]), + ) + } + + #[test] + fn filtered_source_excludes_skipped_cluster_from_counts_and_items() { + let mut source = filtered_source(XtreamCluster::Video); + + assert_eq!(source.get_channel_count(), 2); + assert_eq!(source.get_group_count(), 2); + assert_eq!( + source + .items() + .map(|item| item.as_ref().header.xtream_cluster) + .collect::>(), + vec![XtreamCluster::Live, XtreamCluster::Series] + ); + } + + #[test] + fn filtered_source_excludes_skipped_cluster_when_taking_groups() { + let mut source = filtered_source(XtreamCluster::Video); + + let groups = source.take_groups(); + + assert_eq!(groups.len(), 2); + assert!(groups.iter().all(|group| group.xtream_cluster != XtreamCluster::Video)); + } + + #[test] + fn filtered_source_deduplicates_visible_memory_items_only() { + let duplicated_live = make_cluster_item("same-live", XtreamCluster::Live); + let groups = vec![ + make_cluster_group( + "Live", + XtreamCluster::Live, + vec![duplicated_live.clone(), duplicated_live], + ), + make_cluster_group( + "Video", + XtreamCluster::Video, + vec![make_cluster_item("video-1", XtreamCluster::Video)], + ), + ]; + let mut source = PlaylistSource::filtered( + MemoryPlaylistSource::new(groups).into_source(), + HashSet::from([XtreamCluster::Video]), + ); + let mut duplicates = HashSet::new(); + + source.deduplicate(&mut duplicates); + + let groups = source.take_groups(); + assert_eq!(groups.len(), 1); + assert_eq!(groups[0].channels.len(), 1); + assert_eq!(groups[0].xtream_cluster, XtreamCluster::Live); + } } diff --git a/backend/src/utils/file/file_utils.rs b/backend/src/utils/file/file_utils.rs index 81fab991a..132fd1e88 100644 --- a/backend/src/utils/file/file_utils.rs +++ b/backend/src/utils/file/file_utils.rs @@ -260,7 +260,8 @@ pub async fn persist_file(persist_file: Option, text: &str) { pub fn prepare_persist_path(file_name: &str, date_prefix: &str) -> PathBuf { let now = chrono::Local::now(); - let persist_filename = file_name.replace("{}", format!("{date_prefix}{}", now.format("%Y%m%d_%H%M%S").to_string().as_str()).as_str()); + let timestamp = format!("{date_prefix}{}", now.format("%Y%m%d_%H%M%S")); + let persist_filename = file_name.replace("{}", ×tamp); std::path::PathBuf::from(persist_filename) } diff --git a/backend/src/utils/network/request.rs b/backend/src/utils/network/request.rs index 635bd0b7e..edf5b3091 100644 --- a/backend/src/utils/network/request.rs +++ b/backend/src/utils/network/request.rs @@ -545,7 +545,7 @@ pub async fn send_with_retry_and_provider_policy( let request_builder = send(&attempt_target.request_url); let (base_client, request_result) = request_builder.build_split(); let mut request = request_result.map_err(|err| { - string_to_io_error(format!("Failed to build request: {}", sanitize_sensitive_info(err.to_string().as_str()))) + string_to_io_error(format!("Failed to build request: {}", sanitize_sensitive_info(&err.to_string()))) })?; apply_attempt_to_request(&mut request, &attempt_target)?; @@ -675,7 +675,7 @@ pub async fn send_with_retry_and_provider_policy( last_provider_failure = Some(format!( "connection error while trying {}: {}", sanitize_sensitive_info(attempt_target.request_url.as_str()), - sanitize_sensitive_info(err.to_string().as_str()) + sanitize_sensitive_info(&err.to_string()) )); // Connection errors (Timeout/Connect) trigger failover if provider exists @@ -715,7 +715,7 @@ pub async fn send_with_retry_and_provider_policy( } } - return Err(string_to_io_error(format!("Request error: {}", sanitize_sensitive_info(err.to_string().as_str())))); + return Err(string_to_io_error(format!("Request error: {}", sanitize_sensitive_info(&err.to_string())))); } } } @@ -750,7 +750,7 @@ pub async fn send_with_retry_and_provider_policy( break; } - Err(string_to_io_error("All attempts and providers exhausted".to_string())) + Err(string_to_io_error("All attempts and providers exhausted")) } fn is_failover_redirect(url: &Url, patterns: &[Arc]) -> bool { @@ -792,13 +792,13 @@ pub async fn get_input_epg_content_as_file( "can't download input {} epg url: {} => {}", input.name, sanitize_sensitive_info(url_str), - sanitize_sensitive_info(e.to_string().as_str()) + sanitize_sensitive_info(&e.to_string()) ); Err(TuliproxError::RepositoryNetwork(format!( "can't download input {} epg url: {} => {}", input.name, sanitize_sensitive_info(url_str), - sanitize_sensitive_info(e.to_string().as_str()) + sanitize_sensitive_info(&e.to_string()) ))) } } @@ -856,12 +856,12 @@ pub async fn get_input_text_content( error!( "Failed to download input '{}': {}", input.name, - sanitize_sensitive_info(e.to_string().as_str()) + sanitize_sensitive_info(&e.to_string()) ); Err(TuliproxError::RepositoryNetwork(format!( "Failed to download input '{}': {}", input.name, - sanitize_sensitive_info(e.to_string().as_str()) + sanitize_sensitive_info(&e.to_string()) ))) } } @@ -920,12 +920,12 @@ pub async fn get_input_text_content_as_stream( error!( "Failed to download input '{}': {}", input.name, - sanitize_sensitive_info(e.to_string().as_str()) + sanitize_sensitive_info(&e.to_string()) ); Err(TuliproxError::RepositoryNetwork(format!( "Failed to download input '{}': {}", input.name, - sanitize_sensitive_info(e.to_string().as_str()) + sanitize_sensitive_info(&e.to_string()) ))) } } @@ -1604,7 +1604,7 @@ pub async fn get_input_json_content( Err(e) => Err(TuliproxError::RepositoryNetwork(format!( "can't download input {input} => {sanitized}", input = input.name, - sanitized = sanitize_sensitive_info(e.to_string().as_str()) + sanitized = sanitize_sensitive_info(&e.to_string()) ))), } } @@ -1633,7 +1633,7 @@ pub async fn get_input_json_content_as_stream( Err(e) => Err(TuliproxError::RepositoryNetwork(format!( "can't download input {input} => {sanitized}", input = input.name, - sanitized = sanitize_sensitive_info(e.to_string().as_str()) + sanitized = sanitize_sensitive_info(&e.to_string()) ))), } } diff --git a/frontend/src/app/components/field_explanation.rs b/frontend/src/app/components/field_explanation.rs index d16768fcb..74c4d9884 100644 --- a/frontend/src/app/components/field_explanation.rs +++ b/frontend/src/app/components/field_explanation.rs @@ -3,7 +3,7 @@ use crate::{ i18n::{use_translation, YewI18n}, model::{DialogAction, DialogActions, DialogResult}, services::DialogService, - utils::t_safe, + utils::{join_non_empty_parts, t_safe}, }; use yew::{platform::spawn_local, prelude::*}; @@ -13,7 +13,7 @@ fn normalize_field_id(raw: &str) -> String { .map(|ch| if ch.is_ascii_alphanumeric() || ch == '.' { ch.to_ascii_uppercase() } else { '_' }) .collect::(); - normalized.split('_').filter(|part| !part.is_empty()).collect::>().join("_") + join_non_empty_parts(normalized.split('_'), "_") } fn field_tokens(field_id: &str) -> Vec<&str> { field_id.split('.').filter(|part| !part.is_empty()).collect::>() } diff --git a/frontend/src/app/components/field_id.rs b/frontend/src/app/components/field_id.rs index 08d76d470..31ac251c0 100644 --- a/frontend/src/app/components/field_id.rs +++ b/frontend/src/app/components/field_id.rs @@ -1,3 +1,5 @@ +use crate::utils::join_non_empty_parts; + pub fn resolve_field_id(field_id: &Option, name: &str, label: &str) -> String { let candidate = field_id .as_ref() @@ -28,13 +30,11 @@ pub fn resolve_field_id(field_id: &Option, name: &str, label: &str) -> S } fn normalize_upper_snake(raw: &str) -> String { - raw.chars() + let normalized = raw + .chars() .map(|ch| if ch.is_ascii_alphanumeric() { ch.to_ascii_uppercase() } else { '_' }) - .collect::() - .split('_') - .filter(|part| !part.is_empty()) - .collect::>() - .join("_") + .collect::(); + join_non_empty_parts(normalized.split('_'), "_") } fn to_upper_snake_case(raw: &str) -> String { @@ -63,7 +63,7 @@ fn to_upper_snake_case(raw: &str) -> String { } } - result.split('_').filter(|part| !part.is_empty()).collect::>().join("_") + join_non_empty_parts(result.split('_'), "_") } fn strip_ref_prefix(type_name: &str) -> &str { diff --git a/frontend/src/app/components/setup/setup_helpers.rs b/frontend/src/app/components/setup/setup_helpers.rs index 836623e47..b4581d9af 100644 --- a/frontend/src/app/components/setup/setup_helpers.rs +++ b/frontend/src/app/components/setup/setup_helpers.rs @@ -84,7 +84,14 @@ const SETUP_ERROR_PATTERNS: &[(&[&str], &str)] = &[ ]; fn normalize_setup_error_message(message: &str) -> String { - message.trim().to_lowercase().split_whitespace().collect::>().join(" ") + let mut result = String::new(); + for part in message.trim().to_lowercase().split_whitespace() { + if !result.is_empty() { + result.push(' '); + } + result.push_str(part); + } + result } fn matches_setup_error_pattern(message: &str, pattern_keywords: &[&str]) -> bool { diff --git a/frontend/src/app/components/source_editor/epg_smart_match_form.rs b/frontend/src/app/components/source_editor/epg_smart_match_form.rs index a633162e2..2f8419c30 100644 --- a/frontend/src/app/components/source_editor/epg_smart_match_form.rs +++ b/frontend/src/app/components/source_editor/epg_smart_match_form.rs @@ -46,10 +46,16 @@ impl From for EpgSmartMatchFormData { EpgNamePrefix::Suffix(v) => (NAME_PREFIX_MODE_SUFFIX.to_string(), Some(v)), EpgNamePrefix::Prefix(v) => (NAME_PREFIX_MODE_PREFIX.to_string(), Some(v)), }; - let separator = value - .name_prefix_separator - .as_ref() - .map(|chars| chars.iter().map(char::to_string).collect::>().join(",")); + let separator = value.name_prefix_separator.as_ref().map(|chars| { + let mut result = String::new(); + for ch in chars { + if !result.is_empty() { + result.push(','); + } + result.push(*ch); + } + result + }); Self { enabled: value.enabled, normalize_regex: value.normalize_regex, diff --git a/frontend/src/app/components/source_editor/provider_item_form.rs b/frontend/src/app/components/source_editor/provider_item_form.rs index 3c5a6aef3..5db878e36 100644 --- a/frontend/src/app/components/source_editor/provider_item_form.rs +++ b/frontend/src/app/components/source_editor/provider_item_form.rs @@ -455,7 +455,14 @@ pub fn ProviderItemForm(props: &ProviderItemFormProps) -> Html { { config_field_custom!( translate.t(LABEL_DNS_SCHEMES), dns_state.form.schemes.as_ref().map_or_else(String::new, |schemes| { - schemes.iter().map(|scheme| scheme_to_id(*scheme)).collect::>().join(", ") + let mut result = String::new(); + for scheme in schemes { + if !result.is_empty() { + result.push_str(", "); + } + result.push_str(scheme_to_id(*scheme)); + } + result }) ) } { config_field_bool!(dns_state.form, translate.t(LABEL_DNS_KEEP_VHOST), keep_vhost) } diff --git a/frontend/src/utils/mod.rs b/frontend/src/utils/mod.rs index 97d41d079..7c9eb814f 100644 --- a/frontend/src/utils/mod.rs +++ b/frontend/src/utils/mod.rs @@ -44,3 +44,14 @@ pub fn t_safe(i18n: &YewI18n, key: &str) -> Option { pub fn encoding_for_query(s: &str) -> String { js_sys::encode_uri_component(s).as_string().unwrap_or_else(|| s.to_string()) } + +pub fn join_non_empty_parts<'a>(parts: impl Iterator, separator: &str) -> String { + let mut result = String::new(); + for part in parts.filter(|part| !part.is_empty()) { + if !result.is_empty() { + result.push_str(separator); + } + result.push_str(part); + } + result +} diff --git a/shared/src/error/tuliprox_error.rs b/shared/src/error/tuliprox_error.rs index 84545a24a..35c4bed2c 100644 --- a/shared/src/error/tuliprox_error.rs +++ b/shared/src/error/tuliprox_error.rs @@ -281,4 +281,6 @@ where pub fn str_to_io_error(err: &str) -> std::io::Error { std::io::Error::other(sanitize_sensitive_info(err)) } -pub fn string_to_io_error(err: String) -> std::io::Error { std::io::Error::other(sanitize_sensitive_info(&err)) } +pub fn string_to_io_error(err: impl AsRef) -> std::io::Error { + std::io::Error::other(sanitize_sensitive_info(err.as_ref())) +} diff --git a/shared/src/foundation/filter.rs b/shared/src/foundation/filter.rs index bf4815059..55afb08ec 100644 --- a/shared/src/foundation/filter.rs +++ b/shared/src/foundation/filter.rs @@ -222,7 +222,7 @@ fn get_parser_item_field(expr: &Pair) -> Result if expr.as_rule() == Rule::field { let field_text = expr.as_str(); for item in ItemField::iter() { - if field_text.eq_ignore_ascii_case(item.to_string().as_str()) { + if field_text.eq_ignore_ascii_case(item.as_ref()) { return Ok(item); } } diff --git a/shared/src/foundation/value_provider.rs b/shared/src/foundation/value_provider.rs index 9cbf7af20..415bc2a21 100644 --- a/shared/src/foundation/value_provider.rs +++ b/shared/src/foundation/value_provider.rs @@ -122,7 +122,7 @@ pub fn get_field_value(pli: &PlaylistItem, field: ItemField) -> Arc { } } -pub fn set_field_value(pli: &mut PlaylistItem, field: ItemField, value: String) -> bool { +pub fn set_field_value(pli: &mut PlaylistItem, field: ItemField, value: &str) -> bool { let header = &mut pli.header; match field { ItemField::Group => header.group = value.intern(), diff --git a/shared/src/model/config/input.rs b/shared/src/model/config/input.rs index 8603ecdc8..167d47aa8 100644 --- a/shared/src/model/config/input.rs +++ b/shared/src/model/config/input.rs @@ -71,7 +71,7 @@ macro_rules! check_provider_scheme_url { return Err(TuliproxError::ConfigInput(format!( "Malformed provider URL {}: {}", sanitize_sensitive_info(&$url), - sanitize_sensitive_info(err.to_string().as_str()) + sanitize_sensitive_info(&err.to_string()) ))); } }; diff --git a/shared/src/model/stream_properties.rs b/shared/src/model/stream_properties.rs index 362f46bbd..785535c64 100644 --- a/shared/src/model/stream_properties.rs +++ b/shared/src/model/stream_properties.rs @@ -629,7 +629,7 @@ fn parse_season_field(s: &str) -> Option<(u32, String)> { let season: u32 = parts.next()?.parse().ok()?; - let kind = parts.collect::>().join("_"); + let kind = join_parts(parts, '_'); if kind.is_empty() { return None; } @@ -648,7 +648,7 @@ fn parse_season_episode_field(s: &str) -> Option<(u32, u32, String)> { let season: u32 = parts.next()?.parse().ok()?; let episode: u32 = parts.next()?.parse().ok()?; - let kind = parts.collect::>().join("_"); + let kind = join_parts(parts, '_'); if kind.is_empty() { return None; } @@ -656,6 +656,17 @@ fn parse_season_episode_field(s: &str) -> Option<(u32, u32, String)> { Some((season, episode, kind)) } +fn join_parts<'a>(parts: impl Iterator, separator: char) -> String { + let mut result = String::new(); + for part in parts { + if !result.is_empty() { + result.push(separator); + } + result.push_str(part); + } + result +} + impl VideoStreamProperties { fn from_info_base(info: &XtreamVideoInfo) -> VideoStreamProperties { VideoStreamProperties { diff --git a/shared/src/utils/mod.rs b/shared/src/utils/mod.rs index 62fe5903b..306dac1b0 100644 --- a/shared/src/utils/mod.rs +++ b/shared/src/utils/mod.rs @@ -48,6 +48,13 @@ macro_rules! write_if_some { } pub fn display_vec(vec: &[T]) -> String { - let inner = vec.iter().map(|item| format!("{item}")).collect::>().join(", "); - format!("[{inner}]") + let mut result = String::from("["); + for (idx, item) in vec.iter().enumerate() { + if idx > 0 { + result.push_str(", "); + } + result.push_str(&item.to_string()); + } + result.push(']'); + result } diff --git a/shared/src/utils/size_utils.rs b/shared/src/utils/size_utils.rs index c287cf008..7d3fa307f 100644 --- a/shared/src/utils/size_utils.rs +++ b/shared/src/utils/size_utils.rs @@ -68,12 +68,19 @@ pub fn parse_to_kbps(input: &str) -> Result { } } - u64::from_str(speed_str).map_err(|_| { - format!( - "Invalid speed: {speed_str}, supported units are {}", - units.iter().map(|p| p.0).collect::>().join(",") - ) - }) + u64::from_str(speed_str) + .map_err(|_| format!("Invalid speed: {speed_str}, supported units are {}", join_unit_names(units))) +} + +fn join_unit_names(units: &[(&str, u64)]) -> String { + let mut result = String::new(); + for (idx, (unit, _)) in units.iter().enumerate() { + if idx > 0 { + result.push(','); + } + result.push_str(unit); + } + result } pub fn human_readable_kbps(kbps: u64) -> String {