diff --git a/CHANGELOG.md b/CHANGELOG.md index 8dbde9a28..0873a92a5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,4 +1,7 @@ # Changelog +# 2.1.4 (2025-01-xx) +- !BREAKING CHANGE! unique `input` `name` is now mandatory, because rearranging the `source.yml` could lead to wrong results without a playlist update. + # 2.1.3 (2025-01-26) - Hotfix 2.1.2, forgot to update the stream api code. diff --git a/README.md b/README.md index 4661f321b..491521aa5 100644 --- a/README.md +++ b/README.md @@ -309,7 +309,7 @@ This will replace all occurrences of `!delimiter!` and `!quality!` in the regexp Each input has the following attributes: -- `name` is optional, if set it must be unique, should be set for the webui +- `name` is mandatory, it must be unique. - `type` is optional, default is `m3u`. Valid values are `m3u` and `xtream` - `enabled` is optional, default is true, if you disable the processing is skipped - `persist` is optional, you can skip or leave it blank to avoid persisting the input file. The `{}` in the filename is filled with the current timestamp. diff --git a/bin/build_github_docker.sh b/bin/build_github_docker.sh index 724a8e5f5..d321d1715 100755 --- a/bin/build_github_docker.sh +++ b/bin/build_github_docker.sh @@ -29,7 +29,7 @@ fi BIN_FILE=${WORKING_DIR}/target/${TARGET}/release/m3u-filter cp "${WORKING_DIR}/target/${TARGET}/release/m3u-filter" "${DOCKER_DIR}/" rm -rf "${DOCKER_DIR}/web" -cp -r "${WORKING_DIR}/frontend/build" "${DOCKER_DIR}/web" +cp -r "${FRONTEND_DIR}/build" "${DOCKER_DIR}/web" # Get the version from the binary VERSION=$("$BIN_FILE" -V | sed 's/m3u-filter *//') diff --git a/src/api/api_utils.rs b/src/api/api_utils.rs index 4259ef8bb..4deecfc5a 100644 --- a/src/api/api_utils.rs +++ b/src/api/api_utils.rs @@ -1,34 +1,34 @@ use crate::api::model::app_state::AppState; +use crate::api::model::model_utils::get_stream_response_with_headers; +use crate::api::model::persist_pipe_stream::PersistPipeStream; use crate::api::model::provider_stream; -use crate::api::model::provider_stream::{get_provider_pipe_stream}; +use crate::api::model::provider_stream::get_provider_pipe_stream; +use crate::api::model::provider_stream_factory::BufferStreamOptions; use crate::api::model::request::UserApiRequest; use crate::api::model::shared_stream::SharedStream; +use crate::api::model::stream_error::StreamError; use crate::debug_if_enabled; use crate::model::api_proxy::ProxyUserCredentials; use crate::model::config::{ConfigInput, ConfigTarget}; use crate::model::playlist::PlaylistItemType; +use crate::utils::file_utils::create_new_file_for_write; +use crate::utils::lru_cache::LRUResourceCache; use crate::utils::request_utils; use crate::utils::request_utils::sanitize_sensitive_info; use actix_files::NamedFile; use actix_web::body::{BodyStream, SizedStream}; use actix_web::http::header::{HeaderValue, CACHE_CONTROL}; use actix_web::{HttpRequest, HttpResponse}; -use bytes::Bytes; -use log::{error, log_enabled, trace}; -use std::collections::HashMap; -use std::path::{Path}; -use std::sync::Arc; use async_std::sync::Mutex; +use bytes::Bytes; use futures::stream::BoxStream; -use futures::{TryStreamExt}; +use futures::TryStreamExt; +use log::{error, log_enabled, trace}; use reqwest::StatusCode; +use std::collections::HashMap; +use std::path::Path; +use std::sync::Arc; use url::Url; -use crate::api::model::model_utils::get_stream_response_with_headers; -use crate::api::model::persist_pipe_stream::PersistPipeStream; -use crate::api::model::provider_stream_factory::BufferStreamOptions; -use crate::api::model::stream_error::StreamError; -use crate::utils::file_utils::create_new_file_for_write; -use crate::utils::lru_cache::LRUResourceCache; pub async fn serve_file(file_path: &Path, req: &HttpRequest, mime_type: mime::Mime) -> HttpResponse { if file_path.exists() { @@ -78,18 +78,7 @@ async fn create_broadcast_stream( } } -pub async fn stream_response(app_state: &AppState, stream_url: &str, - req: &HttpRequest, input: Option<&ConfigInput>, - item_type: PlaylistItemType, target: &ConfigTarget) -> HttpResponse { - if log_enabled!(log::Level::Trace) { trace!("Try to open stream {}", sanitize_sensitive_info(stream_url)); } - - let share_stream = is_stream_share_enabled(item_type, target); - if share_stream { - if let Some(value) = shared_stream_response(app_state, stream_url).await { - return value; - } - } - +fn get_stream_options(app_state: &AppState) -> (bool, bool, usize, bool, bool) { let (stream_retry, buffer_enabled, buffer_size) = app_state .config .reverse_proxy @@ -102,10 +91,36 @@ pub async fn stream_response(app_state: &AppState, stream_url: &str, .map_or((false, 0), |buffer| (buffer.enabled, buffer.size)); (stream.retry, buffer_enabled, buffer_size) }); + let pipe_provider_stream = !stream_retry && !buffer_enabled; + let shared_stream_use_own_buffer = !buffer_enabled || pipe_provider_stream; + (stream_retry, buffer_enabled, buffer_size, pipe_provider_stream, shared_stream_use_own_buffer) +} +fn get_stream_content_length(provider_response: Option<&(Vec<(String, String)>, StatusCode)>) -> u64 { + let content_length = provider_response + .as_ref() + .and_then(|(headers, _)| headers.iter().find(|(h, _)| h.eq(actix_web::http::header::CONTENT_LENGTH.as_str()))) + .and_then(|(_, val)| val.parse::().ok()) + .unwrap_or(0); + content_length +} + +pub async fn stream_response(app_state: &AppState, stream_url: &str, + req: &HttpRequest, input: Option<&ConfigInput>, + item_type: PlaylistItemType, target: &ConfigTarget) -> HttpResponse { + if log_enabled!(log::Level::Trace) { trace!("Try to open stream {}", sanitize_sensitive_info(stream_url)); } + + let share_stream = is_stream_share_enabled(item_type, target); + if share_stream { + if let Some(value) = shared_stream_response(app_state, stream_url).await { + return value; + } + } + + let (stream_retry, buffer_enabled, buffer_size, direct_pipe_provider_stream, shared_stream_use_own_buffer) = + get_stream_options(app_state); if let Ok(url) = Url::parse(stream_url) { - let direct_pipe_provider_stream = !stream_retry && !buffer_enabled; let (stream_opt, provider_response) = if direct_pipe_provider_stream { get_provider_pipe_stream(&app_state.http_client, &url, req, input).await } else { @@ -113,14 +128,8 @@ pub async fn stream_response(app_state: &AppState, stream_url: &str, provider_stream::get_provider_reconnect_buffered_stream(&app_state.http_client, &url, req, input, buffer_stream_options).await }; if let Some(stream) = stream_opt { - let content_length = provider_response - .as_ref() - .and_then(|(headers, _)| headers.iter().find(|(h, _)| h.eq(actix_web::http::header::CONTENT_LENGTH.as_str()))) - .and_then(|(_, val)| val.parse::().ok()) - .unwrap_or(0); + let content_length = get_stream_content_length(provider_response.as_ref()); - - let shared_stream_use_own_buffer = !buffer_enabled || direct_pipe_provider_stream; let stream_resp = if share_stream { let shared_headers = provider_response.as_ref().map_or_else(Vec::new, |(h, _)| h.clone()); SharedStream::register(app_state, stream_url, stream, shared_stream_use_own_buffer, shared_headers).await; @@ -145,7 +154,7 @@ pub async fn stream_response(app_state: &AppState, stream_url: &str, async fn shared_stream_response(app_state: &AppState, stream_url: &str) -> Option { if let Some(stream) = create_broadcast_stream(app_state, stream_url).await { debug_if_enabled!("Using shared channel {}", sanitize_sensitive_info(stream_url)); - if let Some((headers,_)) = app_state.shared_streams.lock().await.get(stream_url) { + if let Some((headers, _)) = app_state.shared_streams.lock().await.get(stream_url) { let mut response_builder = get_stream_response_with_headers(Some((headers.clone(), StatusCode::OK)), stream_url); return Some(response_builder.body(BodyStream::new(stream))); } @@ -172,7 +181,7 @@ pub fn get_headers_from_request(req: &HttpRequest, filter: &HeaderFilter) -> Has fn get_add_cache_content(res_url: &str, cache: &Arc>>) -> Box { let resource_url = String::from(res_url); let cache = Arc::clone(cache); - let add_cache_content: Box = Box::new(move|size| { + let add_cache_content: Box = Box::new(move |size| { let res_url = resource_url.clone(); let cache = Arc::clone(&cache); actix_rt::spawn(async move { @@ -212,9 +221,9 @@ pub async fn resource_response(app_state: &AppState, resource_url: &str, req: &H response_builder.insert_header((k.as_str(), v.as_ref())); }); - let byte_stream = response.bytes_stream().map_err(|err|StreamError::reqwest(&err)); + let byte_stream = response.bytes_stream().map_err(|err| StreamError::reqwest(&err)); if let Some(cache) = app_state.cache.as_ref() { - let resource_path = { + let resource_path = { let guard = cache.lock().await; guard.store_path(resource_url) }; @@ -225,7 +234,7 @@ pub async fn resource_response(app_state: &AppState, resource_url: &str, req: &H return response_builder.body(BodyStream::new(stream)); } } - return response_builder.body(BodyStream::new(byte_stream)); + return response_builder.body(BodyStream::new(byte_stream)); } debug_if_enabled!("Failed to open resource got status {} for {}", status, sanitize_sensitive_info(resource_url)); } diff --git a/src/api/model/config.rs b/src/api/model/config.rs index 949eb99ff..da42c0dde 100644 --- a/src/api/model/config.rs +++ b/src/api/model/config.rs @@ -12,7 +12,7 @@ pub struct ServerInputConfig { pub username: Option, pub password: Option, pub persist: Option, - pub name: Option, + pub name: String, pub enabled: bool, } diff --git a/src/api/model/request.rs b/src/api/model/request.rs index 1b9e1e8d8..c001a8481 100644 --- a/src/api/model/request.rs +++ b/src/api/model/request.rs @@ -4,7 +4,7 @@ use serde::{Deserialize, Serialize}; #[derive(Deserialize, Serialize, Debug, Clone)] pub struct PlaylistRequest { pub url: Option, - pub input_id: Option, + pub input_name: Option, } impl From> for PlaylistRequest { diff --git a/src/api/v1_api.rs b/src/api/v1_api.rs index 8f0d71e33..3a98c2fa9 100644 --- a/src/api/v1_api.rs +++ b/src/api/v1_api.rs @@ -113,9 +113,10 @@ async fn playlist_update( } } -fn create_config_input_for_url(url: &str) -> ConfigInput { +fn create_config_input_for_url(name: &str, url: &str) -> ConfigInput { ConfigInput { id: 0, + name: String::from(name), input_type: InputType::M3u, url: String::from(url), enabled: true, @@ -151,11 +152,12 @@ async fn playlist( req: web::Json, app_state: web::Data, ) -> HttpResponse { - if let Some(input_id) = req.input_id { - get_playlist(Arc::clone(&app_state.http_client), app_state.config.get_input_by_id(input_id), &app_state.config).await + if let Some(input_name) = req.input_name.as_ref() { + get_playlist(Arc::clone(&app_state.http_client), app_state.config.get_input_by_name(input_name), &app_state.config).await } else { let url = req.url.as_deref().unwrap_or(""); - let input = create_config_input_for_url(url); + let name = req.input_name.as_deref().unwrap_or(""); + let input = create_config_input_for_url(name, url); get_playlist(Arc::clone(&app_state.http_client), Some(&input), &app_state.config).await } } @@ -165,12 +167,12 @@ async fn config( ) -> HttpResponse { let map_input = |i: &ConfigInput| ServerInputConfig { id: i.id, + name: i.name.clone(), input_type: i.input_type.clone(), url: i.url.clone(), username: i.username.clone(), password: i.password.clone(), persist: i.persist.clone(), - name: i.name.clone(), enabled: i.enabled, }; diff --git a/src/api/xtream_api.rs b/src/api/xtream_api.rs index 721883d8c..e6a501785 100644 --- a/src/api/xtream_api.rs +++ b/src/api/xtream_api.rs @@ -189,7 +189,7 @@ async fn xtream_player_api_stream( let (action_stream_id, stream_ext) = xtream_api_request_separate_number_and_remainder(stream_req.stream_id); let virtual_id: u32 = try_result_bad_request!(action_stream_id.trim().parse()); let (pli, mapping) = try_result_bad_request!(xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, None).await, true, format!("Failed to read xtream item for stream id {}", virtual_id)); - let input = try_option_bad_request!(app_state.config.get_input_by_id(pli.input_id), true, format!("Cant find input for target {target_name}, context {}, stream_id {virtual_id}", stream_req.context)); + let input = try_option_bad_request!(app_state.config.get_input_by_name(pli.input_name.as_str()), true, format!("Cant find input for target {target_name}, context {}, stream_id {virtual_id}", stream_req.context)); if pli.item_type == PlaylistItemType::LiveHls { debug_if_enabled!("Redirecting stream request to {}", sanitize_sensitive_info(&pli.url)); @@ -461,8 +461,8 @@ async fn xtream_get_stream_info_response(app_state: &AppState, user: &ProxyUserC }; if let Ok((pli, _)) = xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, Some(cluster)).await { - let input_id = pli.input_id; - if let Some(input) = app_state.config.get_input_by_id(input_id) { + let input_name = Rc::clone(&pli.input_name); + if let Some(input) = app_state.config.get_input_by_name(input_name.as_str()) { if let Some(info_url) = download::get_xtream_player_api_info_url(input, cluster, pli.provider_id) { // Redirect is only possible for live streams, vod and series info needs to be modified if user.proxy == ProxyType::Redirect && cluster == XtreamCluster::Live { @@ -489,8 +489,8 @@ async fn xtream_get_short_epg(app_state: &AppState, user: &ProxyUserCredentials, }; if let Ok((pli, _)) = xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, None).await { - let input_id: u16 = pli.input_id; - if let Some(input) = app_state.config.get_input_by_id(input_id) { + let input_name = Rc::clone(&pli.input_name); + if let Some(input) = app_state.config.get_input_by_name(input_name.as_str()) { if let Some(action_url) = download::get_xtream_player_api_action_url(input, ACTION_GET_SHORT_EPG) { let mut info_url = format!("{action_url}&{TAG_STREAM_ID}={}", pli.provider_id); if !(limit.is_empty() || limit.eq("0")) { @@ -539,7 +539,7 @@ async fn xtream_player_api_handle_content_action(config: &Config, target_name: & async fn xtream_get_catchup_response(app_state: &AppState, target: &ConfigTarget, stream_id: &str, start: &str, end: &str) -> HttpResponse { let virtual_id: u32 = try_result_bad_request!(FromStr::from_str(stream_id)); let (pli, _) = try_result_bad_request!(xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, Some(XtreamCluster::Live)).await); - let input = try_option_bad_request!(app_state.config.get_input_by_id(pli.input_id)); + let input = try_option_bad_request!(app_state.config.get_input_by_name(pli.input_name.as_str())); let info_url = try_option_bad_request!(download::get_xtream_player_api_action_url(input, ACTION_GET_CATCHUP_TABLE).map(|action_url| format!("{action_url}&{TAG_STREAM_ID}={}&start={start}&end={end}", pli.provider_id))); let content = try_result_bad_request!(download::get_xtream_stream_info_content(Arc::clone(&app_state.http_client), info_url.as_str(), input).await); let mut doc: Map = try_result_bad_request!(serde_json::from_str(&content)); diff --git a/src/model/config.rs b/src/model/config.rs index 418619885..044fd7842 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -22,10 +22,10 @@ use crate::model::mapping::Mapping; use crate::model::mapping::Mappings; use crate::utils::default_utils::{default_as_default, default_as_true, default_as_two_u16}; use crate::utils::file_lock_manager::FileLockManager; -use crate::utils::{config_reader, file_utils}; -use crate::{exit, info_err}; use crate::utils::file_utils::file_reader; use crate::utils::size_utils::parse_size_base_2; +use crate::utils::{config_reader, file_utils}; +use crate::{exit, info_err}; pub const MAPPER_ATTRIBUTE_FIELDS: &[&str] = &[ "name", "title", "group", "id", "chno", "logo", @@ -568,6 +568,7 @@ pub struct InputUserInfo { pub struct ConfigInput { #[serde(skip)] pub id: u16, + pub name: String, #[serde(default, rename = "type")] pub input_type: InputType, #[serde(default)] @@ -585,18 +586,18 @@ pub struct ConfigInput { pub prefix: Option, #[serde(default, skip_serializing_if = "Option::is_none")] pub suffix: Option, - #[serde(default, skip_serializing_if = "Option::is_none")] - pub name: Option, #[serde(default = "default_as_true")] pub enabled: bool, #[serde(default, skip_serializing_if = "Option::is_none")] pub options: Option, - } impl ConfigInput { pub fn prepare(&mut self, id: u16) -> Result<(), M3uFilterError> { self.id = id; + if self.name.trim().is_empty() { + return Err(info_err!("name for input is mandatory".to_string())); + } if self.url.trim().is_empty() { return Err(info_err!("url for input is mandatory".to_string())); } @@ -914,7 +915,7 @@ pub struct CacheConfig { } impl CacheConfig { - fn prepare(&mut self, working_dir: &str, resolve_var: bool) { + fn prepare(&mut self, working_dir: &str, resolve_var: bool) { if self.enabled { let work_path = PathBuf::from(working_dir); if self.dir.is_none() { @@ -1086,10 +1087,10 @@ impl Config { self.t_api_proxy.read().unwrap().as_ref().and_then(|api_proxy| api_proxy.get_user_credentials(username)) } - pub fn get_input_by_id(&self, input_id: u16) -> Option<&ConfigInput> { + pub fn get_input_by_name(&self, input_name: &str) -> Option<&ConfigInput> { for source in &self.sources { for input in &source.inputs { - if input.id == input_id { + if input.name == input_name { return Some(input); } } @@ -1114,6 +1115,60 @@ impl Config { } } + fn check_unique_input_names(&mut self) -> Result<(), M3uFilterError> { + let mut seen_names = HashSet::new(); + for source in &mut self.sources { + for input in &source.inputs { + let input_name = input.name.trim().to_string(); + if input_name.is_empty() { + return create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "input name required"); + } + if seen_names.contains(input_name.as_str()) { + return create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "input names should be unique: {}", input_name); + } + seen_names.insert(input_name); + } + } + Ok(()) + } + + fn check_unique_target_names(&mut self) -> Result, M3uFilterError> { + let mut seen_names = HashSet::new(); + let default_target_name = default_as_default(); + for source in &self.sources { + for target in &source.targets { + // check target name is unique + let target_name = target.name.trim().to_string(); + if target_name.is_empty() { + return create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "target name required"); + } + if !default_target_name.eq_ignore_ascii_case(target_name.as_str()) { + if seen_names.contains(target_name.as_str()) { + return create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "target names should be unique: {}", target_name); + } + seen_names.insert(target_name); + } + } + } + Ok(seen_names) + } + + + fn check_scheduled_targets(&mut self, target_names: &HashSet) -> Result<(), M3uFilterError> { + if let Some(schedules) = &self.schedules { + for schedule in schedules { + if let Some(targets) = &schedule.targets { + for target_name in targets { + if !target_names.contains(target_name) { + return create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "Unknown target name in scheduler: {}", target_name); + } + } + } + } + } + Ok(()) + } + pub fn prepare(&mut self, resolve_var: bool) -> Result<(), M3uFilterError> { let work_dir = if resolve_var { &config_reader::resolve_env_var(&self.working_dir) } else { &self.working_dir }; self.working_dir = file_utils::get_working_path(work_dir); @@ -1139,24 +1194,14 @@ impl Config { } }; // prepare sources and set id's - let mut target_names_check = HashSet::::new(); - let default_target_name = default_as_default(); let mut source_index: u16 = 1; let mut target_index: u16 = 1; + self.check_unique_input_names()?; + let target_names = self.check_unique_target_names()?; + self.check_scheduled_targets(&target_names)?; for source in &mut self.sources { source_index = source.prepare(source_index)?; for target in &mut source.targets { - // check target name is unique - let target_name = target.name.trim().to_string(); - if target_name.is_empty() { - return create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "target name required"); - } - if !default_target_name.eq_ignore_ascii_case(target_name.as_str()) { - if target_names_check.contains(target_name.as_str()) { - return create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "target names should be unique: {}", target_name); - } - target_names_check.insert(target_name); - } // prepare templates let prepare_result = match &self.templates { Some(templ) => target.prepare(target_index, Some(templ)), @@ -1165,18 +1210,6 @@ impl Config { prepare_result?; target_index += 1; } - - if let Some(schedules) = &self.schedules { - for schedule in schedules { - if let Some(targets) = &schedule.targets { - for target_name in targets { - if !target_names_check.contains(target_name) { - return create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "Unknown target name in scheduler: {}", target_name); - } - } - } - } - } } match &mut self.video { @@ -1236,7 +1269,7 @@ impl Config { pub fn get_user_server_info(&self, user: &ProxyUserCredentials) -> ApiProxyServerInfo { let server_info_list = self.t_api_proxy.read().unwrap().as_ref().unwrap().server.clone(); let server_info_name = user.server.as_ref().map_or("default", |server_name| server_name.as_str()); - server_info_list.iter().find(|c| c.name.eq(server_info_name)).map_or_else(|| server_info_list.first().unwrap().clone(), std::clone::Clone::clone) + server_info_list.iter().find(|c| c.name.eq(server_info_name)).map_or_else(|| server_info_list.first().unwrap().clone(), Clone::clone) } } diff --git a/src/model/playlist.rs b/src/model/playlist.rs index e7f1321f8..99984e344 100644 --- a/src/model/playlist.rs +++ b/src/model/playlist.rs @@ -159,8 +159,7 @@ pub struct PlaylistItemHeader { pub item_type: PlaylistItemType, #[serde(default)] pub category_id: u32, - #[serde(default)] - pub input_id: u16, + pub input_name: Rc, } impl PlaylistItemHeader { @@ -288,7 +287,7 @@ pub struct M3uPlaylistItem { pub rec: Rc, pub url: Rc, pub epg_channel_id: Option>, - pub input_id: u16, + pub input_name: Rc, pub item_type: PlaylistItemType, } @@ -390,7 +389,7 @@ pub struct XtreamPlaylistItem { pub additional_properties: Option, pub item_type: PlaylistItemType, pub category_id: u32, - pub input_id: u16, + pub input_name: Rc, } impl XtreamPlaylistItem { @@ -501,7 +500,7 @@ impl PlaylistItem { rec: Rc::clone(&header.rec), url: Rc::clone(&header.url), epg_channel_id: header.epg_channel_id.clone(), - input_id: header.input_id, + input_name: Rc::clone(&header.input_name), item_type: header.item_type, } } @@ -525,7 +524,7 @@ impl PlaylistItem { additional_properties: header.additional_properties.as_ref().and_then(|props| serde_json::to_string(props).ok()), item_type: header.item_type, category_id: header.category_id, - input_id: header.input_id, + input_name: Rc::clone(&header.input_name), } } } diff --git a/src/processing/m3u_parser.rs b/src/processing/m3u_parser.rs index 6eb7bdd0e..00c98ab2b 100644 --- a/src/processing/m3u_parser.rs +++ b/src/processing/m3u_parser.rs @@ -73,11 +73,11 @@ fn skip_digit(it: &mut std::str::Chars) -> Option { } } -fn create_empty_playlistitem_header(input_id: u16, url: &str) -> PlaylistItemHeader { +fn create_empty_playlistitem_header(input_name: &str, url: &str) -> PlaylistItemHeader { PlaylistItemHeader { url: Rc::new(url.to_owned()), category_id: 0, - input_id, + input_name: Rc::new(input_name.to_string()), ..Default::default() } } @@ -94,7 +94,7 @@ macro_rules! process_header_fields { } fn process_header(input: &ConfigInput, video_suffixes: &[&str], content: &str, url: &str) -> PlaylistItemHeader { - let mut plih = create_empty_playlistitem_header(input.id, url); + let mut plih = create_empty_playlistitem_header(input.name.as_str(), url); let mut it = content.chars(); let line_token = token_till(&mut it, ':', false); if line_token.as_deref() == Some("#EXTINF") { diff --git a/src/processing/playlist_processor.rs b/src/processing/playlist_processor.rs index 170ffe077..0d2e3738b 100644 --- a/src/processing/playlist_processor.rs +++ b/src/processing/playlist_processor.rs @@ -31,7 +31,6 @@ use crate::processing::xtream_processor_vod::playlist_resolve_vod; use crate::repository::playlist_repository::persist_playlist; use crate::utils::default_utils::default_as_default; use crate::utils::download; -use crate::utils::request_utils::sanitize_sensitive_info; use crate::{debug_if_enabled, get_errors_notify_message, model::config, notify_err, Config}; fn is_valid(pli: &PlaylistItem, target: &ConfigTarget) -> bool { @@ -294,14 +293,13 @@ fn is_target_enabled(target: &ConfigTarget, user_targets: &ProcessTargets) -> bo async fn process_source(client: Arc, cfg: Arc, source_idx: usize, user_targets: Arc) -> (Vec, Vec, Vec) { let source = cfg.sources.get(source_idx).unwrap(); let mut errors = vec![]; - let mut input_stats = HashMap::::new(); + let mut input_stats = HashMap::::new(); let mut target_stats = Vec::::new(); let mut source_playlists = Vec::with_capacity(128); let enabled_inputs = source.inputs.iter().filter(|item| item.enabled).count(); // Downlod the sources for input in &source.inputs { - let input_id = input.id; - if is_input_enabled(enabled_inputs, input.enabled, input_id, &user_targets) { + if is_input_enabled(enabled_inputs, input.enabled, input.id, &user_targets) { let start_time = Instant::now(); let (mut playlistgroups, mut error_list) = match input.input_type { InputType::M3u => download::get_m3u_playlist(Arc::clone(&client), &cfg, input, &cfg.working_dir).await, @@ -318,7 +316,7 @@ async fn process_source(client: Arc, cfg: Arc, source_i let channel_count = playlistgroups.iter() .map(|group| group.channels.len()) .sum(); - let input_name = input.name.as_ref().map_or_else(|| sanitize_sensitive_info(input.url.as_str()), std::string::ToString::to_string); + let input_name = &input.name; if playlistgroups.is_empty() { info!("Source is empty {input_name}"); errors.push(notify_err!(format!("Source is empty {input_name}"))); @@ -333,8 +331,8 @@ async fn process_source(client: Arc, cfg: Arc, source_i ); } let elapsed = start_time.elapsed().as_secs(); - input_stats.insert(input_id, create_input_stat(group_count, channel_count, error_list.len(), - input.input_type.clone(), &input_name, elapsed)); + input_stats.insert(input_name.to_string(), create_input_stat(group_count, channel_count, error_list.len(), + input.input_type.clone(), input_name, elapsed)); } } if source_playlists.is_empty() { @@ -490,7 +488,7 @@ async fn process_playlist_for_target(client: Arc, playlists: &mut [FetchedPlaylist<'_>], target: &ConfigTarget, cfg: &Config, - stats: &mut HashMap, + stats: &mut HashMap, errors: &mut Vec) -> Result<(), Vec> { let pipe = get_processing_pipe(target); debug_if_enabled!("Processing order is {}", &target.processing_order); @@ -502,7 +500,7 @@ async fn process_playlist_for_target(client: Arc, playlist_resolve_series(Arc::clone(&client), cfg, target, errors, &pipe, provider_fpl, &mut processed_fpl).await; playlist_resolve_vod(Arc::clone(&client), cfg, target, errors, &processed_fpl).await; // stats - let input_stats = stats.get_mut(&processed_fpl.input.id); + let input_stats = stats.get_mut(&processed_fpl.input.name); if let Some(stat) = input_stats { stat.processed_stats.group_count = processed_fpl.playlistgroups.len(); stat.processed_stats.channel_count = processed_fpl.playlistgroups.iter() diff --git a/src/processing/xtream_parser.rs b/src/processing/xtream_parser.rs index d14d680dd..0de2df820 100644 --- a/src/processing/xtream_parser.rs +++ b/src/processing/xtream_parser.rs @@ -63,7 +63,7 @@ pub fn parse_xtream_series_info(info: &Value, group_title: &str, series_name: &s xtream_cluster: XtreamCluster::Series, additional_properties: episode.get_additional_properties(&series_info), category_id: 0, - input_id: input.id, + input_name: Rc::new(input.name.to_string()), ..Default::default() }) }) @@ -101,7 +101,7 @@ pub fn parse_xtream(input: &ConfigInput, streams: &Value) -> Result>, M3uFilterError> { match map_to_xtream_category(categories) { Ok(xtream_categories) => { - let input_id = input.id; + let input_name = Rc::new(input.name.to_string()); let url = input.url.as_str(); let username = input.username.as_ref().map_or("", |v| v); let password = input.password.as_ref().map_or("", |v| v); @@ -136,7 +136,7 @@ pub fn parse_xtream(input: &ConfigInput, xtream_cluster, additional_properties: stream.get_additional_properties(), category_id: 0, - input_id, + input_name: Rc::clone(&input_name), ..Default::default() }), }; diff --git a/src/processing/xtream_processor.rs b/src/processing/xtream_processor.rs index 81a958c3f..0d6b4220e 100644 --- a/src/processing/xtream_processor.rs +++ b/src/processing/xtream_processor.rs @@ -135,7 +135,7 @@ where { Ok(file_path) => file_path, Err(err) => { - let fpl_name = fpl.input.name.as_ref().map_or("?", String::as_str); + let fpl_name = &fpl.input.name; errors.push(notify_err!(format!("Could not create storage path for input {fpl_name}: {err}"))); return processed_info_ids; } diff --git a/src/repository/kodi_repository.rs b/src/repository/kodi_repository.rs index 5f13e870d..566c18570 100644 --- a/src/repository/kodi_repository.rs +++ b/src/repository/kodi_repository.rs @@ -132,7 +132,7 @@ async fn kodi_style_rename(cfg: &Config, strm_item_info: &StrmItemInfo, style: & .filter(|&series_name| name_4.starts_with(series_name)) .and_then(|series_name| trim_string_after_pos(&name_3, series_name.len())); let tmdb_id = if let Some(value) = match strm_item_info.item_type { - PlaylistItemType::Series | PlaylistItemType::Video => get_tmdb_value(cfg, strm_item_info.provider_id, strm_item_info.input_id, input_tmdb_indexes, strm_item_info.item_type).await, + PlaylistItemType::Series | PlaylistItemType::Video => get_tmdb_value(cfg, strm_item_info.provider_id, strm_item_info.input_name.as_str(), input_tmdb_indexes, strm_item_info.item_type).await, _ => None, } { match value { @@ -207,15 +207,15 @@ enum InputTmdbIndexValue { Series(XtreamSeriesEpisode), } -type InputTmdbIndexMap = HashMap>; -async fn get_tmdb_value(cfg: &Config, provider_id: Option, input_id: u16, +type InputTmdbIndexMap = HashMap>; +async fn get_tmdb_value(cfg: &Config, provider_id: Option, input_name: &str, input_indexes: &mut InputTmdbIndexMap, item_type: PlaylistItemType) -> Option { // the tmdb_ids are stored inside record files for xtream input. // we load this record files on request for each input and item_type. match provider_id { None => None, Some(pid) => { - match input_indexes.entry(input_id) { + match input_indexes.entry(input_name.to_string()) { std::collections::hash_map::Entry::Occupied(entry) => { if let Some((_, tree_value)) = entry.get() { match tree_value { @@ -227,7 +227,7 @@ async fn get_tmdb_value(cfg: &Config, provider_id: Option, input_id: u16, } } std::collections::hash_map::Entry::Vacant(entry) => { - if let Some(input) = cfg.get_input_by_id(input_id) { + if let Some(input) = cfg.get_input_by_name(input_name) { if let Ok(Some(tmdb_path)) = get_input_storage_path(input, &cfg.working_dir) .map(|storage_path| xtream_get_record_file_path(&storage_path, item_type)) { if let Ok(file_lock) = cfg.file_locks.read_lock(&tmdb_path).await { @@ -266,7 +266,7 @@ struct StrmItemInfo { item_type: PlaylistItemType, provider_id: Option, virtual_id: u32, - input_id: u16, + input_name: Rc, url: Rc, series_name: Option, release_date: Option, @@ -281,7 +281,7 @@ fn extract_item_info(pli: &PlaylistItem) -> StrmItemInfo { let item_type = header.item_type; let provider_id = header.get_provider_id(); let virtual_id = header.virtual_id; - let input_id = header.input_id; + let input_name = Rc::clone(&header.input_name); let url = Rc::clone(&header.url); let (series_name, release_date, season, episode) = if header.item_type == PlaylistItemType::Series { let series_name = match header.get_field("name") { @@ -295,7 +295,7 @@ fn extract_item_info(pli: &PlaylistItem) -> StrmItemInfo { (series_name, release_date, season, episode) } else { (None, None, None, None) }; - StrmItemInfo { group, title, item_type, provider_id, virtual_id, input_id, url, series_name, release_date, season, episode } + StrmItemInfo { group, title, item_type, provider_id, virtual_id, input_name, url, series_name, release_date, season, episode } } fn prepare_strm_output_directory(cleanup: bool, path: &PathBuf) -> Result<(), M3uFilterError> { diff --git a/src/repository/storage.rs b/src/repository/storage.rs index 5c20a033f..c8780d7bf 100644 --- a/src/repository/storage.rs +++ b/src/repository/storage.rs @@ -51,7 +51,7 @@ pub fn get_target_storage_path(cfg: &Config, target_name: &str) -> Option std::io::Result { - let name = format!("input_{}", input.name.clone().unwrap_or_else(|| format!("{}", input.id))); + let name = format!("input_{}", &input.name); let path = Path::new(working_dir).join(name); // Create the directory and return the path or propagate the error std::fs::create_dir_all(&path).map(|()| path) diff --git a/src/repository/xtream_repository.rs b/src/repository/xtream_repository.rs index 983e4425e..9bef86508 100644 --- a/src/repository/xtream_repository.rs +++ b/src/repository/xtream_repository.rs @@ -812,7 +812,7 @@ pub async fn xtream_update_input_info_file( Err(err) => Err(info_err!(format!("{err}"))), } } - Ok(None) => Err(notify_err!(format!("Could not create storage path for input {}", &input.name.as_ref().map_or("?", |v| v)))), + Ok(None) => Err(notify_err!(format!("Could not create storage path for input {}", &input.name))), Err(err) => Err(notify_err!(format!("Could not create storage path for input {err}"))), } } @@ -824,7 +824,7 @@ pub async fn xtream_update_input_vod_record_from_wal_file( ) -> Result<(), M3uFilterError> { let record_path = get_input_storage_path(input, &cfg.working_dir).map(|storage_path| xtream_get_record_file_path(&storage_path, PlaylistItemType::Video)) .map_err(|err| notify_err!(format!("Error accessing storage path: {err}"))) - .and_then(|opt| opt.ok_or_else(|| notify_err!(format!("Error accessing storage path for input: {}", input.name.clone().unwrap_or_else(|| input.id.to_string())))))?; + .and_then(|opt| opt.ok_or_else(|| notify_err!(format!("Error accessing storage path for input: {}", &input.name))))?; match cfg.file_locks.write_lock(&record_path).await { Ok(_file_lock) => { @@ -866,7 +866,7 @@ pub async fn xtream_update_input_series_record_from_wal_file( ) -> Result<(), M3uFilterError> { let record_path = get_input_storage_path(input, &cfg.working_dir).map(|storage_path| xtream_get_record_file_path(&storage_path, PlaylistItemType::SeriesInfo)) .map_err(|err| notify_err!(format!("Error accessing storage path: {err}"))) - .and_then(|opt| opt.ok_or_else(|| notify_err!(format!("Error accessing storage path for input: {}", input.name.clone().unwrap_or_else(|| input.id.to_string())))))?; + .and_then(|opt| opt.ok_or_else(|| notify_err!(format!("Error accessing storage path for input: {}", &input.name))))?; match cfg.file_locks.write_lock(&record_path).await { Ok(_file_lock) => { let mut reader = file_reader(open_readonly_file(wal_path).map_err(|err| notify_err!(format!("Could not read series wal info {err}")))?); @@ -903,7 +903,7 @@ pub async fn xtream_update_input_series_episodes_record_from_wal_file( ) -> Result<(), M3uFilterError> { let record_path = get_input_storage_path(input, &cfg.working_dir).map(|storage_path| xtream_get_record_file_path(&storage_path, PlaylistItemType::Series)) .map_err(|err| notify_err!(format!("Error accessing storage path: {err}"))) - .and_then(|opt| opt.ok_or_else(|| notify_err!(format!("Error accessing storage path for input: {}", input.name.clone().unwrap_or_else(|| input.id.to_string())))))?; + .and_then(|opt| opt.ok_or_else(|| notify_err!(format!("Error accessing storage path for input: {}", &input.name))))?; match cfg.file_locks.write_lock(&record_path).await { Ok(_file_lock) => { let mut reader = file_reader(open_readonly_file(wal_path).map_err(|err| notify_err!(format!("Could not read series episode wal info {err}")))?); diff --git a/src/utils/download.rs b/src/utils/download.rs index de215c532..17cb4c619 100644 --- a/src/utils/download.rs +++ b/src/utils/download.rs @@ -138,8 +138,7 @@ fn get_skip_cluster(input: &ConfigInput) -> Vec { } } if skip_cluster.len() == 3 { - let name = input.name.as_ref().map_or_else(|| input.id.to_string(), ToString::to_string); - info!("You have skipped all sections from xtream input {name}"); + info!("You have skipped all sections from xtream input {}", &input.name); } skip_cluster }