From 30650be0adcc2cc399b7fd35b99f8428b2233823 Mon Sep 17 00:00:00 2001 From: euzu Date: Tue, 26 Mar 2024 01:04:21 +0100 Subject: [PATCH] Refactored m3u playlist parsing for playlists with missing group and id information in header. --- Cargo.lock | 87 ++++------- Cargo.toml | 22 +-- src/api/download_api.rs | 6 +- src/api/v1_api.rs | 6 +- src/api/xmltv_api.rs | 8 +- src/api/xtream_api.rs | 14 +- src/config_reader.rs | 11 +- src/main.rs | 40 +++-- src/model/config.rs | 9 +- src/processing/m3u_parser.rs | 43 +++-- src/processing/playlist_processor.rs | 10 +- src/processing/playlist_watch.rs | 4 +- src/repository/m3u_repository.rs | 12 +- src/repository/xtream_repository.rs | 190 ++++++++++++----------- src/{ => utils}/download.rs | 19 ++- src/utils/file_utils.rs | 162 +++++++++++++++++++ src/utils/mod.rs | 4 + src/{utils.rs => utils/request_utils.rs} | 171 +------------------- src/utils/string_utils.rs | 14 ++ 19 files changed, 435 insertions(+), 397 deletions(-) rename src/{ => utils}/download.rs (87%) create mode 100644 src/utils/file_utils.rs create mode 100644 src/utils/mod.rs rename src/{utils.rs => utils/request_utils.rs} (59%) create mode 100644 src/utils/string_utils.rs diff --git a/Cargo.lock b/Cargo.lock index 55000f5d5..40109a567 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -289,9 +289,9 @@ dependencies = [ [[package]] name = "anstream" -version = "0.6.4" +version = "0.6.13" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2ab91ebe16eb252986481c5b62f6098f3b698a45e34b5b98200cf20dd2484a44" +checksum = "d96bd03f33fe50a863e394ee9718a706f988b9079b20c3784fb726e7678b62fb" dependencies = [ "anstyle", "anstyle-parse", @@ -303,9 +303,9 @@ dependencies = [ [[package]] name = "anstyle" -version = "1.0.4" +version = "1.0.6" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7079075b41f533b8c61d2a4d073c4676e1f8b249ff94a393b0595db304e0dd87" +checksum = "8901269c6307e8d93993578286ac0edf7f195079ffff5ebdeea6a59ffb7e36bc" [[package]] name = "anstyle-parse" @@ -737,16 +737,26 @@ dependencies = [ ] [[package]] -name = "env_logger" -version = "0.10.0" +name = "env_filter" +version = "0.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "85cdab6a89accf66733ad5a1693a4dcced6aeff64602b634530dd73c1f3ee9f0" +checksum = "a009aa4810eb158359dda09d0c87378e4bbb89b5a801f016885a4707ba24f7ea" dependencies = [ - "humantime", - "is-terminal", "log", "regex", - "termcolor", +] + +[[package]] +name = "env_logger" +version = "0.11.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "38b35839ba51819680ba087cd351788c9a3c476841207e0b8cee0b04722343b9" +dependencies = [ + "anstream", + "anstyle", + "env_filter", + "humantime", + "log", ] [[package]] @@ -825,9 +835,9 @@ checksum = "00b0228411908ca8685dba7fc2cdd70ec9990a6e753e89b6ac91a84c40fbaf4b" [[package]] name = "form_urlencoded" -version = "1.2.0" +version = "1.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a62bc1cf6f830c2ec14a513a9fb124d0a213a629668a4186f329db21fe045652" +checksum = "e13624c2627564efccf4934284bdd98cbaa14e79b0b5a141218e507b3a823456" dependencies = [ "percent-encoding", ] @@ -1128,9 +1138,9 @@ dependencies = [ [[package]] name = "idna" -version = "0.4.0" +version = "0.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7d20d6b07bfbc108882d88ed8e37d39636dcc260e15e30c45e6ba089610b917c" +checksum = "634d9b1461af396cad843f47fdba5597a4f9e6ddd4bfb6ff5d85028c25cb12f6" dependencies = [ "unicode-bidi", "unicode-normalization", @@ -1171,17 +1181,6 @@ version = "2.9.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8f518f335dce6725a761382244631d86cf0ccb2863413590b31338feb467f9c3" -[[package]] -name = "is-terminal" -version = "0.4.9" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "cb0889898416213fab133e1d33a0e5858a48177452750691bde3666d0fdbaf8b" -dependencies = [ - "hermit-abi", - "rustix", - "windows-sys", -] - [[package]] name = "isahc" version = "1.7.2" @@ -1310,9 +1309,9 @@ dependencies = [ [[package]] name = "log" -version = "0.4.20" +version = "0.4.21" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b5e6163cb8c49088c2c36f57875e58ccd8c87c7427f7fbd50ea6710b2f3f2e8f" +checksum = "90ed8c1e510134f979dbc4f070f87d4313098b704861a105fe34231c70a3901c" [[package]] name = "lzma-rs" @@ -1581,9 +1580,9 @@ dependencies = [ [[package]] name = "percent-encoding" -version = "2.3.0" +version = "2.3.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9b2a4787296e9989611394c33f193f676704af1686e70b8f8033ab5ba9a35a94" +checksum = "e3148f5046208a5d56bcfc03053e3ca6334e51da8dfb19b6cdc8b306fae3283e" [[package]] name = "pest" @@ -2215,15 +2214,6 @@ dependencies = [ "windows-sys", ] -[[package]] -name = "termcolor" -version = "1.3.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6093bad37da69aab9d123a8091e4be0aa4a03e4d601ec641c327398315f62b64" -dependencies = [ - "winapi-util", -] - [[package]] name = "thiserror" version = "1.0.50" @@ -2456,9 +2446,9 @@ checksum = "8ecb6da28b8a351d773b68d5825ac39017e680750f980f3a1a85cd8dd28a47c1" [[package]] name = "url" -version = "2.4.1" +version = "2.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "143b538f18257fac9cad154828a57c6bf5157e1aa604d4816b5995bf6de87ae5" +checksum = "31e6302e3bb753d46e83516cae55ae196fc0c309407cf11ab35cc51a4c2a4633" dependencies = [ "form_urlencoded", "idna", @@ -2473,9 +2463,9 @@ checksum = "711b9620af191e0cdc7468a8d14e709c3dcdb115b36f838e601583af800a370a" [[package]] name = "uuid" -version = "1.5.0" +version = "1.8.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "88ad59a7560b41a70d191093a945f0b87bc1deeda46fb237479708a1d6b6cdfc" +checksum = "a183cf7feeba97b4dd1c0d46788634f6221d87fa961b305bed08c851829efcc0" dependencies = [ "getrandom", "rand", @@ -2484,9 +2474,9 @@ dependencies = [ [[package]] name = "uuid-macro-internal" -version = "1.5.0" +version = "1.8.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3d8c6bba9b149ee82950daefc9623b32bb1dacbfb1890e352f6b887bd582adaf" +checksum = "9881bea7cbe687e36c9ab3b778c36cd0487402e270304e8b1296d5085303c1a2" dependencies = [ "proc-macro2", "quote", @@ -2637,15 +2627,6 @@ version = "0.4.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ac3b87c63620426dd9b991e5ce0329eff545bccbbb34f3be09ff6fb6ab51b7b6" -[[package]] -name = "winapi-util" -version = "0.1.6" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f29e6f9198ba0d26b4c9f07dbe6f9ed633e1f3d5b8b414090084349e46a52596" -dependencies = [ - "winapi", -] - [[package]] name = "winapi-x86_64-pc-windows-gnu" version = "0.4.0" diff --git a/Cargo.toml b/Cargo.toml index b774c2b1e..8a6ea671a 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -17,29 +17,29 @@ serde_yaml = "0" serde_json = "1" serde-xml-rs = "0" quick-xml = { version = "0", features = ["serialize"] } -regex = "1.7" +regex = "1.10" clap = { version = "4", features = ["derive"] } -url = "2.3" +url = "2.5" reqwest = { version = "0", features = ["blocking", "json", "stream", "rustls-tls"] } chrono = "0.4" cron = "0" -actix-web = "4.3" -actix-server = "2.2" +actix-web = "4.4" +actix-server = "2.3" actix-files = "0" actix-cors = "0" -actix-rt = "2.8" +actix-rt = "2.9" futures = "0.3" -path-absolutize = "3.0" -pest = "2.5" -pest_derive = "2.5" +path-absolutize = "3.1" +pest = "2.7" +pest_derive = "2.7" enum-iterator = "1.3" unidecode = "0" petgraph = "0" openssl = { version = "*", features = ["vendored"] } #https://docs.rs/openssl/0.10.34/openssl/#vendored mime = "0.3" log = "0.4" -env_logger = "0.10" +env_logger = "0.11" rustelebot = "0.3" bincode = "1.3" -uuid = { version = "1.3.0", features = ["v4", "fast-rng", "macro-diagnostics"] } -lzma-rs = "0.3.0" +uuid = { version = "1.7", features = ["v4", "fast-rng", "macro-diagnostics"] } +lzma-rs = "0.3" \ No newline at end of file diff --git a/src/api/download_api.rs b/src/api/download_api.rs index 63b077851..0298434c5 100644 --- a/src/api/download_api.rs +++ b/src/api/download_api.rs @@ -7,9 +7,9 @@ use actix_web::{HttpResponse, web}; use serde_json::{json, Value}; use crate::api::api_model::{AppState, DownloadQueue, FileDownload, FileDownloadRequest}; use crate::model::config::{VideoDownloadConfig}; -use crate::utils::{bytes_to_megabytes, get_request_headers}; use futures::stream::TryStreamExt; use log::{info}; +use crate::utils::{request_utils}; async fn download_file(active: Arc>>, client: &reqwest::Client) -> Result<(), String> { let file_download = { active.read().unwrap().as_ref().unwrap().clone() }; @@ -37,7 +37,7 @@ async fn download_file(active: Arc>>, client: &reqwe } } None => { - let megabytes = bytes_to_megabytes(downloaded); + let megabytes = request_utils::bytes_to_megabytes(downloaded); info!("Downloaded {}, filesize: {}MB", file_path_str, megabytes); active.write().unwrap().as_mut().unwrap().size = downloaded; return Ok(()); @@ -67,7 +67,7 @@ fn run_download_queue(download_cfg: &VideoDownloadConfig, download_queue: Arc { diff --git a/src/api/v1_api.rs b/src/api/v1_api.rs index 70085cdf8..120774750 100644 --- a/src/api/v1_api.rs +++ b/src/api/v1_api.rs @@ -2,7 +2,6 @@ use std::sync::{Arc}; use actix_web::{HttpResponse, Scope, web}; use serde_json::{json}; use crate::api::api_model::{AppState, PlaylistRequest, ServerConfig, ServerInputConfig, ServerSourceConfig, ServerTargetConfig}; -use crate::download::{get_m3u_playlist, get_xtream_playlist}; use crate::model::config::{Config, ConfigDto, ConfigInput, ConfigInputOptions, ConfigSource, ConfigTarget, InputType, validate_targets}; use log::{error}; use crate::api::download_api::{download_file_info, queue_download_file}; @@ -10,6 +9,7 @@ use crate::config_reader::{read_config, save_api_proxy, save_main_config}; use crate::m3u_filter_error::M3uFilterError; use crate::model::api_proxy::{ApiProxyConfig, ApiProxyServerInfo, TargetUser}; use crate::processing::playlist_processor::exec_processing; +use crate::utils::download; fn _save_config_api_proxy(backup_dir: &str, api_proxy: &mut ApiProxyConfig) -> Option { match save_api_proxy(api_proxy._file_path.as_str(), backup_dir, api_proxy) { @@ -142,8 +142,8 @@ pub(crate) async fn playlist( Some(input) => { let (result, errors) = match input.input_type { - InputType::M3u => get_m3u_playlist(&_app_state.config, &input, &_app_state.config.working_dir).await, - InputType::Xtream => get_xtream_playlist(&input, &_app_state.config.working_dir).await, + InputType::M3u => download::get_m3u_playlist(&_app_state.config, &input, &_app_state.config.working_dir).await, + InputType::Xtream => download::get_xtream_playlist(&input, &_app_state.config.working_dir).await, }; if result.is_empty() { let error_strings: Vec = errors.iter().map(|err| err.to_string()).collect(); diff --git a/src/api/xmltv_api.rs b/src/api/xmltv_api.rs index 02ac0e011..3972e4467 100644 --- a/src/api/xmltv_api.rs +++ b/src/api/xmltv_api.rs @@ -10,7 +10,7 @@ use crate::model::config::{Config, ConfigTarget, InputType}; use crate::model::model_config::TargetType; use crate::repository::m3u_repository::get_m3u_epg_file_path; use crate::repository::xtream_repository::{get_xtream_epg_file_path, get_xtream_storage_path}; -use crate::utils::{get_client_request, path_exists}; +use crate::utils::{file_utils, request_utils}; fn get_epg_path_for_target(config: &Config, target: &ConfigTarget) -> Option { @@ -18,7 +18,7 @@ fn get_epg_path_for_target(config: &Config, target: &ConfigTarget) -> Option { if let Some(epg_path) = get_m3u_epg_file_path(config, &target.get_m3u_filename()) { - if path_exists(&epg_path) { + if file_utils::path_exists(&epg_path) { return Some(epg_path); } else { info!("Cant find epg file for m3u target: {}", epg_path.to_str().unwrap_or("?")) @@ -28,7 +28,7 @@ fn get_epg_path_for_target(config: &Config, target: &ConfigTarget) -> Option { if let Some(storage_path) = get_xtream_storage_path(config, &target.name) { let epg_path = get_xtream_epg_file_path(&storage_path); - if path_exists(&epg_path) { + if file_utils::path_exists(&epg_path) { return Some(epg_path); } else { info!("Cant find epg file for xtream target: {}", epg_path.to_str().unwrap_or("?")) @@ -68,7 +68,7 @@ async fn xmltv_api( debug!("Redirecting epg request to {}", api_url); return HttpResponse::Found().insert_header(("Location", api_url)).finish(); } - let client = get_client_request(input, url, None); + let client = request_utils::get_client_request(input, url, None); if let Ok(response) = client.send().await { if response.status().is_success() { if let Ok(content) = response.text().await { diff --git a/src/api/xtream_api.rs b/src/api/xtream_api.rs index 2269b19c8..8b1169b60 100644 --- a/src/api/xtream_api.rs +++ b/src/api/xtream_api.rs @@ -16,7 +16,7 @@ use crate::model::model_config::{TargetType}; use crate::model::model_playlist::XtreamCluster; use crate::repository::xtream_repository::{COL_CAT_LIVE, COL_CAT_SERIES, COL_CAT_VOD, COL_LIVE, COL_SERIES, COL_VOD, xtream_get_all, xtream_get_stored_stream_info, xtream_persist_stream_info}; -use crate::utils::{get_client_request}; +use crate::utils::request_utils; fn get_xtream_player_api_action_url(input: &ConfigInput, action: &str) -> Option { match input.input_type { @@ -120,7 +120,7 @@ async fn xtream_player_api_stream( let req_headers: HashMap<&str, &[u8]> = req.headers().iter().map(|(k, v)| (k.as_str(), v.as_bytes())).collect(); debug!("Try to open stream {}", &stream_url); if let Ok(url) = Url::parse(&stream_url) { - let client = get_client_request(target_input, url, Some(&req_headers)); + let client = request_utils::get_client_request(target_input, url, Some(&req_headers)); match client.send().await { Ok(response) => { if response.status().is_success() { @@ -216,7 +216,7 @@ async fn xtream_get_stream_info(app_state: &AppState, target_name: &str, stream_ if let Some(info_url) = get_xtream_player_api_info_url(target_input, cluster, stream_id) { if let Ok(url) = Url::parse(&info_url) { - let client = get_client_request(target_input, url, None); + let client = request_utils::get_client_request(target_input, url, None); if let Ok(response) = client.send().await { debug!("{}", response.status()); if response.status().is_success() { @@ -271,16 +271,16 @@ async fn xtream_get_short_epg(app_state: &AppState, user: &UserCredentials, targ return HttpResponse::Found().insert_header(("Location", info_url)).finish(); } - let client = get_client_request(target_input, url, None); + let client = request_utils::get_client_request(target_input, url, None); if let Ok(response) = client.send().await { if response.status().is_success() { - match response.text().await { + return match response.text().await { Ok(content) => { - return HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body(content); + HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body(content) } Err(err) => { error!("Failed to download epg {}", err.to_string()); - return HttpResponse::NoContent().finish(); + HttpResponse::NoContent().finish() } } } diff --git a/src/config_reader.rs b/src/config_reader.rs index a3a5e084d..89628ddcc 100644 --- a/src/config_reader.rs +++ b/src/config_reader.rs @@ -6,12 +6,13 @@ use serde::Serialize; use crate::model::api_proxy::ApiProxyConfig; use crate::model::config::{Config, ConfigDto}; use crate::model::mapping::Mappings; -use crate::{create_m3u_filter_error_result, handle_m3u_filter_error_result, utils}; +use crate::{create_m3u_filter_error_result, handle_m3u_filter_error_result}; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::multi_file_reader::MultiFileReader; +use crate::utils::file_utils; pub(crate) fn read_mappings(args_mapping: Option, cfg: &mut Config) -> Result<(), M3uFilterError> { - let mappings_file: String = args_mapping.unwrap_or(utils::get_default_mappings_path(cfg._config_path.as_str())); + let mappings_file: String = args_mapping.unwrap_or(file_utils::get_default_mappings_path(cfg._config_path.as_str())); match read_mapping(mappings_file.as_str()) { Ok(mappings) => { @@ -25,7 +26,7 @@ pub(crate) fn read_mappings(args_mapping: Option, cfg: &mut Config) -> R } pub(crate) fn read_api_proxy_config(args_api_proxy_config: Option, cfg: &mut Config) { - let api_proxy_config_file: String = args_api_proxy_config.unwrap_or(utils::get_default_api_proxy_config_path(cfg._config_path.as_str())); + let api_proxy_config_file: String = args_api_proxy_config.unwrap_or(file_utils::get_default_api_proxy_config_path(cfg._config_path.as_str())); let api_proxy_config = read_api_proxy(api_proxy_config_file.as_str()); match api_proxy_config { None => { @@ -64,7 +65,7 @@ pub(crate) fn read_config(config_path: &str, config_file: &str, sources_file: &s pub(crate) fn read_mapping(mapping_file: &str) -> Result, M3uFilterError> { let mapping_file = std::path::PathBuf::from(mapping_file); - match utils::open_file(&mapping_file) { + match file_utils::open_file(&mapping_file) { Ok(file) => { let mapping: Result = serde_yaml::from_reader(file); match mapping { @@ -86,7 +87,7 @@ pub(crate) fn read_mapping(mapping_file: &str) -> Result, M3uFi } pub(crate) fn read_api_proxy(api_proxy_file: &str) -> Option { - match utils::open_file(&std::path::PathBuf::from(api_proxy_file)) { + match file_utils::open_file(&std::path::PathBuf::from(api_proxy_file)) { Ok(file) => { let mapping: Result = serde_yaml::from_reader(file); match mapping { diff --git a/src/main.rs b/src/main.rs index cc19e2d82..4c363690d 100644 --- a/src/main.rs +++ b/src/main.rs @@ -13,19 +13,19 @@ use log::{error, info, LevelFilter}; use crate::config_reader::{read_api_proxy_config, read_config, read_mappings}; use crate::model::config::{Config, ProcessTargets, validate_targets}; use crate::processing::playlist_processor::exec_processing; +use crate::utils::file_utils; mod m3u_filter_error; mod config_reader; mod model; mod filter; mod repository; -mod download; -mod utils; mod messaging; mod test; mod api; mod processing; mod multi_file_reader; +mod utils; #[derive(Parser)] #[command(name = "m3u-filter")] @@ -72,14 +72,22 @@ const VERSION: &str = env!("CARGO_PKG_VERSION"); fn main() { let args = Args::parse(); - init_logger(&args.log_level.unwrap_or("info".to_string())); + init_logger(args.log_level.as_ref().unwrap_or(&"info".to_string())); - let config_path: String = args.config_path.unwrap_or(utils::get_default_config_path()); - let config_file: String = args.config_file.unwrap_or(utils::get_default_config_file_path(&config_path)); - let sources_file: String = args.source_file.unwrap_or(utils::get_default_sources_file_path(&config_path)); + let config_path: String = args.config_path.unwrap_or(file_utils::get_default_config_path()); + let config_file: String = args.config_file.unwrap_or(file_utils::get_default_config_file_path(&config_path)); + let sources_file: String = args.source_file.unwrap_or(file_utils::get_default_sources_file_path(&config_path)); let mut cfg = read_config(config_path.as_str(), config_file.as_str(), sources_file.as_str()).unwrap_or_else(|err| exit!("{}", err)); + + if args.log_level.is_none() { + if let Some(log_level) = &cfg.log_level { + info!("Setting log level to: {}", get_log_level(log_level.as_str())); + log::set_max_level(get_log_level(log_level.as_str())); + } + } + let targets = validate_targets(&args.target, &cfg.sources).unwrap_or_else(|err| exit!("{}", err)); info!("Version: {}", VERSION); @@ -116,16 +124,20 @@ fn start_in_server_mode(cfg: Arc, targets: Arc) { }; } +fn get_log_level(log_level: &str) -> LevelFilter { + match log_level.to_lowercase().as_str() { + "trace" => LevelFilter::Trace, + "debug" => LevelFilter::Debug, + "info" => LevelFilter::Info, + "warn" => LevelFilter::Warn, + "error" => LevelFilter::Error, + _ => LevelFilter::Info, + } +} + fn init_logger(log_level: &str) { let mut log_builder = Builder::from_default_env(); // Set the log level based on the parsed value - match log_level.to_lowercase().as_str() { - "trace" => log_builder.filter_level(LevelFilter::Trace), - "debug" => log_builder.filter_level(LevelFilter::Debug), - "info" => log_builder.filter_level(LevelFilter::Info), - "warn" => log_builder.filter_level(LevelFilter::Warn), - "error" => log_builder.filter_level(LevelFilter::Error), - _ => log_builder.filter_level(LevelFilter::Info), - }; + log_builder.filter_level(get_log_level(log_level)); log_builder.init(); } diff --git a/src/model/config.rs b/src/model/config.rs index 2b699505d..beccda499 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -15,8 +15,7 @@ use crate::model::api_proxy::{ApiProxyConfig, UserCredentials}; use crate::model::mapping::Mapping; use crate::model::mapping::Mappings; use crate::model::model_config::{default_as_false, default_as_true, default_as_zero, ItemField, ProcessingOrder, SortOrder, TargetType}; -use crate::utils; -use crate::utils::get_working_path; +use crate::utils::file_utils; fn default_as_frm() -> ProcessingOrder { ProcessingOrder::Frm } @@ -26,7 +25,6 @@ fn default_as_empty_map() -> HashMap { HashMap::new() } fn default_as_empty_list() -> Vec { vec![] } - #[macro_export] macro_rules! create_m3u_filter_error_result { ($kind: expr, $($arg:tt)*) => { @@ -571,6 +569,7 @@ pub(crate) struct Config { pub video: Option, pub schedule: Option, pub messaging: Option, + pub log_level: Option, #[serde(skip_serializing, skip_deserializing)] pub _api_proxy: Arc>>, #[serde(skip_serializing, skip_deserializing)] @@ -659,7 +658,7 @@ impl Config { } pub fn prepare(&mut self) -> Result<(), M3uFilterError> { - self.working_dir = get_working_path(&self.working_dir); + self.working_dir = file_utils::get_working_path(&self.working_dir); if self.backup_dir.is_none() { self.backup_dir = Some(PathBuf::from(&self.working_dir).join(".backup").into_os_string().to_string_lossy().to_string()); } @@ -733,7 +732,7 @@ impl Config { if wrpb.is_relative() { let mut wrpb2 = std::path::PathBuf::from(&self.working_dir).join(&self.api.web_root); if !wrpb2.exists() { - wrpb2 = utils::get_exe_path().join(&self.api.web_root); + wrpb2 = file_utils::get_exe_path().join(&self.api.web_root); } if !wrpb2.exists() { let cwd = std::env::current_dir(); diff --git a/src/processing/m3u_parser.rs b/src/processing/m3u_parser.rs index 6ab8f373e..0a1dbab7f 100644 --- a/src/processing/m3u_parser.rs +++ b/src/processing/m3u_parser.rs @@ -4,6 +4,7 @@ use std::rc::Rc; use crate::model::config::Config; use crate::model::model_config::default_as_empty_rc_str; use crate::model::model_playlist::{default_playlist_item_type, default_stream_cluster, PlaylistGroup, PlaylistItem, PlaylistItemHeader, PlaylistItemType, XtreamCluster}; +use crate::utils::string_utils; fn token_value(it: &mut std::str::Chars) -> String { if let Some(oc) = it.next() { @@ -116,8 +117,10 @@ fn process_header(video_suffixes: &Vec<&str>, content: &String, url: String) -> } c = it.next(); } - if plih.group.is_empty() { - plih.group = Rc::new(String::from("Unknown")); + if plih.id.is_empty() { + if let Some(chanid) = extract_id_from_url(url.as_str()) { + plih.id = Rc::new(chanid); + } } plih.epg_channel_id = Some(Rc::clone(&plih.id)); } @@ -143,6 +146,12 @@ fn process_header(video_suffixes: &Vec<&str>, content: &String, url: String) -> plih } +fn extract_id_from_url(url: &str) -> Option { + if let Some(filename) = url.split('/').last() { + return filename.rsplit('.').next().map(|stem| stem.to_string()); + } + None +} pub(crate) fn parse_m3u(cfg: &Config, lines: &Vec) -> Vec { let mut groups: std::collections::HashMap, Vec> = std::collections::HashMap::new(); @@ -150,6 +159,7 @@ pub(crate) fn parse_m3u(cfg: &Config, lines: &Vec) -> Vec let mut header: Option = None; let mut group: Option = None; + let mut playlist = Vec::new(); let video_suffixes = cfg.video.as_ref().unwrap().extensions.iter().map(|ext| ext.as_str()).collect(); for line in lines { if line.starts_with("#EXTINF") { @@ -170,20 +180,31 @@ pub(crate) fn parse_m3u(cfg: &Config, lines: &Vec) -> Vec item.header.borrow_mut().group = Rc::new(group_value); } } - let key = Rc::clone(&item.header.borrow().group); - // let key2 = String::from(&item.header.group); - match groups.entry(Rc::clone(&key)) { - std::collections::hash_map::Entry::Vacant(e) => { - e.insert(vec![item]); - sort_order.push(Rc::clone(&key)); - } - std::collections::hash_map::Entry::Occupied(mut e) => { e.get_mut().push(item); } - } + playlist.push(item); } header = None; group = None; } + for item in &playlist { + if item.header.borrow().group.is_empty() { + let current_title = item.header.borrow().title.to_owned(); + item.header.borrow_mut().group = Rc::new(string_utils::get_title_group(current_title.as_str())); + } + } + + playlist.drain(..).for_each(|item| { + let key = Rc::clone(&item.header.borrow().group); + // let key2 = String::from(&item.header.group); + match groups.entry(Rc::clone(&key)) { + std::collections::hash_map::Entry::Vacant(e) => { + e.insert(vec![item]); + sort_order.push(Rc::clone(&key)); + } + std::collections::hash_map::Entry::Occupied(mut e) => { e.get_mut().push(item); } + } + }); + let mut result: Vec = vec![]; for (grp_id, (key, channels)) in (1_u32..).zip(groups.into_iter()) { let cluster = channels.first().map(|pli| pli.header.borrow().xtream_cluster.clone()); diff --git a/src/processing/playlist_processor.rs b/src/processing/playlist_processor.rs index 6f0373f8f..50d49475d 100644 --- a/src/processing/playlist_processor.rs +++ b/src/processing/playlist_processor.rs @@ -11,7 +11,6 @@ use log::{debug, error, info}; use unidecode::unidecode; use crate::{Config, get_errors_notify_message, model::config, valid_property}; -use crate::download::{get_m3u_playlist, get_xmltv, get_xtream_playlist, get_xtream_playlist_series}; use crate::filter::{get_field_value, MockValueProcessor, set_field_value, ValueProvider}; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::messaging::{MsgKind, send_message}; @@ -26,6 +25,7 @@ use crate::processing::xmltv_parser::flatten_tvguide; use crate::repository::epg_repository::write_epg; use crate::repository::m3u_repository::{write_m3u_playlist, write_strm_playlist}; use crate::repository::xtream_repository::write_xtream_playlist; +use crate::utils::download; fn filter_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option> { debug!("Filtering {} groups", playlist.len()); @@ -284,11 +284,11 @@ async fn process_source(cfg: Arc, source_idx: usize, user_targets: Arc

get_m3u_playlist(&cfg, input, &cfg.working_dir).await, - InputType::Xtream => get_xtream_playlist(input, &cfg.working_dir).await, + InputType::M3u => download::get_m3u_playlist(&cfg, input, &cfg.working_dir).await, + InputType::Xtream => download::get_xtream_playlist(input, &cfg.working_dir).await, }; let (tvguide, mut tvguide_errors) = if error_list.is_empty() { - get_xmltv(&cfg, input, &cfg.working_dir).await + download::get_xmltv(&cfg, input, &cfg.working_dir).await } else { (None, vec![]) }; @@ -433,7 +433,7 @@ pub(crate) async fn process_playlist<'a>(playlists: &mut [FetchedPlaylist<'a>], (false, 0) }; if resolve_series { - let mut series_playlist = get_xtream_playlist_series(fpl, errors, resolve_series_delay).await; + let mut series_playlist = download::get_xtream_playlist_series(fpl, errors, resolve_series_delay).await; // original content saved into original list for plg in &series_playlist { fpl.update_playlist(plg); diff --git a/src/processing/playlist_watch.rs b/src/processing/playlist_watch.rs index bf4d0890f..e401c3696 100644 --- a/src/processing/playlist_watch.rs +++ b/src/processing/playlist_watch.rs @@ -5,7 +5,7 @@ use regex::Regex; use crate::messaging::{MsgKind, send_message}; use crate::model::config::Config; use crate::model::model_playlist::PlaylistGroup; -use crate::utils; +use crate::utils::file_utils; pub(crate) fn process_group_watch(cfg: &Config, target_name: &str, pl: &PlaylistGroup) { let mut new_tree = BTreeSet::new(); @@ -18,7 +18,7 @@ pub(crate) fn process_group_watch(cfg: &Config, target_name: &str, pl: &Playlist let filename_re = Regex::new(r"[^A-Za-z0-9_-]").unwrap(); let file_name = format!("watch_{}_{}", target_name, &pl.title); let watch_filename = format!("{}.bin", filename_re.replace_all(&file_name, "_")).to_string(); - match utils::get_file_path(&cfg.working_dir, Some(std::path::PathBuf::from(&watch_filename))) { + match file_utils::get_file_path(&cfg.working_dir, Some(std::path::PathBuf::from(&watch_filename))) { Some(path) => { let save_path = path.clone(); let mut changed = false; diff --git a/src/repository/m3u_repository.rs b/src/repository/m3u_repository.rs index c3674c976..e0aa1f4a6 100644 --- a/src/repository/m3u_repository.rs +++ b/src/repository/m3u_repository.rs @@ -4,11 +4,11 @@ use std::io::Write; use chrono::Datelike; use log::error; -use crate::{create_m3u_filter_error_result, utils}; +use crate::{create_m3u_filter_error_result}; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::{Config, ConfigTarget}; use crate::model::model_playlist::{PlaylistGroup, PlaylistItemType}; -use crate::utils::add_prefix_to_filename; +use crate::utils::file_utils; fn check_write(res: std::io::Result<()>) -> Result<(), std::io::Error> { match res { @@ -83,12 +83,12 @@ fn kodi_style_rename(name: &String, style: &KodiStyle) -> String { } pub(crate) fn get_m3u_file_path(cfg: &Config, filename: &Option) -> Option { - utils::get_file_path(&cfg.working_dir, Some(std::path::PathBuf::from(&filename.as_ref().unwrap()))) + file_utils::get_file_path(&cfg.working_dir, Some(std::path::PathBuf::from(&filename.as_ref().unwrap()))) } pub(crate) fn get_m3u_epg_file_path(cfg: &Config, filename: &Option) -> Option { - utils::get_file_path(&cfg.working_dir, Some(std::path::PathBuf::from(&filename.as_ref().unwrap()))) - .map(|path| add_prefix_to_filename(&path, "epg_", Some("xml"))) + file_utils::get_file_path(&cfg.working_dir, Some(std::path::PathBuf::from(&filename.as_ref().unwrap()))) + .map(|path| file_utils::add_prefix_to_filename(&path, "epg_", Some("xml"))) } @@ -145,7 +145,7 @@ pub(crate) fn write_strm_playlist(target: &ConfigTarget, cfg: &Config, new_playl let cleanup = target.options.as_ref().map_or(false, |o| o.cleanup); let kodi_style = target.options.as_ref().map_or(false, |o| o.kodi_style); - if let Some(path) = utils::get_file_path(&cfg.working_dir, Some(std::path::PathBuf::from(&filename.as_ref().unwrap()))) { + if let Some(path) = file_utils::get_file_path(&cfg.working_dir, Some(std::path::PathBuf::from(&filename.as_ref().unwrap()))) { if cleanup { let _ = std::fs::remove_dir_all(&path); } diff --git a/src/repository/xtream_repository.rs b/src/repository/xtream_repository.rs index bbd2b6345..36bf8150c 100644 --- a/src/repository/xtream_repository.rs +++ b/src/repository/xtream_repository.rs @@ -10,9 +10,10 @@ use serde::Serialize; use serde_json::{json, Map, Value}; use crate::model::config::{Config, ConfigInput, ConfigTarget}; use crate::model::model_playlist::{PlaylistGroup, PlaylistItemHeader, PlaylistItemType, XtreamCluster}; -use crate::{create_m3u_filter_error_result, utils}; +use crate::{create_m3u_filter_error_result}; use crate::api::api_model::AppState; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; +use crate::utils::file_utils; type IndexTree = BTreeMap; @@ -41,7 +42,7 @@ const SERIES_STREAM_FIELDS: &[&str] = &[ pub(crate) fn get_xtream_storage_path(cfg: &Config, target_name: &str) -> Option { - utils::get_file_path(&cfg.working_dir, Some(std::path::PathBuf::from(target_name.replace(' ', "_")))) + file_utils::get_file_path(&cfg.working_dir, Some(std::path::PathBuf::from(target_name.replace(' ', "_")))) } pub(crate) fn get_xtream_epg_file_path(path: &Path) -> PathBuf { @@ -105,7 +106,7 @@ fn write_xtream_info(app_state: &AppState, target_name: &str, stream_id: i32, cl write_index(&idx_path, index_tree)?; } Err(err) => { - return Err(err) + return Err(err); } } } @@ -133,6 +134,7 @@ pub(crate) fn write_xtream_playlist(target: &ConfigTarget, cfg: &Config, playlis let mut series_map = HashMap::::new(); let mut channel_num: i32 = 0; + let mut errors = Vec::new(); for plg in playlist { if !&plg.channels.is_empty() { match &plg.xtream_cluster { @@ -148,91 +150,93 @@ pub(crate) fn write_xtream_playlist(target: &ConfigTarget, cfg: &Config, playlis for pli in &plg.channels { let header = &pli.header.borrow(); - if header.item_type == PlaylistItemType::Series { - // we skip resolved series, because this is only necessary when writing m3u files - continue; - } - channel_num += 1; - let mut document = serde_json::Map::from_iter([ - ("category_id".to_string(), Value::String(format!("{}", &plg.id))), - ("category_ids".to_string(), Value::Array(Vec::from([Value::Number(serde_json::Number::from(plg.id.to_owned()))]))), - ("name".to_string(), Value::String(header.name.as_ref().clone())), - ("num".to_string(), Value::Number(serde_json::Number::from(channel_num))), - ("title".to_string(), Value::String(header.title.as_ref().clone())), - ("stream_icon".to_string(), Value::String(header.logo.as_ref().clone())), - ]); + if let Ok(stream_id) = header.id.parse::() { + if header.item_type == PlaylistItemType::Series { + // we skip resolved series, because this is only necessary when writing m3u files + continue; + } + channel_num += 1; + let mut document = serde_json::Map::from_iter([ + ("category_id".to_string(), Value::String(format!("{}", &plg.id))), + ("category_ids".to_string(), Value::Array(Vec::from([Value::Number(serde_json::Number::from(plg.id.to_owned()))]))), + ("name".to_string(), Value::String(header.name.as_ref().clone())), + ("num".to_string(), Value::Number(serde_json::Number::from(channel_num))), + ("title".to_string(), Value::String(header.title.as_ref().clone())), + ("stream_icon".to_string(), Value::String(header.logo.as_ref().clone())), + ]); - let stream_id = header.id.parse::().unwrap(); - let stream_id_value = Value::Number(serde_json::Number::from(stream_id)); - match header.xtream_cluster { - XtreamCluster::Live => { - document.insert("stream_id".to_string(), stream_id_value); - if skip_live_direct_source { - document.insert("direct_source".to_string(), Value::String("".to_string())); - } else { - document.insert("direct_source".to_string(), Value::String(header.url.as_ref().clone())); + let stream_id_value = Value::Number(serde_json::Number::from(stream_id)); + match header.xtream_cluster { + XtreamCluster::Live => { + document.insert("stream_id".to_string(), stream_id_value); + if skip_live_direct_source { + document.insert("direct_source".to_string(), Value::String("".to_string())); + } else { + document.insert("direct_source".to_string(), Value::String(header.url.as_ref().clone())); + } + document.insert("thumbnail".to_string(), Value::String(header.logo_small.as_ref().clone())); + document.insert("custom_sid".to_string(), Value::String("".to_string())); + document.insert("epg_channel_id".to_string(), match &header.epg_channel_id { + None => Value::Null, + Some(epg_id) => Value::String(epg_id.as_ref().clone()) + }); } - document.insert("thumbnail".to_string(), Value::String(header.logo_small.as_ref().clone())); - document.insert("custom_sid".to_string(), Value::String("".to_string())); - document.insert("epg_channel_id".to_string(), match &header.epg_channel_id { - None => Value::Null, - Some(epg_id) => Value::String(epg_id.as_ref().clone()) - }); - } - XtreamCluster::Video => { - document.insert("stream_id".to_string(), stream_id_value); - if skip_video_direct_source { - document.insert("direct_source".to_string(), Value::String("".to_string())); - } else { - document.insert("direct_source".to_string(), Value::String(header.url.as_ref().clone())); + XtreamCluster::Video => { + document.insert("stream_id".to_string(), stream_id_value); + if skip_video_direct_source { + document.insert("direct_source".to_string(), Value::String("".to_string())); + } else { + document.insert("direct_source".to_string(), Value::String(header.url.as_ref().clone())); + } + document.insert("custom_sid".to_string(), Value::String("".to_string())); } - document.insert("custom_sid".to_string(), Value::String("".to_string())); - } - XtreamCluster::Series => { - document.insert("series_id".to_string(), stream_id_value); - } - }; + XtreamCluster::Series => { + document.insert("series_id".to_string(), stream_id_value); + } + }; - if let Some(add_props) = &header.additional_properties { - for (field_name, field_value) in add_props { - document.insert(field_name.to_string(), field_value.to_owned()); + if let Some(add_props) = &header.additional_properties { + for (field_name, field_value) in add_props { + document.insert(field_name.to_string(), field_value.to_owned()); + } } + + match header.xtream_cluster { + XtreamCluster::Live => { + append_mandatory_fields(&mut document, LIVE_STREAM_FIELDS); + } + XtreamCluster::Video => { + append_mandatory_fields(&mut document, VIDEO_STREAM_FIELDS); + } + XtreamCluster::Series => { + append_prepared_series_properties(header, &mut document); + append_mandatory_fields(&mut document, SERIES_STREAM_FIELDS); + append_release_date(&mut document); + } + }; + + match header.xtream_cluster { + XtreamCluster::Live => {} + XtreamCluster::Series => { + series_map.insert(stream_id, serde_json::to_string(&document).unwrap()); + } + XtreamCluster::Video => { + vod_map.insert(stream_id, serde_json::to_string(&document).unwrap()); + } + } + + match header.xtream_cluster { + XtreamCluster::Live => &mut live_col, + XtreamCluster::Series => &mut series_col, + XtreamCluster::Video => &mut vod_col, + }.push(Value::Object(document)); + } else { + errors.push(format!("Channel does not have an id: {}", &header.title)); } - - match header.xtream_cluster { - XtreamCluster::Live => { - append_mandatory_fields(&mut document, LIVE_STREAM_FIELDS); - } - XtreamCluster::Video => { - append_mandatory_fields(&mut document, VIDEO_STREAM_FIELDS); - } - XtreamCluster::Series => { - append_prepared_series_properties(header, &mut document); - append_mandatory_fields(&mut document, SERIES_STREAM_FIELDS); - append_release_date(&mut document); - } - }; - - match header.xtream_cluster { - XtreamCluster::Live => {} - XtreamCluster::Series => { - series_map.insert(stream_id, serde_json::to_string(&document).unwrap()); - } - XtreamCluster::Video => { - vod_map.insert(stream_id, serde_json::to_string(&document).unwrap()); - } - } - - match header.xtream_cluster { - XtreamCluster::Live => &mut live_col, - XtreamCluster::Series => &mut series_col, - XtreamCluster::Video => &mut vod_col, - }.push(Value::Object(document)); } } } - let mut errors = Vec::new(); for (col_path, data) in [ (get_collection_path(&path, COL_CAT_LIVE), &cat_live_col), (get_collection_path(&path, COL_CAT_VOD), &cat_vod_col), @@ -378,22 +382,22 @@ pub(crate) async fn xtream_persist_stream_info( .map(|o| o.xtream_info_cache).unwrap_or(false); if cache_info { if let Some(path) = get_xtream_storage_path(&app_state.config, target_name) { - let lock = app_state.shared_locks.get_lock(target_name); - let shared_lock = lock.write().unwrap(); - let mut index_tree = { - let (col_path, idx_path) = get_info_collection_and_idx_path(&path, cluster); - if idx_path.exists() && col_path.exists() { - load_index(&idx_path).unwrap_or_default() - } else { - IndexTree::new() - } - }; - match write_xtream_info(app_state, target_name, stream_id, cluster, content, - &mut index_tree) { - Ok(_) => {} - Err(err) => { error!("{}", err.to_string()); } + let lock = app_state.shared_locks.get_lock(target_name); + let shared_lock = lock.write().unwrap(); + let mut index_tree = { + let (col_path, idx_path) = get_info_collection_and_idx_path(&path, cluster); + if idx_path.exists() && col_path.exists() { + load_index(&idx_path).unwrap_or_default() + } else { + IndexTree::new() } - drop(shared_lock); + }; + match write_xtream_info(app_state, target_name, stream_id, cluster, content, + &mut index_tree) { + Ok(_) => {} + Err(err) => { error!("{}", err.to_string()); } + } + drop(shared_lock); } } } \ No newline at end of file diff --git a/src/download.rs b/src/utils/download.rs similarity index 87% rename from src/download.rs rename to src/utils/download.rs index 2c09dd7b7..9dad88ba6 100644 --- a/src/download.rs +++ b/src/utils/download.rs @@ -2,23 +2,22 @@ use std::path::PathBuf; use std::sync::atomic::{AtomicU32}; use std::thread::sleep; use log::debug; -use crate::{utils}; use crate::m3u_filter_error::M3uFilterError; use crate::model::config::{Config, ConfigInput}; use crate::model::model_playlist::{FetchedPlaylist, PlaylistGroup, PlaylistItem, PlaylistItemType, XtreamCluster}; use crate::model::xmltv::TVGuide; use crate::processing::{m3u_parser, xmltv_parser, xtream_parser}; use crate::processing::xtream_parser::parse_xtream_series_info; -use crate::utils::add_prefix_to_filename; +use crate::utils::{file_utils, request_utils}; fn prepare_file_path(input: &ConfigInput, working_dir: &String, action: &str) -> Option { let persist_file: Option = match &input.persist { - Some(persist_path) => utils::prepare_persist_path(persist_path.as_str(), action), + Some(persist_path) => file_utils::prepare_persist_path(persist_path.as_str(), action), _ => None }; if persist_file.is_some() { - let file_path = utils::get_file_path(working_dir, persist_file); + let file_path = file_utils::get_file_path(working_dir, persist_file); debug!("persist to file: {:?}", match &file_path { Some(fp) => fp.display().to_string(), _ => "".to_string() @@ -32,7 +31,7 @@ fn prepare_file_path(input: &ConfigInput, working_dir: &String, action: &str) -> pub(crate) async fn get_m3u_playlist(cfg: &Config, input: &ConfigInput, working_dir: &String) -> (Vec, Vec) { let url = input.url.to_owned(); let persist_file_path = prepare_file_path(input, working_dir, ""); - match utils::get_input_text_content(input, working_dir, &url, persist_file_path).await { + match request_utils::get_input_text_content(input, working_dir, &url, persist_file_path).await { Ok(text) => { let lines = text.lines().map(String::from).collect(); (m3u_parser::parse_m3u(cfg, &lines), vec![]) @@ -56,7 +55,7 @@ pub(crate) async fn get_xtream_playlist_series<'a>(fpl: &mut FetchedPlaylist<'a> (fetch_series, header.url.to_string()) }; if fetch_series { - match utils::get_input_json_content(fpl.input, series_info_url.as_str(), None).await { + match request_utils::get_input_json_content(fpl.input, series_info_url.as_str(), None).await { Ok(series_content) => { match parse_xtream_series_info(&series_content, pli.header.borrow().group.as_str(), input) { Ok(series_info) => { @@ -107,9 +106,9 @@ pub(crate) async fn get_xtream_playlist(input: &ConfigInput, working_dir: &Strin let category_file_path = prepare_file_path(input, working_dir, format!("{}_", category).as_str()); let stream_file_path = prepare_file_path(input, working_dir, format!("{}_", stream).as_str()); - match utils::get_input_json_content(input, category_url.as_str(), category_file_path).await { + match request_utils::get_input_json_content(input, category_url.as_str(), category_file_path).await { Ok(category_content) => { - match utils::get_input_json_content(input, stream_url.as_str(), stream_file_path).await { + match request_utils::get_input_json_content(input, stream_url.as_str(), stream_file_path).await { Ok(stream_content) => { match xtream_parser::parse_xtream(&category_id_cnt, xtream_cluster, @@ -140,8 +139,8 @@ pub(crate) async fn get_xmltv(_cfg: &Config, input: &ConfigInput, working_dir: & None => (None, vec![]), Some(url) => { debug!("Getting epg file path for url: {}", url); - let persist_file_path = prepare_file_path(input, working_dir, "").map(|path| add_prefix_to_filename(&path, "epg_", Some("xml"))); - match utils::get_input_text_content(input, working_dir, url, persist_file_path).await { + let persist_file_path = prepare_file_path(input, working_dir, "").map(|path| file_utils::add_prefix_to_filename(&path, "epg_", Some("xml"))); + match request_utils::get_input_text_content(input, working_dir, url, persist_file_path).await { Ok(xml_content) => { (xmltv_parser::parse_tvguide(xml_content.as_str()), vec![]) } diff --git a/src/utils/file_utils.rs b/src/utils/file_utils.rs new file mode 100644 index 000000000..414daf6af --- /dev/null +++ b/src/utils/file_utils.rs @@ -0,0 +1,162 @@ +use std::fs; +use std::io::{Write}; +use std::path::{Path, PathBuf}; +use log::{debug, error}; +use path_absolutize::*; + +#[macro_export] +macro_rules! exit { + ($($arg:tt)*) => {{ + error!($($arg)*); + std::process::exit(1); + }}; +} + + +pub(crate) fn get_exe_path() -> PathBuf { + let default_path = std::path::PathBuf::from("./"); + let current_exe = std::env::current_exe(); + match current_exe { + Ok(exe) => { + match fs::read_link(&exe) { + Ok(f) => f.parent().map_or(default_path, |p| p.to_path_buf()), + Err(_) => return exe.parent().map_or(default_path, |p| p.to_path_buf()) + } + } + Err(_) => default_path + } +} + +fn get_default_path(file: &str) -> String { + let path: PathBuf = get_exe_path(); + let default_path = path.join(file); + String::from(if default_path.exists() { + default_path.to_str().unwrap_or(file) + } else { + file + }) +} + +fn get_default_file_path(config_path: &str, file: &str) -> String { + let path: PathBuf = PathBuf::from(config_path); + let default_path = path.join(file); + String::from(if default_path.exists() { + default_path.to_str().unwrap_or(file) + } else { + file + }) +} + +pub(crate) fn get_default_config_path() -> String { + get_default_path("config") +} + +pub(crate) fn get_default_config_file_path(config_path: &str) -> String { + get_default_file_path(config_path, "config.yml") +} + +pub(crate) fn get_default_sources_file_path(config_path: &str) -> String { + get_default_file_path(config_path, "source.yml") +} + +pub(crate) fn get_default_mappings_path(config_path: &str) -> String { + get_default_file_path(config_path, "mapping.yml") +} + +pub(crate) fn get_default_api_proxy_config_path(config_path: &str) -> String { + get_default_file_path(config_path, "api-proxy.yml") +} + +pub(crate) fn get_working_path(wd: &String) -> String { + let current_dir = std::env::current_dir().unwrap(); + if wd.is_empty() { + String::from(current_dir.to_str().unwrap_or(".")) + } else { + let work_path = std::path::PathBuf::from(wd); + let wdpath = match fs::metadata(&work_path) { + Ok(md) => { + if md.is_dir() && !md.permissions().readonly() { + match work_path.canonicalize() { + Ok(ap) => Some(ap), + Err(_) => None + } + } else { + error!("Path not found {:?}", &work_path); + None + } + } + Err(_) => None, + }; + let rp: PathBuf = match wdpath { + Some(d) => d, + None => current_dir.join(wd) + }; + match rp.canonicalize() { + Ok(ap) => String::from(ap.to_str().unwrap_or("./")), + Err(_) => { + error!("Path not found {:?}", &rp); + String::from("./") + } + } + } +} + +pub(crate) fn open_file(file_name: &Path) -> Result { + fs::File::open(file_name) +} + +pub(crate) fn persist_file(persist_file: Option, text: &String) { + if let Some(path_buf) = persist_file { + let filename = &path_buf.to_str().unwrap_or("?"); + match fs::File::create(&path_buf) { + Ok(mut file) => match file.write_all(text.as_bytes()) { + Ok(_) => debug!("persisted: {}", filename), + Err(e) => error!("failed to persist file {}, {}", filename, e) + }, + Err(e) => error!("failed to persist file {}, {}", filename, e) + } + } +} + +pub(crate) fn prepare_persist_path(file_name: &str, date_prefix: &str) -> Option { + let now = chrono::Local::now(); + let filename = file_name.replace("{}", format!("{}{}", date_prefix, now.format("%Y%m%d_%H%M%S").to_string().as_str()).as_str()); + Some(std::path::PathBuf::from(filename)) +} + +pub(crate) fn get_file_path(wd: &String, path: Option) -> Option { + match path { + Some(p) => { + if p.is_relative() { + let pb = PathBuf::from(wd); + match pb.join(&p).absolutize() { + Ok(os) => Some(PathBuf::from(os)), + Err(e) => { + error!("path is not relative {:?}", e); + Some(p) + } + } + } else { + Some(p) + } + } + None => None + } +} + +pub(crate) fn add_prefix_to_filename(path: &Path, prefix: &str, ext: Option<&str>) -> PathBuf { + let file_name = path.file_name().unwrap_or_default(); + let new_file_name = format!("{}{}", prefix, file_name.to_string_lossy()); + let result = path.with_file_name(new_file_name); + match ext { + None => result, + Some(extension) => result.with_extension(extension) + } +} + +pub(crate) fn path_exists(file_path: &Path) -> bool { + if let Ok(metadata) = fs::metadata(file_path) { + return metadata.is_file(); + } + false +} diff --git a/src/utils/mod.rs b/src/utils/mod.rs new file mode 100644 index 000000000..79c00e8ce --- /dev/null +++ b/src/utils/mod.rs @@ -0,0 +1,4 @@ +pub (crate) mod file_utils; +pub (crate) mod request_utils; +pub (crate) mod download; +pub (crate) mod string_utils; \ No newline at end of file diff --git a/src/utils.rs b/src/utils/request_utils.rs similarity index 59% rename from src/utils.rs rename to src/utils/request_utils.rs index 197b88bef..be3e43cfd 100644 --- a/src/utils.rs +++ b/src/utils/request_utils.rs @@ -1,113 +1,16 @@ use std::collections::{HashMap, HashSet}; use std::fs; -use std::io::{Read, Write}; -use std::path::{Path, PathBuf}; +use std::io::{Read}; +use std::path::{PathBuf}; use log::{debug, error, Level}; -use path_absolutize::*; use reqwest::header::{HeaderMap, HeaderName, HeaderValue}; use crate::create_m3u_filter_error_result; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::{ConfigInput}; +use crate::utils::file_utils::{get_file_path, open_file, persist_file}; -#[macro_export] -macro_rules! exit { - ($($arg:tt)*) => {{ - error!($($arg)*); - std::process::exit(1); - }}; -} - - -pub(crate) fn get_exe_path() -> PathBuf { - let default_path = std::path::PathBuf::from("./"); - let current_exe = std::env::current_exe(); - match current_exe { - Ok(exe) => { - match fs::read_link(&exe) { - Ok(f) => f.parent().map_or(default_path, |p| p.to_path_buf()), - Err(_) => return exe.parent().map_or(default_path, |p| p.to_path_buf()) - } - } - Err(_) => default_path - } -} - -fn get_default_path(file: &str) -> String { - let path: PathBuf = get_exe_path(); - let default_path = path.join(file); - String::from(if default_path.exists() { - default_path.to_str().unwrap_or(file) - } else { - file - }) -} - -fn get_default_file_path(config_path: &str, file: &str) -> String { - let path: PathBuf = PathBuf::from(config_path); - let default_path = path.join(file); - String::from(if default_path.exists() { - default_path.to_str().unwrap_or(file) - } else { - file - }) -} - -pub(crate) fn get_default_config_path() -> String { - get_default_path("config") -} - -pub(crate) fn get_default_config_file_path(config_path: &str) -> String { - get_default_file_path(config_path, "config.yml") -} - -pub(crate) fn get_default_sources_file_path(config_path: &str) -> String { - get_default_file_path(config_path, "source.yml") -} - -pub(crate) fn get_default_mappings_path(config_path: &str) -> String { - get_default_file_path(config_path, "mapping.yml") -} - -pub(crate) fn get_default_api_proxy_config_path(config_path: &str) -> String { - get_default_file_path(config_path, "api-proxy.yml") -} - -pub(crate) fn get_working_path(wd: &String) -> String { - let current_dir = std::env::current_dir().unwrap(); - if wd.is_empty() { - String::from(current_dir.to_str().unwrap_or(".")) - } else { - let work_path = std::path::PathBuf::from(wd); - let wdpath = match fs::metadata(&work_path) { - Ok(md) => { - if md.is_dir() && !md.permissions().readonly() { - match work_path.canonicalize() { - Ok(ap) => Some(ap), - Err(_) => None - } - } else { - error!("Path not found {:?}", &work_path); - None - } - } - Err(_) => None, - }; - let rp: PathBuf = match wdpath { - Some(d) => d, - None => current_dir.join(wd) - }; - match rp.canonicalize() { - Ok(ap) => String::from(ap.to_str().unwrap_or("./")), - Err(_) => { - error!("Path not found {:?}", &rp); - String::from("./") - } - } - } -} - -pub(crate) fn open_file(file_name: &Path) -> Result { - fs::File::open(file_name) +pub(crate) fn bytes_to_megabytes(bytes: u64) -> u64 { + bytes / 1_048_576 } pub(crate) async fn get_input_text_content(input: &ConfigInput, working_dir: &String, url_str: &str, persist_filepath: Option) -> Result { @@ -170,47 +73,6 @@ pub(crate) async fn get_input_text_content(input: &ConfigInput, working_dir: &St } } - -fn persist_file(persist_file: Option, text: &String) { - if let Some(path_buf) = persist_file { - let filename = &path_buf.to_str().unwrap_or("?"); - match fs::File::create(&path_buf) { - Ok(mut file) => match file.write_all(text.as_bytes()) { - Ok(_) => debug!("persisted: {}", filename), - Err(e) => error!("failed to persist file {}, {}", filename, e) - }, - Err(e) => error!("failed to persist file {}, {}", filename, e) - } - } -} - -pub(crate) fn prepare_persist_path(file_name: &str, date_prefix: &str) -> Option { - let now = chrono::Local::now(); - let filename = file_name.replace("{}", format!("{}{}", date_prefix, now.format("%Y%m%d_%H%M%S").to_string().as_str()).as_str()); - Some(std::path::PathBuf::from(filename)) -} - -pub(crate) fn get_file_path(wd: &String, path: Option) -> Option { - match path { - Some(p) => { - if p.is_relative() { - let pb = PathBuf::from(wd); - match pb.join(&p).absolutize() { - Ok(os) => Some(PathBuf::from(os)), - Err(e) => { - error!("path is not relative {:?}", e); - Some(p) - } - } - } else { - Some(p) - } - } - None => None - } -} - - pub(crate) fn get_client_request(input: &ConfigInput, url: url::Url, custom_headers: Option<&HashMap<&str, &[u8]>>) -> reqwest::RequestBuilder { let mut request = reqwest::Client::new().get(url); let headers = get_request_headers(&input.headers, custom_headers); @@ -297,25 +159,4 @@ async fn download_text_content(input: &ConfigInput, url: url::Url, persist_filep } Err(e) => Err(e.to_string()) } -} - -pub(crate) fn bytes_to_megabytes(bytes: u64) -> u64 { - bytes / 1_048_576 -} - -pub(crate) fn add_prefix_to_filename(path: &Path, prefix: &str, ext: Option<&str>) -> PathBuf { - let file_name = path.file_name().unwrap_or_default(); - let new_file_name = format!("{}{}", prefix, file_name.to_string_lossy()); - let result = path.with_file_name(new_file_name); - match ext { - None => result, - Some(extension) => result.with_extension(extension) - } -} - -pub(crate) fn path_exists(file_path: &Path) -> bool { - if let Ok(metadata) = fs::metadata(file_path) { - return metadata.is_file(); - } - false -} +} \ No newline at end of file diff --git a/src/utils/string_utils.rs b/src/utils/string_utils.rs new file mode 100644 index 000000000..53136c400 --- /dev/null +++ b/src/utils/string_utils.rs @@ -0,0 +1,14 @@ +// other implementations like calculating text_distance on all titles took too much time +// we keep it now as simple as possible and less memory intensive. +pub (crate) fn get_title_group(text: &str) -> String { + let alphabetic_only: String = text.chars().map(|c| if c.is_alphanumeric() { c } else { ' ' }).collect(); + let parts = alphabetic_only.split_whitespace(); + let mut combination = "".to_string(); + for p in parts.into_iter() { + combination = format!("{} {}", combination, p).trim().to_string(); + if combination.len() > 2 { + return combination; + } + } + text.to_string() +}