diff --git a/CHANGELOG.md b/CHANGELOG.md index 076aa49ae..c06d85a08 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,8 @@ # 3.1.0 (2025-05-xx) - !BREAKING_CHANGE! epg config. Added priority, url config is now renamed to sources, priority is `optional` +The `priority` value determines the importance or order of processing. Lower numbers mean higher priority. That is: +A `priority` of `0` is higher than `1`. **Negative numbers** are allowed and represent even higher priority ```yaml epg: sources: diff --git a/src/api/api_utils.rs b/src/api/api_utils.rs index 2d34c181b..03ebb5a6f 100644 --- a/src/api/api_utils.rs +++ b/src/api/api_utils.rs @@ -19,7 +19,7 @@ use crate::tools::atomic_once_flag::AtomicOnceFlag; use crate::tools::lru_cache::LRUResourceCache; use crate::utils::{DASH_EXT, HLS_EXT}; use crate::utils::default_grace_period_millis; -use crate::utils::file_utils::create_new_file_for_write; +use crate::utils::create_new_file_for_write; use crate::utils::request; use crate::utils::request::{extract_extension_from_url, replace_url_extension, sanitize_sensitive_info}; use crate::utils::human_readable_byte_size; diff --git a/src/api/config_watch.rs b/src/api/config_watch.rs index c836b8d21..a5df5f46c 100644 --- a/src/api/config_watch.rs +++ b/src/api/config_watch.rs @@ -1,6 +1,7 @@ use crate::api::model::app_state::AppState; use crate::tuliprox_error::{TuliproxError, TuliproxErrorKind}; -use crate::utils::config_reader; +use crate::utils; +use crate::utils::is_directory; use log::{debug, error, info}; use notify::event::{AccessKind, AccessMode}; use notify::{recommended_watcher, Event, EventKind, RecursiveMode, Watcher}; @@ -17,7 +18,7 @@ enum ConfigFile { impl ConfigFile { fn load_mappping(app_state: &Arc) -> Result<(), TuliproxError> { - match config_reader::read_mappings(app_state.config.t_mapping_file_path.as_str(), true) { + match utils::read_mappings(app_state.config.t_mapping_file_path.as_str(), true) { Ok(Some(mappings_cfg)) => { app_state.config.set_mappings(&mappings_cfg); info!("Loaded mapping file {}", app_state.config.t_mapping_file_path.as_str()); @@ -35,7 +36,7 @@ impl ConfigFile { } fn load_api_proxy(app_state: &Arc) -> Result<(), TuliproxError> { - match config_reader::read_api_proxy_config(&app_state.config) { + match utils::read_api_proxy_config(&app_state.config) { Ok(()) => { info!("Api Proxy File: {:?}", &app_state.config.t_api_proxy_file_path); } @@ -60,10 +61,11 @@ pub async fn exec_config_watch(app_state: &Arc) -> Result<(), Tuliprox let (tx, rx) = mpsc::channel::>(); let mut files = HashMap::new(); - files.insert(PathBuf::from(&app_state.config.t_config_file_path), ConfigFile::Config); - files.insert(PathBuf::from(&app_state.config.t_api_proxy_file_path), ConfigFile::ApiProxy); - files.insert(PathBuf::from(&app_state.config.t_mapping_file_path), ConfigFile::Mapping); - files.insert(PathBuf::from(&app_state.config.t_sources_file_path), ConfigFile::Sources); + [(&app_state.config.t_config_file_path, ConfigFile::Config), + (&app_state.config.t_api_proxy_file_path, ConfigFile::ApiProxy), + (&app_state.config.t_mapping_file_path, ConfigFile::Mapping), + (&app_state.config.t_sources_file_path, ConfigFile::Sources) + ].into_iter().for_each(|(path, config_file)| { files.insert(PathBuf::from(path), (config_file, is_directory(path))); }); // Use recommended_watcher() to automatically select the best implementation // for your platform. The `EventHandler` passed to this constructor can be a @@ -74,7 +76,8 @@ pub async fn exec_config_watch(app_state: &Arc) -> Result<(), Tuliprox // Add a path to be watched. All files and directories at that path and // below will be monitored for changes. let path = Path::new(app_state.config.t_config_path.as_str()); - watcher.watch(path, RecursiveMode::NonRecursive).map_err(|err| TuliproxError::new(TuliproxErrorKind::Info, format!("Failed to start config file watcher {err}")))?; + let recursive_mode = if utils::is_directory(&app_state.config.t_mapping_file_path) { RecursiveMode::Recursive } else { RecursiveMode::NonRecursive }; + watcher.watch(path, recursive_mode).map_err(|err| TuliproxError::new(TuliproxErrorKind::Info, format!("Failed to start config file watcher {err}")))?; info!("Watching config file changes {path:?}"); let watcher_app_state = Arc::clone(app_state); @@ -85,10 +88,18 @@ pub async fn exec_config_watch(app_state: &Arc) -> Result<(), Tuliprox Ok(event) => { if let EventKind::Access(AccessKind::Close(AccessMode::Write)) = event.kind { for path in event.paths { - if let Some(config_file) = files.get(&path) { + if let Some((config_file, _is_dir)) = files.get(&path) { if let Err(err) = config_file.reload(&path, &watcher_app_state) { error!("Failed to reload config file {path:?}: {err}"); } + } else if recursive_mode == RecursiveMode::Recursive && path.extension().is_some_and(|ext| ext == "yml") { + for (key, (config_file, is_dir)) in files.iter() { + if *is_dir && path.starts_with(key) { + if let Err(err) = config_file.reload(&path, &watcher_app_state) { + error!("Failed to reload config file {path:?}: {err}"); + } + } + } } } } diff --git a/src/api/endpoints/api_playlist_utils.rs b/src/api/endpoints/api_playlist_utils.rs index f25714367..00460fd37 100644 --- a/src/api/endpoints/api_playlist_utils.rs +++ b/src/api/endpoints/api_playlist_utils.rs @@ -1,13 +1,13 @@ use crate::model::{Config, ConfigInput, ConfigTarget, InputType, TargetType}; use crate::model::{M3uPlaylistItem, PlaylistGroup, PlaylistItemType, XtreamCluster}; use crate::repository::{m3u_repository, xtream_repository}; -use crate::utils::file_lock_manager::FileReadGuard; use crate::utils::{m3u, xtream}; use axum::response::IntoResponse; use serde::Serialize; use serde_json::{json, Value}; use std::sync::Arc; use indexmap::IndexMap; +use crate::utils; #[derive(serde::Serialize, serde::Deserialize)] struct PlaylistResponseGroup { @@ -54,7 +54,7 @@ where .collect() } -fn group_playlist_items_by_cluster(params: Option<(FileReadGuard, +fn group_playlist_items_by_cluster(params: Option<(utils::FileReadGuard, impl Iterator)>) -> (Vec, Vec, Vec) { if params.is_none() { diff --git a/src/api/endpoints/v1_api.rs b/src/api/endpoints/v1_api.rs index dfe221766..61081c6ca 100644 --- a/src/api/endpoints/v1_api.rs +++ b/src/api/endpoints/v1_api.rs @@ -13,10 +13,9 @@ use crate::model::{validate_targets, Config, ConfigDto, ConfigInput, ConfigInput use crate::model::{ApiProxyConfig, ApiProxyServerInfo, ProxyUserCredentials, TargetUser}; use crate::processing::processor::playlist; use crate::repository::user_repository::store_api_user; -use crate::utils::config_reader; use crate::utils::ip_checker::get_ips; use crate::utils::request::sanitize_sensitive_info; -use crate::VERSION; +use crate::{utils, VERSION}; use axum::response::IntoResponse; use log::error; use serde_json::json; @@ -24,7 +23,7 @@ use std::collections::{BTreeMap, HashSet}; use std::sync::Arc; fn intern_save_config_api_proxy(backup_dir: &str, api_proxy: &ApiProxyConfig, file_path: &str) -> Option { - match config_reader::save_api_proxy(file_path, backup_dir, api_proxy) { + match utils::save_api_proxy(file_path, backup_dir, api_proxy) { Ok(()) => {} Err(err) => { error!("Failed to save api_proxy.yml {err}"); @@ -35,7 +34,7 @@ fn intern_save_config_api_proxy(backup_dir: &str, api_proxy: &ApiProxyConfig, fi } fn intern_save_config_main(file_path: &str, backup_dir: &str, cfg: &ConfigDto) -> Option { - match config_reader::save_main_config(file_path, backup_dir, cfg) { + match utils::save_main_config(file_path, backup_dir, cfg) { Ok(()) => {} Err(err) => { error!("Failed to save config.yml {err}"); @@ -287,10 +286,10 @@ async fn config( sources: config.sources.iter().map(map_source).collect(), proxy: config.proxy.clone(), ipcheck: config.ipcheck.clone(), - api_proxy: config_reader::read_api_proxy(&app_state.config, false), + api_proxy: utils::read_api_proxy(&app_state.config, false), }; - let mut result = match config_reader::read_config(app_state.config.t_config_path.as_str(), + let mut result = match utils::read_config(app_state.config.t_config_path.as_str(), app_state.config.t_config_file_path.as_str(), app_state.config.t_sources_file_path.as_str(), app_state.config.t_api_proxy_file_path.as_str(), diff --git a/src/api/endpoints/xmltv_api.rs b/src/api/endpoints/xmltv_api.rs index 44ae29249..b77f8b9c5 100644 --- a/src/api/endpoints/xmltv_api.rs +++ b/src/api/endpoints/xmltv_api.rs @@ -17,8 +17,7 @@ use crate::model::{Config, ConfigTarget, TargetOutput}; use crate::repository::m3u_repository::m3u_get_epg_file_path; use crate::repository::storage::get_target_storage_path; use crate::repository::xtream_repository::{xtream_get_epg_file_path, xtream_get_storage_path}; -use crate::utils::file_utils; -use crate::utils::file_utils::file_reader; +use crate::utils; pub fn get_empty_epg_response() -> impl axum::response::IntoResponse + Send { axum::response::Response::builder() @@ -46,7 +45,7 @@ fn time_correct(date_time: &str, correction: &TimeDelta) -> String { } fn get_epg_path_for_target_of_type(target_name: &str, epg_path: PathBuf) -> Option { - if file_utils::path_exists(&epg_path) { + if utils::path_exists(&epg_path) { return Some(epg_path); } trace!("Cant find epg file for {target_name} target: {}", epg_path.to_str().unwrap_or("?")); @@ -108,7 +107,7 @@ async fn serve_epg(epg_path: &Path, user: &ProxyUserCredentials) -> impl axum::r } fn serve_epg_with_timeshift(epg_file: File, offset_minutes: i32) -> impl axum::response::IntoResponse + Send { - let reader = file_reader(epg_file); + let reader = utils::file_reader(epg_file); let encoder = GzEncoder::new(Vec::with_capacity(4096), Compression::default()); let mut xml_reader = Reader::from_reader(reader); let mut xml_writer = Writer::new(encoder); diff --git a/src/main.rs b/src/main.rs index 7a48c5f7d..1d9896c58 100644 --- a/src/main.rs +++ b/src/main.rs @@ -12,9 +12,7 @@ include_modules!(); use crate::auth::password::generate_password; use crate::model::{validate_targets, Config, HealthcheckConfig, Healthcheck, ProcessTargets}; use crate::processing::processor::playlist; -use crate::utils::config_reader::config_file_reader; -use utils::config_reader; -use crate::utils::file_utils; +use crate::utils::{config_file_reader, resolve_env_var}; use crate::utils::request::{create_client, set_sanitize_sensitive_info}; use chrono::{DateTime, Utc}; use clap::Parser; @@ -93,10 +91,10 @@ fn main() { return; } - let config_path: String = file_utils::resolve_directory_path(&args.config_path.unwrap_or_else(file_utils::get_default_config_path)); - let config_file: String = args.config_file.unwrap_or_else(|| file_utils::get_default_config_file_path(&config_path)); - let api_proxy_file = args.api_proxy.unwrap_or_else(|| file_utils::get_default_api_proxy_config_path(config_path.as_str())); - let mappings_file = args.mapping_file.unwrap_or_else(|| file_utils::get_default_mappings_path(config_path.as_str())); + let config_path: String = utils::resolve_directory_path(&resolve_env_var(&args.config_path.unwrap_or_else(utils::get_default_config_path))); + let config_file: String = resolve_env_var(&args.config_file.unwrap_or_else(|| utils::get_default_config_file_path(&config_path))); + let api_proxy_file = resolve_env_var(&args.api_proxy.unwrap_or_else(|| utils::get_default_api_proxy_config_path(config_path.as_str()))); + let mappings_file = resolve_env_var(&args.mapping_file.unwrap_or_else(|| utils::get_default_mappings_path(config_path.as_str()))); init_logger(args.log_level.as_ref(), config_file.as_str()); @@ -109,8 +107,8 @@ fn main() { healthcheck(config_file.as_str()); } - let sources_file: String = args.source_file.unwrap_or_else(|| file_utils::get_default_sources_file_path(&config_path)); - let cfg = config_reader::read_config(config_path.as_str(), config_file.as_str(), + let sources_file: String = args.source_file.unwrap_or_else(|| utils::get_default_sources_file_path(&config_path)); + let cfg = utils::read_config(config_path.as_str(), config_file.as_str(), sources_file.as_str(), api_proxy_file.as_str(), mappings_file.as_str(), true).unwrap_or_else(|err| exit!("{}", err)); @@ -122,7 +120,7 @@ fn main() { let targets = validate_targets(args.target.as_ref(), &cfg.sources).unwrap_or_else(|err| exit!("{}", err)); - match config_reader::read_mappings(mappings_file.as_str(), true) { + match utils::read_mappings(mappings_file.as_str(), true) { Ok(Some(mappings)) => cfg.set_mappings(&mappings), Ok(None) => {}, Err(err) => exit!("{err}"), @@ -145,7 +143,7 @@ fn main() { let rt = tokio::runtime::Runtime::new().unwrap(); let () = rt.block_on(async { if args.server { - match config_reader::read_api_proxy_config(&cfg) { + match utils::read_api_proxy_config(&cfg) { Ok(()) => {} Err(err) => exit!("{err}"), } diff --git a/src/model/api_proxy.rs b/src/model/api_proxy.rs index 9ad8cfcc9..1d8f75bb0 100644 --- a/src/model/api_proxy.rs +++ b/src/model/api_proxy.rs @@ -2,8 +2,7 @@ use crate::api::model::app_state::AppState; use crate::tuliprox_error::{create_tuliprox_error_result, info_err, TuliproxError, TuliproxErrorKind}; use crate::model::{Config, ClusterFlags}; use crate::repository::user_repository::{backup_api_user_db_file, get_api_user_db_path, load_api_user, merge_api_user}; -use crate::utils::default_as_true; -use crate::utils::config_reader; +use crate::utils::{default_as_true, save_api_proxy}; use chrono::Local; use enum_iterator::Sequence; use log::debug; @@ -14,6 +13,7 @@ use std::fs; use std::str::FromStr; use serde::{Deserialize, Deserializer, Serialize, Serializer}; use crate::model::PlaylistItemType; +use crate::utils; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)] pub enum UserConnectionPermission { @@ -413,7 +413,7 @@ impl ApiProxyConfig { let api_proxy_file = cfg.t_api_proxy_file_path.as_str(); let backup_dir = cfg.backup_dir.as_ref().unwrap().as_str(); self.user = vec![]; - if let Err(err) = config_reader::save_api_proxy(api_proxy_file, backup_dir, self) { + if let Err(err) = utils::save_api_proxy(api_proxy_file, backup_dir, self) { errors.push(format!("Error saving api proxy file: {err}")); } } @@ -447,7 +447,7 @@ impl ApiProxyConfig { } let api_proxy_file = cfg.t_api_proxy_file_path.as_str(); let backup_dir = cfg.backup_dir.as_ref().unwrap().as_str(); - if let Err(err) = config_reader::save_api_proxy(api_proxy_file, backup_dir, self) { + if let Err(err) = save_api_proxy(api_proxy_file, backup_dir, self) { errors.push(format!("Error saving api proxy file: {err}")); } else { backup_api_user_db_file(cfg, &user_db_path); diff --git a/src/model/config.rs b/src/model/config.rs index 3af072e2e..cc750cc56 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -16,8 +16,6 @@ use crate::model::{ConfigInput, ConfigInputOptions, ConfigSource, ConfigTarget, use crate::tuliprox_error::create_tuliprox_error_result; use crate::tuliprox_error::{TuliproxError, TuliproxErrorKind}; use crate::utils::exit; -use crate::utils::file_lock_manager::FileLockManager; -use crate::utils::file_utils; use crate::utils::{default_as_default, default_connect_timeout_secs, default_grace_period_millis, default_grace_period_timeout_secs}; use crate::utils::{parse_size_base_2, parse_to_kbps}; @@ -38,6 +36,7 @@ macro_rules! valid_property { }}; } pub use valid_property; +use crate::utils; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Sequence, Eq, PartialEq)] pub enum ItemField { @@ -390,7 +389,7 @@ pub struct Config { #[serde(skip)] pub t_api_proxy_file_path: String, #[serde(skip)] - pub file_locks: Arc, + pub file_locks: Arc, #[serde(skip)] pub t_channel_unavailable_video: Option>>, #[serde(skip)] @@ -645,7 +644,7 @@ impl Config { */ pub fn prepare(&mut self, include_computed: bool) -> Result<(), TuliproxError> { let work_dir = &self.working_dir; - self.working_dir = file_utils::resolve_directory_path(work_dir); + self.working_dir = utils::resolve_directory_path(work_dir); if include_computed { self.t_access_token_secret = generate_secret(); self.t_encrypt_secret = <&[u8] as TryInto<[u8; 16]>>::try_into(&generate_secret()[0..16]).map_err(|err| TuliproxError::new(TuliproxErrorKind::Info, err.to_string()))?; @@ -761,8 +760,8 @@ impl Config { if let Some(custom_stream_response) = self.custom_stream_response.as_ref() { fn load_and_set_file(path: Option<&String>, working_dir: &str) -> Option>> { path.as_ref() - .map(|file| file_utils::make_absolute_path(file, working_dir)) - .and_then(|absolute_path| match file_utils::read_file_as_bytes(&PathBuf::from(&absolute_path)) { + .map(|file| utils::make_absolute_path(file, working_dir)) + .and_then(|absolute_path| match utils::read_file_as_bytes(&PathBuf::from(&absolute_path)) { Ok(data) => Some(Arc::new(data)), Err(err) => { error!("Failed to load file: {absolute_path} {err}"); @@ -779,7 +778,7 @@ impl Config { fn prepare_api_web_root(&mut self) { if !self.api.web_root.is_empty() { - self.api.web_root = file_utils::make_absolute_path(&self.api.web_root, &self.working_dir); + self.api.web_root = utils::make_absolute_path(&self.api.web_root, &self.working_dir); } } diff --git a/src/model/config_input.rs b/src/model/config_input.rs index 3f1375e6e..2dbb97d57 100644 --- a/src/model/config_input.rs +++ b/src/model/config_input.rs @@ -1,6 +1,5 @@ use crate::tuliprox_error::{create_tuliprox_error_result, handle_tuliprox_error_result_list, info_err, TuliproxError, TuliproxErrorKind}; use crate::model::{EpgConfig, EpgSource}; -use crate::utils::config_reader::csv_read_inputs; use crate::utils::default_as_true; use crate::utils::get_trimmed_string; use crate::utils::request::{get_base_url_from_str, get_credentials_from_url, get_credentials_from_url_str}; @@ -10,6 +9,7 @@ use std::collections::HashMap; use std::fmt::Display; use std::str::FromStr; use url::Url; +use crate::utils; macro_rules! check_input_credentials { ($this:ident, $input_type:expr) => { @@ -342,7 +342,7 @@ impl ConfigInput { InputType::Xtream }; - match csv_read_inputs(self) { + match utils::csv_read_inputs(self) { Ok(mut batch_aliases) => { if !batch_aliases.is_empty() { batch_aliases.reverse(); diff --git a/src/model/config_web_auth.rs b/src/model/config_web_auth.rs index b283614d9..ab85be8d4 100644 --- a/src/model/config_web_auth.rs +++ b/src/model/config_web_auth.rs @@ -4,8 +4,7 @@ use std::io::BufRead; use std::path::PathBuf; use crate::auth::user::UserCredential; use crate::tuliprox_error::{TuliproxError, TuliproxErrorKind, create_tuliprox_error_result}; -use crate::utils::file_utils; -use crate::utils::file_utils::file_reader; +use crate::utils; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] #[serde(deny_unknown_fields)] @@ -22,20 +21,20 @@ pub struct WebAuthConfig { impl WebAuthConfig { pub fn prepare(&mut self, config_path: &str) -> Result<(), TuliproxError> { - let userfile_name = self.userfile.as_ref().map_or_else(|| file_utils::get_default_user_file_path(config_path), std::borrow::ToOwned::to_owned); + let userfile_name = self.userfile.as_ref().map_or_else(|| utils::get_default_user_file_path(config_path), std::borrow::ToOwned::to_owned); self.userfile = Some(userfile_name.clone()); let mut userfile_path = PathBuf::from(&userfile_name); - if !file_utils::path_exists(&userfile_path) { + if !utils::path_exists(&userfile_path) { userfile_path = PathBuf::from(config_path).join(&userfile_name); - if !file_utils::path_exists(&userfile_path) { + if !utils::path_exists(&userfile_path) { return create_tuliprox_error_result!(TuliproxErrorKind::Info, "Could not find userfile {}", &userfile_name); } } if let Ok(file) = File::open(&userfile_path) { let mut users = vec![]; - let reader = file_reader(file); + let reader = utils::file_reader(file); // TODO maybe print out errors for credentials in reader.lines().map_while(Result::ok) { let mut parts = credentials.split(':'); diff --git a/src/model/mapping.rs b/src/model/mapping.rs index 0dbe42114..a17250491 100644 --- a/src/model/mapping.rs +++ b/src/model/mapping.rs @@ -452,18 +452,19 @@ pub struct Mapping { pub id: String, #[serde(default)] pub match_as_ascii: bool, - pub mapper: Vec, + pub mapper: Option>, pub counter: Option>, #[serde(skip_serializing, skip_deserializing)] pub t_counter: Option>, - } impl Mapping { pub fn prepare(&mut self, templates: Option<&Vec>, tags: Option<&Vec>) -> Result<(), TuliproxError> { - for mapper in &mut self.mapper { - handle_tuliprox_error_result!(TuliproxErrorKind::Info, mapper.prepare(templates, tags)); + if let Some(mapper_list) = &mut self.mapper { + for mapper in mapper_list { + handle_tuliprox_error_result!(TuliproxErrorKind::Info, mapper.prepare(templates, tags)); + } } if let Some(counter_def_list) = &self.counter { diff --git a/src/processing/playlist_watch.rs b/src/processing/playlist_watch.rs index 741a6e44e..f0420fff5 100644 --- a/src/processing/playlist_watch.rs +++ b/src/processing/playlist_watch.rs @@ -5,9 +5,8 @@ use log::{error, info}; use crate::messaging::{MsgKind, send_message}; use crate::model::Config; use crate::model::PlaylistGroup; +use crate::utils; use crate::utils::{bincode_deserialize, bincode_serialize}; -use crate::utils::file_utils; -use crate::utils::file_utils::sanitize_filename; pub fn process_group_watch(client: &Arc, cfg: &Config, target_name: &str, pl: &PlaylistGroup) { let mut new_tree = BTreeSet::new(); @@ -17,8 +16,8 @@ pub fn process_group_watch(client: &Arc, cfg: &Config, target_n new_tree.insert(title); }); - let watch_filename = format!("{}/{}.bin", sanitize_filename(target_name), sanitize_filename(&pl.title)); - match file_utils::get_file_path(&cfg.working_dir, Some(std::path::PathBuf::from(&watch_filename))) { + let watch_filename = format!("{}/{}.bin", utils::sanitize_filename(target_name), utils::sanitize_filename(&pl.title)); + match utils::get_file_path(&cfg.working_dir, Some(std::path::PathBuf::from(&watch_filename))) { Some(path) => { let save_path = path.as_path(); let mut changed = false; diff --git a/src/processing/processor/playlist.rs b/src/processing/processor/playlist.rs index 81f543c51..ea6473e0e 100644 --- a/src/processing/processor/playlist.rs +++ b/src/processing/processor/playlist.rs @@ -10,7 +10,6 @@ use std::thread; use tokio::sync::Mutex; use crate::foundation::filter::{get_field_value, set_field_value, MockValueProcessor, ValueProvider}; -use crate::tuliprox_error::{get_errors_notify_message, notify_err, TuliproxError, TuliproxErrorKind}; use crate::messaging::{send_message, MsgKind}; use crate::model::{ConfigTarget, InputType, ItemField, ProcessTargets, ProcessingOrder}; use crate::model::{CounterModifier, Mapping, MappingValueProcessor}; @@ -21,8 +20,9 @@ use crate::processing::processor::affix::apply_affixes; use crate::processing::processor::xtream_series::playlist_resolve_series; use crate::processing::processor::xtream_vod::playlist_resolve_vod; use crate::repository::playlist_repository::persist_playlist; -use crate::utils::default_as_default; +use crate::tuliprox_error::{get_errors_notify_message, notify_err, TuliproxError, TuliproxErrorKind}; use crate::utils::debug_if_enabled; +use crate::utils::default_as_default; use deunicode::deunicode; use log::{debug, error, info, log_enabled, trace, warn, Level}; use std::time::Instant; @@ -130,25 +130,27 @@ macro_rules! apply_pattern { } fn map_channel(mut channel: PlaylistItem, mapping: &Mapping) -> PlaylistItem { - if !mapping.mapper.is_empty() { - let header = &channel.header; - let channel_name = if mapping.match_as_ascii { deunicode(&header.name) } else { header.name.to_string() }; - if mapping.match_as_ascii && log_enabled!(Level::Trace) { trace!("Decoded {} for matching to {}", &header.name, &channel_name); } - // let ref_chan = &mut channel; - let ref_chan = &mut channel; - let mut mock_processor = MockValueProcessor {}; - for m in &mapping.mapper { - let provider = ValueProvider { pli: &ref_chan.clone() }; - let mut processor = MappingValueProcessor { pli: ref_chan, mapper: m }; - match &m.t_filter { - Some(filter) => { - if filter.filter(&provider, &mut mock_processor) { + if let Some(mapper) = &mapping.mapper { + if !mapper.is_empty() { + let header = &channel.header; + let channel_name = if mapping.match_as_ascii { deunicode(&header.name) } else { header.name.to_string() }; + if mapping.match_as_ascii && log_enabled!(Level::Trace) { trace!("Decoded {} for matching to {}", &header.name, &channel_name); } + // let ref_chan = &mut channel; + let ref_chan = &mut channel; + let mut mock_processor = MockValueProcessor {}; + for m in mapper { + let provider = ValueProvider { pli: &ref_chan.clone() }; + let mut processor = MappingValueProcessor { pli: ref_chan, mapper: m }; + match &m.t_filter { + Some(filter) => { + if filter.filter(&provider, &mut mock_processor) { + apply_pattern!(&m.t_pattern, &provider, &mut processor); + } + } + _ => { apply_pattern!(&m.t_pattern, &provider, &mut processor); } } - _ => { - apply_pattern!(&m.t_pattern, &provider, &mut processor); - } } } } @@ -159,7 +161,8 @@ fn map_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option if let Some(mappings) = target.t_mapping.load().as_ref() { let new_playlist: Vec = playlist.iter().map(|playlist_group| { let mut grp = playlist_group.clone(); - mappings.iter().filter(|&mapping| !mapping.mapper.is_empty()).for_each(|mapping| + mappings.iter().filter(|&mapping| mapping.mapper.as_ref().map_or(false, |v| !v.is_empty())) + .for_each(|mapping| grp.channels = grp.channels.drain(..).map(|chan| map_channel(chan, mapping)).collect()); grp }).collect(); diff --git a/src/processing/processor/xtream.rs b/src/processing/processor/xtream.rs index 36d7d6f74..353bdb89d 100644 --- a/src/processing/processor/xtream.rs +++ b/src/processing/processor/xtream.rs @@ -12,7 +12,7 @@ use std::sync::Arc; use crate::repository::bplustree::BPlusTree; use crate::repository::storage_const; use crate::repository::xtream_repository::xtream_get_record_file_path; -use crate::utils::file_utils::append_or_crate_file; +use crate::utils; use crate::utils::xtream; pub(in crate::processing) async fn playlist_resolve_download_playlist_item(client: Arc, pli: &PlaylistItem, input: &ConfigInput, errors: &mut Vec, resolve_delay: u16, cluster: XtreamCluster) -> Option { @@ -47,7 +47,7 @@ pub(in crate::processing) fn create_resolve_episode_wal_files(cfg: &Config, inpu match get_input_storage_path(&input.name, &cfg.working_dir) { Ok(storage_path) => { let info_path = storage_path.join(format!("{}.{}", crate::model::XC_FILE_SERIES_EPISODE_RECORD, storage_const::FILE_SUFFIX_WAL)); - let info_file = append_or_crate_file(&info_path).ok()?; + let info_file = utils::append_or_crate_file(&info_path).ok()?; Some((info_file, info_path)) } Err(_) => None @@ -64,8 +64,8 @@ pub(in crate::processing) fn create_resolve_info_wal_files(cfg: &Config, input: } { let content_path = storage_path.join(format!("{file_prefix}_content.{}", storage_const::FILE_SUFFIX_WAL)); let info_path = storage_path.join(format!("{file_prefix}_record.{}", storage_const::FILE_SUFFIX_WAL)); - let content_file = append_or_crate_file(&content_path).ok()?; - let info_file = append_or_crate_file(&info_path).ok()?; + let content_file = utils::append_or_crate_file(&content_path).ok()?; + let info_file = utils::append_or_crate_file(&info_path).ok()?; return Some((content_file, info_file, content_path, info_path)); } None diff --git a/src/processing/processor/xtream_series.rs b/src/processing/processor/xtream_series.rs index ae4d0647f..711bdee27 100644 --- a/src/processing/processor/xtream_series.rs +++ b/src/processing/processor/xtream_series.rs @@ -16,8 +16,8 @@ use std::sync::Arc; use std::time::Instant; use log::{info, log_enabled, Level}; use crate::model::{XtreamSeriesEpisode, XtreamSeriesInfoEpisode}; +use crate::utils; use crate::utils::bincode_serialize; -use crate::utils::file_utils::file_writer; create_resolve_options_function_for_xtream_target!(series); @@ -63,8 +63,8 @@ async fn playlist_resolve_series_info(client: Arc, cfg: &Config let Some((wal_content_file, wal_record_file, wal_content_path, wal_record_path)) = create_resolve_info_wal_files(cfg, fpl.input, XtreamCluster::Series) else { return !processed_info_ids.is_empty(); }; - let mut content_writer = file_writer(&wal_content_file); - let mut record_writer = file_writer(&wal_record_file); + let mut content_writer = utils::file_writer(&wal_content_file); + let mut record_writer = utils::file_writer(&wal_record_file); let mut content_updated = false; // TODO merge both filters to one @@ -153,7 +153,7 @@ async fn process_series_info( errors.push(notify_err!("Could not create wal file for series episodes record".to_string())); return result; }; - let mut wal_writer = file_writer(&wal_file); + let mut wal_writer = utils::file_writer(&wal_file); for plg in fpl .playlistgroups diff --git a/src/processing/processor/xtream_vod.rs b/src/processing/processor/xtream_vod.rs index 2b47dc677..b5a287b02 100644 --- a/src/processing/processor/xtream_vod.rs +++ b/src/processing/processor/xtream_vod.rs @@ -13,7 +13,7 @@ use std::io::{BufWriter, Write}; use std::sync::Arc; use std::time::Instant; use log::{info, log_enabled, Level}; -use crate::utils::file_utils::file_writer; +use crate::utils; create_resolve_options_function_for_xtream_target!(vod); @@ -71,8 +71,8 @@ pub async fn playlist_resolve_vod(client: Arc, cfg: &Config, ta else { return; }; let mut processed_info_ids = read_processed_vod_info_ids(cfg, errors, fpl).await; - let mut content_writer = file_writer(&wal_content_file); - let mut record_writer = file_writer(&wal_record_file); + let mut content_writer = utils::file_writer(&wal_content_file); + let mut record_writer = utils::file_writer(&wal_record_file); let mut content_updated = false; // TODO merge both filters to one diff --git a/src/repository/bplustree.rs b/src/repository/bplustree.rs index bf83bcc6f..614677767 100644 --- a/src/repository/bplustree.rs +++ b/src/repository/bplustree.rs @@ -4,12 +4,12 @@ use std::marker::PhantomData; use std::mem::size_of; use std::path::Path; use crate::tuliprox_error::{str_to_io_error, to_io_error}; -use crate::utils::file_utils::{file_reader, file_writer, open_read_write_file, rename_or_copy}; use log::error; use ruzstd::decoding::StreamingDecoder; use ruzstd::encoding::{compress_to_vec, CompressionLevel}; use serde::{Deserialize, Serialize}; use tempfile::NamedTempFile; +use crate::utils; use crate::utils::{bincode_deserialize, bincode_serialize}; const BLOCK_SIZE: usize = 4096; @@ -474,13 +474,13 @@ where pub fn store(&mut self, filepath: &Path) -> io::Result { if self.dirty { let tempfile = NamedTempFile::new()?; - let mut file = file_writer(&tempfile); //create_new_file_for_write(&tempfile)?); + let mut file = utils::file_writer(&tempfile); //create_new_file_for_write(&tempfile)?); let mut buffer = vec![0u8; BLOCK_SIZE]; match self.root.serialize_to_block(&mut file, &mut buffer, 0u64) { Ok(result) => { file.flush()?; drop(file); - if let Err(err) = rename_or_copy(tempfile.path(), filepath, false) { + if let Err(err) = utils::rename_or_copy(tempfile.path(), filepath, false) { return Err(str_to_io_error(&format!("Temp file rename/copy did not work {} {err}", tempfile.path().to_string_lossy()))); } self.dirty = false; @@ -497,7 +497,7 @@ where pub fn load(filepath: &Path) -> io::Result { let file = is_file_valid(File::open(filepath)?)?; - let mut reader = file_reader(file); + let mut reader = utils::file_reader(file); let mut buffer = vec![0u8; BLOCK_SIZE]; let (root, _) = BPlusTreeNode::deserialize_from_block(&mut reader, &mut buffer, 0, true)?; Ok(Self::new_with_root(root)) @@ -582,7 +582,7 @@ where pub fn try_from_file(file: File) -> io::Result { let file = is_file_valid(file)?; Ok(Self { - file: file_reader(file), + file: utils::file_reader(file), _marker_k: PhantomData, _marker_v: PhantomData, }) @@ -620,7 +620,7 @@ where if !filepath.exists() { return Err(io::Error::new(io::ErrorKind::NotFound, format!("File not found {}", filepath.to_str().unwrap_or("?")))); } - let file = is_file_valid(open_read_write_file(filepath)?)?; + let file = is_file_valid(utils::open_read_write_file(filepath)?)?; Ok(Self { file, _marker_k: PhantomData, @@ -629,7 +629,7 @@ where } pub fn query(&mut self, key: &K) -> Option { - let mut reader = file_reader(&mut self.file); + let mut reader = utils::file_reader(&mut self.file); query_tree(&mut reader, key) } @@ -643,7 +643,7 @@ where pub fn update(&mut self, key: &K, value: V) -> io::Result { let mut offset = 0; let mut buffer = vec![0u8; BLOCK_SIZE]; - let mut reader = file_reader(&mut self.file); + let mut reader = utils::file_reader(&mut self.file); loop { match BPlusTreeNode::::deserialize_from_block(&mut reader, &mut buffer, offset, false) { Ok((mut node, pointers)) => { diff --git a/src/repository/indexed_document.rs b/src/repository/indexed_document.rs index 7bcfe96c2..613061b17 100644 --- a/src/repository/indexed_document.rs +++ b/src/repository/indexed_document.rs @@ -5,13 +5,12 @@ use std::marker::PhantomData; use std::path::{Path, PathBuf}; use crate::repository::bplustree::{BPlusTree, BPlusTreeQuery}; -use crate::utils::file_utils; use log::error; use serde::{Deserialize, Serialize}; use tempfile::NamedTempFile; use crate::tuliprox_error::{str_to_io_error, to_io_error}; +use crate::utils; use crate::utils::{bincode_deserialize, bincode_serialize}; -use crate::utils::file_utils::{create_new_file_for_read_write, file_reader, file_writer, open_read_write_file, open_readonly_file, rename_or_copy}; const BLOCK_SIZE: usize = 4096; const LEN_SIZE: usize = 4; @@ -94,9 +93,9 @@ where fn new_with_mode(main_path: PathBuf, index_path: PathBuf, append: bool) -> Result { let append_mode = append && main_path.exists(); let mut main_file = if append_mode { - open_read_write_file(&main_path) + utils::open_read_write_file(&main_path) } else { - create_new_file_for_read_write(&main_path) + utils::create_new_file_for_read_write(&main_path) }?; // Retrieve file size and convert to `u32` for `main_offset`, if possible @@ -200,7 +199,7 @@ where let encoded_bytes_len = SizeType::try_from(encoded_bytes.len()).map_err(to_io_error)?; self.main_file.write_all(&encoded_bytes_len.to_le_bytes())?; - match file_utils::check_write(&self.main_file.write_all(&encoded_bytes)) { + match utils::check_write(&self.main_file.write_all(&encoded_bytes)) { Ok(()) => { if new_record_appended { self.index_tree.insert(doc_id, self.main_offset); @@ -249,11 +248,11 @@ where { pub fn new(main_path: &Path, index_path: &Path) -> Result { if main_path.exists() && index_path.exists() { - let main_file = open_readonly_file(main_path)?; + let main_file = utils::open_readonly_file(main_path)?; let index_tree = IndexedDocumentIndex::::load(index_path)?; Ok(Self { - main_file: file_reader(main_file), + main_file: utils::file_reader(main_file), index_tree, t_type: PhantomData, }) @@ -309,7 +308,7 @@ where Ok(file) => { Ok(Self { main_path: main_path.to_path_buf(), - main_file: file_reader(file), + main_file: utils::file_reader(file), offsets, index: 0, failed: false, @@ -430,7 +429,7 @@ where if main_path.exists() && index_path.exists() { // Attempt to open the main file in the specified mode (append or not) - let main_file = open_read_write_file(&main_path)?; + let main_file = utils::open_read_write_file(&main_path)?; // Retrieve file size and convert to `u32` for `main_file`, if possible let size = main_file @@ -468,7 +467,7 @@ where self.index_tree.traverse(|keys, values| { keys.iter().zip(values.iter()).for_each(|(key, &offset)| key_offset.push((key.clone(), offset))); }); - let mut gc_writer = file_writer(&gc_file); + let mut gc_writer = utils::file_writer(&gc_file); let fragmented_byte = 0u8.to_le_bytes(); gc_writer.write_all(&fragmented_byte)?; @@ -499,7 +498,7 @@ where gc_writer.flush()?; } - rename_or_copy(gc_path, &self.main_path, false)?; + utils::rename_or_copy(gc_path, &self.main_path, false)?; self.index_tree.store(&self.index_path)?; Ok(()) @@ -515,7 +514,7 @@ mod tests { use crate::model::XtreamPlaylistItem; // use crate::model::XtreamPlaylistItem; use crate::repository::indexed_document::{IndexedDocumentGarbageCollector, IndexedDocumentIterator, IndexedDocumentWriter}; - use crate::utils::config_reader::resolve_env_var; + use crate::utils::resolve_env_var; // Example usage with a simple struct #[derive(Serialize, Deserialize, Debug, Clone, PartialEq, Eq)] diff --git a/src/repository/kodi_repository.rs b/src/repository/kodi_repository.rs index 0f4faf36e..f27b8de5a 100644 --- a/src/repository/kodi_repository.rs +++ b/src/repository/kodi_repository.rs @@ -11,8 +11,7 @@ use crate::repository::storage::{ensure_target_storage_path, get_input_storage_p use crate::repository::storage_const; use crate::repository::xtream_repository::{xtream_get_record_file_path, InputVodInfoRecord}; use crate::utils::{KodiStyle, CONSTANTS}; -use crate::utils::file_lock_manager::FileReadGuard; -use crate::utils::file_utils; +use crate::utils::FileReadGuard; use crate::utils::request::extract_extension_from_url; use chrono::Datelike; use filetime::{set_file_times, FileTime}; @@ -24,6 +23,7 @@ use std::path::{Path, PathBuf}; use std::sync::Arc; use tokio::fs::{create_dir_all, remove_dir, remove_file, File}; use tokio::io::{AsyncBufReadExt, AsyncReadExt, AsyncWriteExt, BufReader, BufWriter}; +use crate::utils; fn sanitize_for_filename(text: &str, underscore_whitespace: bool) -> String { text.trim() @@ -580,7 +580,7 @@ pub async fn kodi_write_strm_playlist( return Ok(()); } - let Some(root_path) = file_utils::get_file_path( + let Some(root_path) = utils::get_file_path( &cfg.working_dir, Some(std::path::PathBuf::from(&target_output.directory)), ) else { diff --git a/src/repository/m3u_playlist_iterator.rs b/src/repository/m3u_playlist_iterator.rs index 577ea5de4..ccc5172e2 100644 --- a/src/repository/m3u_playlist_iterator.rs +++ b/src/repository/m3u_playlist_iterator.rs @@ -8,7 +8,7 @@ use crate::repository::m3u_repository::m3u_get_file_paths; use crate::repository::storage::ensure_target_storage_path; use crate::repository::storage_const; use crate::repository::user_repository::user_get_bouquet_filter; -use crate::utils::file_lock_manager::FileReadGuard; +use crate::utils::FileReadGuard; use std::collections::HashSet; #[allow(clippy::struct_excessive_bools)] diff --git a/src/repository/m3u_repository.rs b/src/repository/m3u_repository.rs index 0d674e164..2693aaacd 100644 --- a/src/repository/m3u_repository.rs +++ b/src/repository/m3u_repository.rs @@ -6,15 +6,13 @@ use crate::model::{M3uPlaylistItem, PlaylistGroup, PlaylistItem, PlaylistItemTyp use crate::repository::indexed_document::{IndexedDocumentDirectAccess, IndexedDocumentIterator, IndexedDocumentWriter}; use crate::repository::m3u_playlist_iterator::{M3uPlaylistM3uTextIterator}; use crate::repository::storage::{get_target_storage_path}; -use crate::utils::file_lock_manager::FileReadGuard; -use crate::utils::file_utils; -use crate::utils::file_utils::file_writer; use log::error; use std::fs::File; use std::io::{Error, Write}; use std::path::{Path, PathBuf}; use std::sync::Arc; use crate::repository::storage_const; +use crate::utils; macro_rules! cant_write_result { ($path:expr, $err:expr) => { @@ -30,15 +28,15 @@ pub fn m3u_get_file_paths(target_path: &Path) -> (PathBuf, PathBuf) { pub fn m3u_get_epg_file_path(target_path: &Path) -> PathBuf { let path = target_path.join(PathBuf::from(format!("{}.{}", storage_const::FILE_M3U, storage_const::FILE_SUFFIX_DB))); - file_utils::add_prefix_to_filename(&path, "epg_", Some("xml")) + utils::add_prefix_to_filename(&path, "epg_", Some("xml")) } fn persist_m3u_playlist_as_text(cfg: &Config, target: &ConfigTarget, target_output: &M3uTargetOutput, m3u_playlist: &Vec) { if let Some(filename) = target_output.filename.as_ref() { - if let Some(m3u_filename) = file_utils::get_file_path(&cfg.working_dir, Some(PathBuf::from(filename))) { + if let Some(m3u_filename) = utils::get_file_path(&cfg.working_dir, Some(PathBuf::from(filename))) { match File::create(&m3u_filename) { Ok(file) => { - let mut buf_writer = file_writer(&file); + let mut buf_writer = utils::file_writer(&file); let _ = buf_writer.write(b"#EXTM3U\n"); for m3u in m3u_playlist { let _ = buf_writer.write(m3u.to_m3u(target.options.as_ref(), false).to_string().as_bytes()); @@ -101,7 +99,7 @@ pub async fn m3u_get_item_for_stream_id(stream_id: u32, cfg: &Config, target: &C } } -pub async fn iter_raw_m3u_playlist(config: &Arc, target: &ConfigTarget) -> Option<(FileReadGuard, impl Iterator)> { +pub async fn iter_raw_m3u_playlist(config: &Arc, target: &ConfigTarget) -> Option<(utils::FileReadGuard, impl Iterator)> { let target_path = get_target_storage_path(config, target.name.as_str())?; let (m3u_path, idx_path) = m3u_get_file_paths(&target_path); if !m3u_path.exists() || !idx_path.exists() { diff --git a/src/repository/playlist_repository.rs b/src/repository/playlist_repository.rs index 94c411ec3..b308f4c46 100644 --- a/src/repository/playlist_repository.rs +++ b/src/repository/playlist_repository.rs @@ -9,9 +9,9 @@ use crate::repository::m3u_repository::m3u_write_playlist; use crate::repository::storage::{ensure_target_storage_path, get_target_id_mapping_file}; use crate::repository::target_id_mapping::TargetIdMapping; use crate::repository::xtream_repository::xtream_write_playlist; -use crate::utils::file_lock_manager::FileWriteGuard; use crate::utils::request::{is_dash_url, is_hls_url}; use std::path::Path; +use crate::utils; pub async fn persist_playlist(playlist: &mut [PlaylistGroup], epg: Option<&Epg>, target: &ConfigTarget, cfg: &Config) -> Result<(), Vec> { @@ -72,7 +72,7 @@ pub async fn persist_playlist(playlist: &mut [PlaylistGroup], epg: Option<&Epg>, if errors.is_empty() { Ok(()) } else { Err(errors) } } -pub async fn get_target_id_mapping(cfg: &Config, target_path: &Path) -> (TargetIdMapping, FileWriteGuard) { +pub async fn get_target_id_mapping(cfg: &Config, target_path: &Path) -> (TargetIdMapping, utils::FileWriteGuard) { let target_id_mapping_file = get_target_id_mapping_file(target_path); let file_lock = cfg.file_locks.write_lock(&target_id_mapping_file).await; (TargetIdMapping::new(&target_id_mapping_file), file_lock) diff --git a/src/repository/storage.rs b/src/repository/storage.rs index ad3c87bbc..b9a35406d 100644 --- a/src/repository/storage.rs +++ b/src/repository/storage.rs @@ -5,7 +5,7 @@ use crate::model::{Config}; use crate::model::UUIDType; use crate::tuliprox_error::{notify_err}; use crate::repository::storage_const; -use crate::utils::file_utils; +use crate::utils; #[inline] pub fn hash_bytes(bytes: &[u8]) -> UUIDType { @@ -68,7 +68,7 @@ pub fn ensure_target_storage_path(cfg: &Config, target_name: &str) -> Result Option { - file_utils::get_file_path(&cfg.working_dir, Some(std::path::PathBuf::from(target_name.replace(' ', "_")))) + utils::get_file_path(&cfg.working_dir, Some(std::path::PathBuf::from(target_name.replace(' ', "_")))) } pub fn get_input_storage_path(input_name: &str, working_dir: &str) -> std::io::Result { diff --git a/src/repository/user_repository.rs b/src/repository/user_repository.rs index a4324f9bf..17e450327 100644 --- a/src/repository/user_repository.rs +++ b/src/repository/user_repository.rs @@ -6,14 +6,13 @@ use crate::model::PlaylistXtreamCategory; use crate::repository::bplustree::BPlusTree; use crate::repository::storage_const; use crate::repository::xtream_repository::xtream_get_playlist_categories; -use crate::utils::file_utils; use crate::utils::json_write_documents_to_file; use chrono::Local; use log::error; use std::collections::{HashMap, HashSet}; use std::io::Error; use std::path::{Path, PathBuf}; - +use crate::utils; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] struct StoredProxyUserCredentialsDeprecated { @@ -209,7 +208,7 @@ pub fn load_api_user(cfg: &Config) -> Result, Error> { } pub fn get_user_storage_path(cfg: &Config, username: &str) -> Option { - cfg.user_config_dir.as_ref().and_then(|ucd| file_utils::get_file_path(ucd, Some(std::path::PathBuf::from(username)))) + cfg.user_config_dir.as_ref().and_then(|ucd| utils::get_file_path(ucd, Some(std::path::PathBuf::from(username)))) } fn ensure_user_storage_path(cfg: &Config, username: &str) -> Option { diff --git a/src/repository/xtream_playlist_iterator.rs b/src/repository/xtream_playlist_iterator.rs index 4c1dbe2a3..5f4d97b1f 100644 --- a/src/repository/xtream_playlist_iterator.rs +++ b/src/repository/xtream_playlist_iterator.rs @@ -9,7 +9,7 @@ use crate::model::XtreamMappingOptions; use crate::repository::indexed_document::{IndexedDocumentIterator}; use crate::repository::user_repository::user_get_bouquet_filter; use crate::repository::xtream_repository::{xtream_get_file_paths, xtream_get_storage_path}; -use crate::utils::file_lock_manager::FileReadGuard; +use crate::utils::FileReadGuard; pub struct XtreamPlaylistIterator { reader: IndexedDocumentIterator, diff --git a/src/repository/xtream_repository.rs b/src/repository/xtream_repository.rs index 410fa8b78..d23d0c72f 100644 --- a/src/repository/xtream_repository.rs +++ b/src/repository/xtream_repository.rs @@ -12,9 +12,9 @@ use crate::repository::storage_const; use crate::repository::target_id_mapping::VirtualIdRecord; use crate::repository::xtream_playlist_iterator::XtreamPlaylistJsonIterator; use crate::utils::bincode_deserialize; -use crate::utils::file_lock_manager::FileReadGuard; -use crate::utils::file_utils::file_reader; -use crate::utils::file_utils::open_readonly_file; +use crate::utils::FileReadGuard; +use crate::utils::file_reader; +use crate::utils::open_readonly_file; use crate::utils::generate_playlist_uuid; use crate::utils::{get_u32_from_serde_value, json_iter_array, json_write_documents_to_file}; use bytes::Bytes; diff --git a/src/tools/lru_cache.rs b/src/tools/lru_cache.rs index 34d1eaa2f..f22f17fea 100644 --- a/src/tools/lru_cache.rs +++ b/src/tools/lru_cache.rs @@ -1,5 +1,5 @@ use crate::repository::storage::hash_string_as_hex; -use crate::utils::file_utils::traverse_dir; +use crate::utils::traverse_dir; use crate::utils::human_readable_byte_size; use log::{debug, error, info, trace}; use std::collections::{HashMap, VecDeque}; diff --git a/src/utils/compression/compressed_file_reader.rs b/src/utils/compression/compressed_file_reader.rs index 50335a4a2..5326c7613 100644 --- a/src/utils/compression/compressed_file_reader.rs +++ b/src/utils/compression/compressed_file_reader.rs @@ -2,7 +2,7 @@ use std::io::{BufRead, BufReader, Read, Seek, SeekFrom}; use std::path::Path; use flate2::bufread::{GzDecoder, ZlibDecoder}; use crate::utils::compression::compression_utils::{is_deflate, is_gzip}; -use crate::utils::file_utils::{file_reader, open_readonly_file}; +use crate::utils::{file_reader, open_readonly_file}; pub struct CompressedFileReader { reader: BufReader>, diff --git a/src/utils/file/config_reader.rs b/src/utils/file/config_reader.rs index cb8836f27..49c5d64e3 100644 --- a/src/utils/file/config_reader.rs +++ b/src/utils/file/config_reader.rs @@ -1,22 +1,18 @@ use crate::model::ApiProxyConfig; -use crate::model::Mappings; -use crate::model::{Config, ConfigDto, ConfigInput, ConfigInputAlias, InputType}; -use crate::tuliprox_error::{create_tuliprox_error, create_tuliprox_error_result, handle_tuliprox_error_result, info_err, str_to_io_error, to_io_error, TuliproxError, TuliproxErrorKind}; -use crate::utils::env_resolving_reader::EnvResolvingReader; -use crate::utils::file_utils::{file_reader, resolve_relative_path}; -use crate::utils::request::{get_credentials_from_url, get_local_file_content}; +use crate::model::{Config, ConfigDto}; +use crate::tuliprox_error::{create_tuliprox_error, create_tuliprox_error_result, to_io_error, TuliproxError, TuliproxErrorKind}; +use crate::utils::{open_file, EnvResolvingReader}; +use crate::utils::{MultiFileReader, file_reader}; use crate::utils::sys_utils::exit; use crate::utils::CONSTANTS; -use crate::utils::{file_utils, multi_file_reader}; use chrono::Local; use log::{error, info, warn}; use serde::Serialize; use std::env; use std::fs::File; -use std::io::{self, BufRead, BufReader, Cursor, Error, Read}; -use std::path::PathBuf; +use std::io::{self, BufReader, Read}; +use std::path::{PathBuf}; use std::sync::Arc; -use url::Url; enum EitherReader { Left(L), @@ -42,18 +38,6 @@ pub fn config_file_reader(file: File, resolve_env: bool) -> impl Read } } -pub fn read_mappings(mappings_file: &str, resolve_env: bool) -> Result, TuliproxError> { - match read_mapping(mappings_file, resolve_env) { - Ok(mappings) => { - match mappings { - None => Ok(None), - Some(mappings_cfg) => Ok(Some(mappings_cfg)) - } - } - Err(err) => Err(err), - } -} - pub fn read_api_proxy_config(cfg: &Config) -> Result<(), TuliproxError> { let api_proxy_config = read_api_proxy(cfg, true); match api_proxy_config { @@ -70,7 +54,7 @@ pub fn read_api_proxy_config(cfg: &Config) -> Result<(), TuliproxError> { pub fn read_config(config_path: &str, config_file: &str, sources_file: &str, api_proxy_file: &str, mappings_file: &str, include_computed: bool) -> Result { let files = vec![std::path::PathBuf::from(config_file), std::path::PathBuf::from(sources_file)]; - match multi_file_reader::MultiFileReader::new(&files) { + match MultiFileReader::new(&files) { Ok(reader) => { match serde_yaml::from_reader::<_, Config>(reader) { Ok(mut result) => { @@ -93,27 +77,9 @@ pub fn read_config(config_path: &str, config_file: &str, sources_file: &str, api } } -pub fn read_mapping(mapping_file: &str, resolve_var: bool) -> Result, TuliproxError> { - let mapping_file = std::path::PathBuf::from(mapping_file); - if let Ok(file) = file_utils::open_file(&mapping_file) { - let maybe_mapping: Result = serde_yaml::from_reader(config_file_reader(file, resolve_var)); - return match maybe_mapping { - Ok(mut mapping) => { - handle_tuliprox_error_result!(TuliproxErrorKind::Info, mapping.prepare()); - Ok(Some(mapping)) - } - Err(err) => { - Err(info_err!(err.to_string())) - } - }; - } - warn!("cant read mapping file: {}", mapping_file.to_str().unwrap_or("?")); - Ok(None) -} - pub fn read_api_proxy(config: &Config, resolve_env: bool) -> Option { let api_proxy_file = config.t_api_proxy_file_path.as_str(); - file_utils::open_file(&std::path::PathBuf::from(api_proxy_file)).map_or(None, |file| { + open_file(&std::path::PathBuf::from(api_proxy_file)).map_or(None, |file| { let maybe_api_proxy: Result = serde_yaml::from_reader(config_file_reader(file, resolve_env)); match maybe_api_proxy { Ok(mut api_proxy) => { @@ -179,242 +145,9 @@ pub fn resolve_env_var(value: &str) -> String { }).to_string() } -const CSV_SEPARATOR: char = ';'; -const HEADER_PREFIX: char = '#'; -const FIELD_MAX_CON: &str = "max_connections"; -const FIELD_PRIO: &str = "priority"; -const FIELD_URL: &str = "url"; -const FIELD_NAME: &str = "name"; -const FIELD_USERNAME: &str = "username"; -const FIELD_PASSWORD: &str = "password"; -const FIELD_UNKNOWN: &str = "?"; -const DEFAULT_COLUMNS: &[&str] = &[FIELD_URL, FIELD_MAX_CON, FIELD_PRIO, FIELD_NAME, FIELD_USERNAME, FIELD_PASSWORD]; - -fn csv_assign_mandatory_fields(alias: &mut ConfigInputAlias, input_type: InputType) { - if !alias.url.is_empty() { - match Url::parse(alias.url.as_str()) { - Ok(url) => { - let (username, password) = get_credentials_from_url(&url); - if username.is_none() || password.is_none() { - // xtream url - if input_type == InputType::XtreamBatch { - alias.url = url.origin().ascii_serialization().to_string(); - } else if input_type == InputType::M3uBatch && alias.username.is_some() && alias.password.is_some() { - alias.url = format!("{}/get_php?username={}&password={}&type=m3u_plus", - url.origin().ascii_serialization(), - alias.username.as_deref().unwrap_or(""), - alias.password.as_deref().unwrap_or("") - ); - } - } else { - if input_type == InputType::XtreamBatch { - alias.url = url.origin().ascii_serialization().to_string(); - } - // m3u url - alias.username = username; - alias.password = password; - } - - if alias.name.is_empty() { - let username = alias.username.as_deref().unwrap_or_default(); - let domain: Vec<&str> = url.domain().unwrap_or_default().split('.').collect(); - if domain.len() > 1 { - alias.name = format!("{}_{username}", domain[domain.len() - 2]); - } else { - alias.name = username.to_string(); - } - } - } - Err(_err) => {} - } - } -} - -fn csv_assign_config_input_column(config_input: &mut ConfigInputAlias, header: &str, raw_value: &str) -> Result<(), io::Error> { - let value = raw_value.trim(); - if !value.is_empty() { - match header { - FIELD_URL => { - let url = Url::parse(value.trim()).map_err(to_io_error)?; - config_input.url = url.to_string(); - } - FIELD_MAX_CON => { - let max_connections = value.parse::().unwrap_or(1); - config_input.max_connections = max_connections; - } - FIELD_PRIO => { - let priority = value.parse::().unwrap_or(0); - config_input.priority = priority; - } - FIELD_NAME => { - config_input.name = value.to_string(); - } - FIELD_USERNAME => { - config_input.username = Some(value.to_string()); - } - FIELD_PASSWORD => { - config_input.password = Some(value.to_string()); - } - _ => {} - } - } - Ok(()) -} - -pub fn csv_read_inputs_from_reader(batch_input_type: InputType, reader: impl BufRead) -> Result, io::Error> { - let input_type = match batch_input_type { - InputType::M3uBatch | InputType::M3u => InputType::M3uBatch, - InputType::XtreamBatch | InputType::Xtream => InputType::XtreamBatch - }; - let mut result = vec![]; - let mut default_columns = vec![]; - default_columns.extend_from_slice(DEFAULT_COLUMNS); - let mut header_defined = false; - for line in reader.lines() { - let line = line?; - if line.is_empty() { - continue; - } - if line.starts_with(HEADER_PREFIX) { - if !header_defined { - header_defined = true; - default_columns = line[1..].split(CSV_SEPARATOR).map(|s| { - match s { - FIELD_URL => FIELD_URL, - FIELD_MAX_CON => FIELD_MAX_CON, - FIELD_PRIO => FIELD_PRIO, - FIELD_NAME => FIELD_NAME, - FIELD_USERNAME => FIELD_USERNAME, - FIELD_PASSWORD => FIELD_PASSWORD, - _ => { - error!("Field {s} is unsupported for csv input"); - FIELD_UNKNOWN - } - } - }).collect(); - } - continue; - } - - let mut config_input = ConfigInputAlias { - id: 0, - name: String::new(), - url: String::new(), - username: None, - password: None, - priority: 0, - max_connections: 1, - t_base_url: String::default(), - }; - - let columns: Vec<&str> = line.split(CSV_SEPARATOR).collect(); - for (&header, &value) in default_columns.iter().zip(columns.iter()) { - if let Err(err) = csv_assign_config_input_column(&mut config_input, header, value) { - error!("Could not parse input line: {line} err: {err}"); - } - } - csv_assign_mandatory_fields(&mut config_input, input_type); - result.push(config_input); - } - Ok(result) -} - - -pub fn csv_read_inputs(input: &ConfigInput) -> Result, io::Error> { - let file_uri = input.url.to_string(); - let file_path = get_csv_file_path(&file_uri)?; - match get_local_file_content(&file_path) { - Ok(content) => { - csv_read_inputs_from_reader(input.input_type, EnvResolvingReader::new(file_reader(Cursor::new(content)))) - } - Err(err) => { - Err(err) - } - } -} - -fn get_csv_file_path(file_uri: &String) -> Result { - if file_uri.contains("://") { - if let Ok(url) = file_uri.parse::() { - if url.scheme() == "file" { - return match url.to_file_path() { - Ok(path) => Ok(path), - Err(()) => Err(str_to_io_error(&format!("Could not open {file_uri}"))), - }; - } - } - Err(str_to_io_error(&format!("Only file:// is supported {file_uri}"))) - } else { - resolve_relative_path(file_uri) - } -} - #[cfg(test)] mod tests { - use crate::model::InputType; - use crate::utils::config_reader::{csv_read_inputs_from_reader, resolve_env_var}; - use std::io::{BufReader, Cursor}; - const M3U_BATCH: &str = r#" -#url;name;max_connections;priority -http://hd.providerline.com:8080/get.php?username=user1&password=user1&type=m3u_plus;input_1 -http://hd.providerline.com/get.php?username=user2&password=user2&type=m3u_plus;input_2;1;2 -http://hd.providerline.com/get.php?username=user3&password=user3&type=m3u_plus;input_3;1;2 -http://hd.providerline.com/get.php?username=user4&password=user4&type=m3u_plus;input_4 -"#; - - const XTREAM_BATCH: &str = r#" -#name;username;password;url;max_connections -input_1;desanocra;eyCG8SN523KQ;http://provider_1.tv:80;1 -input_2;desanocra;eyCG8SN523KQ;http://provider_2.tv:8080;1 -"#; - - #[test] - fn test_read_inputs_xtream_as_m3u() { - let reader = BufReader::new(Cursor::new(XTREAM_BATCH)); - let result = csv_read_inputs_from_reader(InputType::M3uBatch, reader); - assert_eq!(result.is_ok(), true); - let aliases = result.unwrap(); - assert_eq!(aliases.is_empty(), false); - for config in aliases { - assert_eq!(config.url.contains("username"), true); - } - } - - #[test] - fn test_read_inputs_m3u_as_m3u() { - let reader = BufReader::new(Cursor::new(M3U_BATCH)); - let result = csv_read_inputs_from_reader(InputType::M3uBatch, reader); - assert_eq!(result.is_ok(), true); - let aliases = result.unwrap(); - assert_eq!(aliases.is_empty(), false); - for config in aliases { - assert_eq!(config.url.contains("username"), true); - } - } - - #[test] - fn test_read_inputs_xtream_as_xtream() { - let reader = BufReader::new(Cursor::new(XTREAM_BATCH)); - let result = csv_read_inputs_from_reader(InputType::XtreamBatch, reader); - assert_eq!(result.is_ok(), true); - let aliases = result.unwrap(); - assert_eq!(aliases.is_empty(), false); - for config in aliases { - assert_eq!(config.url.contains("username"), false); - } - } - - #[test] - fn test_read_inputs_m3u_as_xtream() { - let reader = BufReader::new(Cursor::new(M3U_BATCH)); - let result = csv_read_inputs_from_reader(InputType::XtreamBatch, reader); - assert_eq!(result.is_ok(), true); - let aliases = result.unwrap(); - assert_eq!(aliases.is_empty(), false); - for config in aliases { - assert_eq!(config.url.contains("username"), false); - } - } + use crate::utils::resolve_env_var; #[test] fn test_resolve() { diff --git a/src/utils/file/csv_input_reader.rs b/src/utils/file/csv_input_reader.rs new file mode 100644 index 000000000..18b91ce1c --- /dev/null +++ b/src/utils/file/csv_input_reader.rs @@ -0,0 +1,256 @@ +use crate::model::{ConfigInput, ConfigInputAlias, InputType}; +use crate::tuliprox_error::{str_to_io_error, to_io_error}; +use crate::utils::EnvResolvingReader; +use crate::utils::{file_reader, resolve_relative_path}; +use crate::utils::request::{get_credentials_from_url, get_local_file_content}; +use log::error; +use std::io; +use std::io::{BufRead, Cursor, Error}; +use std::path::PathBuf; +use url::Url; + +const CSV_SEPARATOR: char = ';'; +const HEADER_PREFIX: char = '#'; +const FIELD_MAX_CON: &str = "max_connections"; +const FIELD_PRIO: &str = "priority"; +const FIELD_URL: &str = "url"; +const FIELD_NAME: &str = "name"; +const FIELD_USERNAME: &str = "username"; +const FIELD_PASSWORD: &str = "password"; +const FIELD_UNKNOWN: &str = "?"; +const DEFAULT_COLUMNS: &[&str] = &[FIELD_URL, FIELD_MAX_CON, FIELD_PRIO, FIELD_NAME, FIELD_USERNAME, FIELD_PASSWORD]; + +fn csv_assign_mandatory_fields(alias: &mut ConfigInputAlias, input_type: InputType) { + if !alias.url.is_empty() { + match Url::parse(alias.url.as_str()) { + Ok(url) => { + let (username, password) = get_credentials_from_url(&url); + if username.is_none() || password.is_none() { + // xtream url + if input_type == InputType::XtreamBatch { + alias.url = url.origin().ascii_serialization().to_string(); + } else if input_type == InputType::M3uBatch && alias.username.is_some() && alias.password.is_some() { + alias.url = format!("{}/get_php?username={}&password={}&type=m3u_plus", + url.origin().ascii_serialization(), + alias.username.as_deref().unwrap_or(""), + alias.password.as_deref().unwrap_or("") + ); + } + } else { + if input_type == InputType::XtreamBatch { + alias.url = url.origin().ascii_serialization().to_string(); + } + // m3u url + alias.username = username; + alias.password = password; + } + + if alias.name.is_empty() { + let username = alias.username.as_deref().unwrap_or_default(); + let domain: Vec<&str> = url.domain().unwrap_or_default().split('.').collect(); + if domain.len() > 1 { + alias.name = format!("{}_{username}", domain[domain.len() - 2]); + } else { + alias.name = username.to_string(); + } + } + } + Err(_err) => {} + } + } +} + +fn csv_assign_config_input_column(config_input: &mut ConfigInputAlias, header: &str, raw_value: &str) -> Result<(), io::Error> { + let value = raw_value.trim(); + if !value.is_empty() { + match header { + FIELD_URL => { + let url = Url::parse(value.trim()).map_err(to_io_error)?; + config_input.url = url.to_string(); + } + FIELD_MAX_CON => { + let max_connections = value.parse::().unwrap_or(1); + config_input.max_connections = max_connections; + } + FIELD_PRIO => { + let priority = value.parse::().unwrap_or(0); + config_input.priority = priority; + } + FIELD_NAME => { + config_input.name = value.to_string(); + } + FIELD_USERNAME => { + config_input.username = Some(value.to_string()); + } + FIELD_PASSWORD => { + config_input.password = Some(value.to_string()); + } + _ => {} + } + } + Ok(()) +} + +pub fn csv_read_inputs_from_reader(batch_input_type: InputType, reader: impl BufRead) -> Result, io::Error> { + let input_type = match batch_input_type { + InputType::M3uBatch | InputType::M3u => InputType::M3uBatch, + InputType::XtreamBatch | InputType::Xtream => InputType::XtreamBatch + }; + let mut result = vec![]; + let mut default_columns = vec![]; + default_columns.extend_from_slice(DEFAULT_COLUMNS); + let mut header_defined = false; + for line in reader.lines() { + let line = line?; + if line.is_empty() { + continue; + } + if line.starts_with(HEADER_PREFIX) { + if !header_defined { + header_defined = true; + default_columns = line[1..].split(CSV_SEPARATOR).map(|s| { + match s { + FIELD_URL => FIELD_URL, + FIELD_MAX_CON => FIELD_MAX_CON, + FIELD_PRIO => FIELD_PRIO, + FIELD_NAME => FIELD_NAME, + FIELD_USERNAME => FIELD_USERNAME, + FIELD_PASSWORD => FIELD_PASSWORD, + _ => { + error!("Field {s} is unsupported for csv input"); + FIELD_UNKNOWN + } + } + }).collect(); + } + continue; + } + + let mut config_input = ConfigInputAlias { + id: 0, + name: String::new(), + url: String::new(), + username: None, + password: None, + priority: 0, + max_connections: 1, + t_base_url: String::default(), + }; + + let columns: Vec<&str> = line.split(CSV_SEPARATOR).collect(); + for (&header, &value) in default_columns.iter().zip(columns.iter()) { + if let Err(err) = csv_assign_config_input_column(&mut config_input, header, value) { + error!("Could not parse input line: {line} err: {err}"); + } + } + csv_assign_mandatory_fields(&mut config_input, input_type); + result.push(config_input); + } + Ok(result) +} + + +pub fn csv_read_inputs(input: &ConfigInput) -> Result, io::Error> { + let file_uri = input.url.to_string(); + let file_path = get_csv_file_path(&file_uri)?; + match get_local_file_content(&file_path) { + Ok(content) => { + csv_read_inputs_from_reader(input.input_type, EnvResolvingReader::new(file_reader(Cursor::new(content)))) + } + Err(err) => { + Err(err) + } + } +} + +fn get_csv_file_path(file_uri: &String) -> Result { + if file_uri.contains("://") { + if let Ok(url) = file_uri.parse::() { + if url.scheme() == "file" { + return match url.to_file_path() { + Ok(path) => Ok(path), + Err(()) => Err(str_to_io_error(&format!("Could not open {file_uri}"))), + }; + } + } + Err(str_to_io_error(&format!("Only file:// is supported {file_uri}"))) + } else { + resolve_relative_path(file_uri) + } +} + +#[cfg(test)] +mod tests { + use crate::model::InputType; + use crate::utils::config_reader::resolve_env_var; + use crate::utils::file::csv_input_reader::csv_read_inputs_from_reader; + use std::io::{BufReader, Cursor}; + + const M3U_BATCH: &str = r#" +#url;name;max_connections;priority +http://hd.providerline.com:8080/get.php?username=user1&password=user1&type=m3u_plus;input_1 +http://hd.providerline.com/get.php?username=user2&password=user2&type=m3u_plus;input_2;1;2 +http://hd.providerline.com/get.php?username=user3&password=user3&type=m3u_plus;input_3;1;2 +http://hd.providerline.com/get.php?username=user4&password=user4&type=m3u_plus;input_4 +"#; + + const XTREAM_BATCH: &str = r#" +#name;username;password;url;max_connections +input_1;desanocra;eyCG8SN523KQ;http://provider_1.tv:80;1 +input_2;desanocra;eyCG8SN523KQ;http://provider_2.tv:8080;1 +"#; + + #[test] + fn test_read_inputs_xtream_as_m3u() { + let reader = BufReader::new(Cursor::new(XTREAM_BATCH)); + let result = csv_read_inputs_from_reader(InputType::M3uBatch, reader); + assert_eq!(result.is_ok(), true); + let aliases = result.unwrap(); + assert_eq!(aliases.is_empty(), false); + for config in aliases { + assert_eq!(config.url.contains("username"), true); + } + } + + #[test] + fn test_read_inputs_m3u_as_m3u() { + let reader = BufReader::new(Cursor::new(M3U_BATCH)); + let result = csv_read_inputs_from_reader(InputType::M3uBatch, reader); + assert_eq!(result.is_ok(), true); + let aliases = result.unwrap(); + assert_eq!(aliases.is_empty(), false); + for config in aliases { + assert_eq!(config.url.contains("username"), true); + } + } + + #[test] + fn test_read_inputs_xtream_as_xtream() { + let reader = BufReader::new(Cursor::new(XTREAM_BATCH)); + let result = csv_read_inputs_from_reader(InputType::XtreamBatch, reader); + assert_eq!(result.is_ok(), true); + let aliases = result.unwrap(); + assert_eq!(aliases.is_empty(), false); + for config in aliases { + assert_eq!(config.url.contains("username"), false); + } + } + + #[test] + fn test_read_inputs_m3u_as_xtream() { + let reader = BufReader::new(Cursor::new(M3U_BATCH)); + let result = csv_read_inputs_from_reader(InputType::XtreamBatch, reader); + assert_eq!(result.is_ok(), true); + let aliases = result.unwrap(); + assert_eq!(aliases.is_empty(), false); + for config in aliases { + assert_eq!(config.url.contains("username"), false); + } + } + + #[test] + fn test_resolve() { + let resolved = resolve_env_var("${env:HOME}"); + assert_eq!(resolved, std::env::var("HOME").unwrap()); + } +} \ No newline at end of file diff --git a/src/utils/file/env_resolving_reader.rs b/src/utils/file/env_resolving_reader.rs index 7345646eb..daa7b1781 100644 --- a/src/utils/file/env_resolving_reader.rs +++ b/src/utils/file/env_resolving_reader.rs @@ -1,5 +1,5 @@ use std::io::{self, BufRead, BufReader, Read, Cursor}; -use crate::utils::config_reader::resolve_env_var; +use crate::utils::resolve_env_var; pub struct EnvResolvingReader { inner: BufReader, diff --git a/src/utils/file/file_utils.rs b/src/utils/file/file_utils.rs index 2fb60d280..321601dac 100644 --- a/src/utils/file/file_utils.rs +++ b/src/utils/file/file_utils.rs @@ -293,3 +293,7 @@ pub fn resolve_relative_path(relative: &str) -> std::io::Result { let current_dir = env::current_dir()?; Ok(current_dir.join(relative)) } + +pub fn is_directory(path: &str) -> bool { + PathBuf::from(path).is_dir() +} \ No newline at end of file diff --git a/src/utils/file/mapping_reader.rs b/src/utils/file/mapping_reader.rs new file mode 100644 index 000000000..59666f1f5 --- /dev/null +++ b/src/utils/file/mapping_reader.rs @@ -0,0 +1,138 @@ +use std::collections::HashMap; +use crate::foundation::filter::PatternTemplate; +use crate::model::{Mapping, MappingDefinition, MappingTag, Mappings}; +use crate::tuliprox_error::{create_tuliprox_error_result, handle_tuliprox_error_result, info_err, TuliproxError, TuliproxErrorKind}; +use crate::utils::traverse_dir; +use crate::utils::{config_file_reader, open_file}; +use log::{warn}; +use std::path::{Path, PathBuf}; + +fn read_mapping(mapping_file: &Path, resolve_var: bool, prepare_mappings: bool) -> Result, TuliproxError> { + if let Ok(file) = open_file(mapping_file) { + let maybe_mapping: Result = serde_yaml::from_reader(config_file_reader(file, resolve_var)); + return match maybe_mapping { + Ok(mut mapping) => { + if prepare_mappings { + handle_tuliprox_error_result!(TuliproxErrorKind::Info, mapping.prepare()); + } + Ok(Some(mapping)) + } + Err(err) => { + Err(info_err!(err.to_string())) + } + }; + } + warn!("cant read mapping file: {}", mapping_file.to_str().unwrap_or("?")); + Ok(None) +} + +fn read_mappings_from_file(mappings_file: &Path, resolve_env: bool) -> Result, TuliproxError> { + match read_mapping(mappings_file, resolve_env, true) { + Ok(mappings) => { + match mappings { + None => Ok(None), + Some(mappings_cfg) => Ok(Some(mappings_cfg)) + } + } + Err(err) => Err(err), + } +} + + +fn merge_mappings(mappings: Vec) -> Vec { + let mut map: HashMap = HashMap::new(); + + for mut m in mappings { + let entry = map.entry(m.id.clone()).or_insert_with(|| Mapping { + id: m.id.clone(), + ..Default::default() + }); + + // Logic for match_as_ascii: true, if one of them is true + entry.match_as_ascii |= m.match_as_ascii; + + if let Some(mut mapper) = m.mapper.take() { + entry.mapper.get_or_insert(vec![]).append(&mut mapper); + } + + if let Some(mut counters) = m.counter.take() { + entry.counter.get_or_insert(vec![]).append(&mut counters); + } + } + + map.into_values().collect() +} +fn merge_mapping_definitions(mappings: Vec) -> Result, TuliproxError> { + let mut merged_templates: Vec = Vec::new(); + let mut merged_tags: Vec = Vec::new(); + let mut merged_mapping: Vec = Vec::new(); + + for mapping in mappings { + if let Some(mut templates) = mapping.mappings.templates { + merged_templates.append(&mut templates); + } + + if let Some(mut tags) = mapping.mappings.tags { + merged_tags.append(&mut tags); + } + + merged_mapping.extend(mapping.mappings.mapping); + } + + let mut result = Mappings { + mappings: MappingDefinition { + templates: if merged_templates.is_empty() { None } else { Some(merged_templates) }, + tags: if merged_tags.is_empty() { None } else { Some(merged_tags) }, + mapping: merge_mappings(merged_mapping) + } + }; + result.prepare()?; + Ok(Some(result)) +} + +fn read_mappings_from_directory(path: &Path, resolve_env: bool) -> Result, TuliproxError> { + let mut files = vec![]; + let mut visit = |entry: &std::fs::DirEntry, metadata: &std::fs::Metadata| { + if metadata.is_file() { + let file_path = entry.path(); + if file_path.extension().is_some_and(|ext| ext == "yml") { + files.push(file_path); + } + } + }; + traverse_dir(path, &mut visit).map_err(|err| TuliproxError::new(TuliproxErrorKind::Info, format!("Failed to read mappings {err}")))?; + + files.sort_by(|a,b| a.file_name().cmp(&b.file_name())); + + let mut mappings = vec![]; + for file_path in files { + match read_mapping(&file_path, resolve_env, false) { + Ok(Some(mapping)) => mappings.push(mapping), + Ok(None) => {} + Err(err) => return create_tuliprox_error_result!(TuliproxErrorKind::Info, "Failed to read mapping file {file_path:?}: {err:?}"), + } + } + + if mappings.is_empty() { + return Ok(None); + } + merge_mapping_definitions(mappings) +} + +pub fn read_mappings(mappings_file: &str, resolve_env: bool) -> Result, TuliproxError> { + let path = PathBuf::from(mappings_file); + match std::fs::metadata(&path) { + Ok(metadata) => { + if metadata.is_file() { + read_mappings_from_file(&path, resolve_env) + } else if metadata.is_dir() { + read_mappings_from_directory(&path, resolve_env) + } else { + Err(TuliproxError::new(TuliproxErrorKind::Info, format!("cant read mappings file: {mappings_file}"))) + } + } + Err(_err) => { + Err(TuliproxError::new(TuliproxErrorKind::Info, format!("cant read mappings file: {mappings_file}"))) + } + } +} \ No newline at end of file diff --git a/src/utils/file/mod.rs b/src/utils/file/mod.rs index ffdbe06f8..e79c6819b 100644 --- a/src/utils/file/mod.rs +++ b/src/utils/file/mod.rs @@ -1,5 +1,15 @@ -pub mod file_utils; -pub mod multi_file_reader; -pub mod file_lock_manager; -pub mod config_reader; -pub mod env_resolving_reader; \ No newline at end of file +mod file_utils; +mod multi_file_reader; +mod file_lock_manager; +mod config_reader; +mod env_resolving_reader; +mod mapping_reader; +mod csv_input_reader; + +pub use self::file_utils::*; +pub use self::multi_file_reader::*; +pub use self::file_lock_manager::*; +pub use self::config_reader::*; +pub use self::mapping_reader::*; +pub use self::env_resolving_reader::*; +pub use self::csv_input_reader::*; \ No newline at end of file diff --git a/src/utils/file/multi_file_reader.rs b/src/utils/file/multi_file_reader.rs index a412cfd07..3dd20f851 100644 --- a/src/utils/file/multi_file_reader.rs +++ b/src/utils/file/multi_file_reader.rs @@ -1,8 +1,8 @@ use std::fs::File; use std::io::{self, ErrorKind, Read}; use std::path::{PathBuf}; -use crate::utils::env_resolving_reader::EnvResolvingReader; -use crate::utils::file_utils::file_reader; +use crate::utils::EnvResolvingReader; +use crate::utils::file_reader; pub struct MultiFileReader { files: Vec, diff --git a/src/utils/json_utils.rs b/src/utils/json_utils.rs index 647322f30..ca57cc9a6 100644 --- a/src/utils/json_utils.rs +++ b/src/utils/json_utils.rs @@ -6,7 +6,7 @@ use std::path::Path; use serde::de::DeserializeOwned; use serde::{Deserialize, Serialize}; use serde_json::{self, Deserializer, Value}; -use crate::utils::file_utils::{file_reader, file_writer}; +use crate::utils::{file_reader, file_writer}; fn read_skipping_ws(mut reader: impl Read) -> io::Result { loop { diff --git a/src/utils/logging.rs b/src/utils/logging.rs index d9c9fb837..c474759c9 100644 --- a/src/utils/logging.rs +++ b/src/utils/logging.rs @@ -2,7 +2,7 @@ use std::fs::File; use env_logger::Builder; use log::{error, info, LevelFilter}; use crate::model::LogLevelConfig; -use crate::utils::config_reader::config_file_reader; +use crate::utils::config_file_reader; const LOG_ERROR_LEVEL_MOD: &[&str] = &[ "reqwest::async_impl::client", diff --git a/src/utils/network/epg.rs b/src/utils/network/epg.rs index 10d239c4a..db09f1a48 100644 --- a/src/utils/network/epg.rs +++ b/src/utils/network/epg.rs @@ -2,8 +2,7 @@ use crate::tuliprox_error::TuliproxError; use crate::model::{Config, ConfigInput, PersistedEpgSource}; use crate::model::TVGuide; use crate::repository::storage::short_hash; -use crate::utils::file_utils; -use crate::utils::file_utils::prepare_file_path; +use crate::utils::{add_prefix_to_filename, prepare_file_path}; use crate::utils::request; use log::debug; use std::path::PathBuf; @@ -14,7 +13,7 @@ async fn download_epg_file(url: &str, client: &Arc, input: &Con debug!("Getting epg file path for url: {url}"); let file_prefix = short_hash(url); let persist_file_path = prepare_file_path(input.persist.as_deref(), working_dir, "") - .map(|path| file_utils::add_prefix_to_filename(&path, format!("{file_prefix}_epg_").as_str(), Some("xml"))); + .map(|path| add_prefix_to_filename(&path, format!("{file_prefix}_epg_").as_str(), Some("xml"))); request::get_input_text_content_as_file(Arc::clone(client), input, working_dir, url, persist_file_path).await } diff --git a/src/utils/network/m3u.rs b/src/utils/network/m3u.rs index c439294ee..0bd447c6f 100644 --- a/src/utils/network/m3u.rs +++ b/src/utils/network/m3u.rs @@ -3,7 +3,7 @@ use crate::tuliprox_error::TuliproxError; use crate::model::{Config, ConfigInput}; use crate::model::PlaylistGroup; use crate::processing::parser::m3u; -use crate::utils::file_utils::prepare_file_path; +use crate::utils::prepare_file_path; use crate::utils::request; pub async fn get_m3u_playlist(client: Arc, cfg: &Config, input: &ConfigInput, working_dir: &str) -> (Vec, Vec) { diff --git a/src/utils/network/request.rs b/src/utils/network/request.rs index 923baf448..a9f17d90b 100644 --- a/src/utils/network/request.rs +++ b/src/utils/network/request.rs @@ -22,7 +22,7 @@ use crate::repository::storage::{get_input_storage_path, short_hash}; use crate::repository::storage_const; use crate::utils::compression::compression_utils::{is_deflate, is_gzip}; use crate::utils::{debug_if_enabled, filter_request_header}; -use crate::utils::file_utils::{get_file_path, persist_file}; +use crate::utils::{get_file_path, persist_file}; use crate::utils::{CONSTANTS, DASH_EXT, DASH_EXT_FRAGMENT, DASH_EXT_QUERY, ENCODING_DEFLATE, ENCODING_GZIP, HLS_EXT, HLS_EXT_FRAGMENT, HLS_EXT_QUERY}; pub const fn bytes_to_megabytes(bytes: u64) -> u64 { diff --git a/src/utils/network/xtream.rs b/src/utils/network/xtream.rs index 9ae9d89b6..a6444d81f 100644 --- a/src/utils/network/xtream.rs +++ b/src/utils/network/xtream.rs @@ -151,8 +151,8 @@ pub async fn get_xtream_playlist(client: Arc, input: &ConfigInp if !skip_cluster.contains(xtream_cluster) { let category_url = format!("{base_url}&action={category}"); let stream_url = format!("{base_url}&action={stream}"); - let category_file_path = crate::utils::file::file_utils::prepare_file_path(input.persist.as_deref(), working_dir, format!("{category}_").as_str()); - let stream_file_path = crate::utils::file::file_utils::prepare_file_path(input.persist.as_deref(), working_dir, format!("{stream}_").as_str()); + let category_file_path = crate::utils::prepare_file_path(input.persist.as_deref(), working_dir, format!("{category}_").as_str()); + let stream_file_path = crate::utils::prepare_file_path(input.persist.as_deref(), working_dir, format!("{stream}_").as_str()); match futures::join!( request::get_input_json_content(Arc::clone(&client), input, category_url.as_str(), category_file_path),