diff --git a/Cargo.toml b/Cargo.toml index 44e0abdd3..ddf8d2d90 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -5,7 +5,7 @@ resolver = "2" [profile.release] debug = false opt-level = 'z' # Optimize for size. -lto = true # Enable Link Time Optimization +lto = "fat" # Enable Link Time Optimization codegen-units = 1 # Reduce number of codegen units to increase optimizations. panic = 'abort' # Abort on panic strip = true diff --git a/backend/src/processing/input_cache.rs b/backend/src/processing/input_cache.rs index d6a550b21..5c354bcc7 100644 --- a/backend/src/processing/input_cache.rs +++ b/backend/src/processing/input_cache.rs @@ -27,8 +27,8 @@ pub struct InputStatus { pub clusters: HashMap, } -pub fn resolve_input_storage_path(working_dir: &str, input_name: &str) -> PathBuf { - if let Ok(path) = get_input_storage_path(input_name, working_dir) { path } else { +pub async fn resolve_input_storage_path(working_dir: &str, input_name: &str) -> PathBuf { + if let Ok(path) = get_input_storage_path(input_name, working_dir).await { path } else { build_input_storage_path(input_name, working_dir) } } diff --git a/backend/src/processing/processor/playlist.rs b/backend/src/processing/processor/playlist.rs index feafa8545..94e307321 100644 --- a/backend/src/processing/processor/playlist.rs +++ b/backend/src/processing/processor/playlist.rs @@ -282,7 +282,7 @@ async fn playlist_download_from_input(client: &reqwest::Client, app_config: &Arc let working_dir = &config.working_dir; // Check Status - let storage_path = input_cache::resolve_input_storage_path(working_dir, &input.name); + let storage_path = input_cache::resolve_input_storage_path(working_dir, &input.name).await; let mut status = input_cache::load_input_status(&storage_path); let cache_duration = input.cache_duration_seconds; diff --git a/backend/src/processing/processor/xtream_series.rs b/backend/src/processing/processor/xtream_series.rs index 618ddd515..e55c80dd4 100644 --- a/backend/src/processing/processor/xtream_series.rs +++ b/backend/src/processing/processor/xtream_series.rs @@ -26,7 +26,7 @@ async fn playlist_resolve_series_info(app_config: &Arc, client: &reqw let input = fpl.input; let working_dir = &app_config.config.load().working_dir; - let storage_path = match get_input_storage_path(&input.name, working_dir) { + let storage_path = match get_input_storage_path(&input.name, working_dir).await { Ok(storage_path) => storage_path, Err(err) => { error!("Can't resolve series info, input storage directory for input '{}' failed: {err}", input.name); diff --git a/backend/src/processing/processor/xtream_vod.rs b/backend/src/processing/processor/xtream_vod.rs index b5dd56c67..ca08df03e 100644 --- a/backend/src/processing/processor/xtream_vod.rs +++ b/backend/src/processing/processor/xtream_vod.rs @@ -27,7 +27,7 @@ pub async fn playlist_resolve_vod(app_config: &Arc, let input = fpl.input; let working_dir = &app_config.config.load().working_dir; - let storage_path = match get_input_storage_path(&input.name, working_dir) { + let storage_path = match get_input_storage_path(&input.name, working_dir).await { Ok(storage_path) => storage_path, Err(err) => { error!("Can't resolve vod, input storage directory for input '{}' failed: {err}", input.name); diff --git a/backend/src/repository/m3u_playlist_iterator.rs b/backend/src/repository/m3u_playlist_iterator.rs index 389c547fd..b1ab6e54b 100644 --- a/backend/src/repository/m3u_playlist_iterator.rs +++ b/backend/src/repository/m3u_playlist_iterator.rs @@ -12,7 +12,7 @@ use crate::utils::FileReadGuard; use std::collections::HashSet; use std::iter::Peekable; use log::error; -use shared::utils::Internable; +use shared::utils::{extract_extension_from_url, Internable}; #[allow(clippy::struct_excessive_bools)] pub struct M3uPlaylistIterator { @@ -41,7 +41,7 @@ impl M3uPlaylistIterator { let m3u_output = target.get_m3u_output().ok_or_else(|| info_err!("Unexpected failure, missing m3u target output for target {}", target.name))?; let config = cfg.config.load(); - let target_path = ensure_target_storage_path(&config, target.name.as_str())?; + let target_path = ensure_target_storage_path(&config, target.name.as_str()).await?; let m3u_path = m3u_get_file_path_for_db(&target_path); let file_lock = cfg.file_locks.read_lock(&m3u_path).await; @@ -107,7 +107,7 @@ impl M3uPlaylistIterator { + 32; // separators and id if typed { cap += stream_type.len() + 1; } - if typed { + let rewritten_url = if typed { shared::concat_string!( cap = cap; &self.base_url, "/", prefix_path, "/", stream_type, "/", @@ -119,7 +119,9 @@ impl M3uPlaylistIterator { &self.base_url, "/", prefix_path, "/", &self.username, "/", &self.password, "/", &m3u_pli.virtual_id.to_string() ) - } + }; + + extract_extension_from_url(&m3u_pli.url).map(|ext| shared::concat_string!(&rewritten_url, ext)).unwrap_or(rewritten_url) } fn get_stream_url(&self, m3u_pli: &M3uPlaylistItem, typed: bool) -> String { diff --git a/backend/src/repository/m3u_repository.rs b/backend/src/repository/m3u_repository.rs index 3adb0ba9e..5f2d2ba23 100644 --- a/backend/src/repository/m3u_repository.rs +++ b/backend/src/repository/m3u_repository.rs @@ -5,13 +5,13 @@ use crate::repository::bplustree::{BPlusTree, BPlusTreeQuery}; use crate::repository::m3u_playlist_iterator::M3uPlaylistM3uTextIterator; use crate::repository::playlist_repository::get_input_m3u_playlist_file_path; use crate::repository::storage::{get_input_storage_path, get_target_storage_path}; -use crate::repository::storage_const; +use crate::repository::{storage_const}; use crate::repository::xtream_repository::CategoryKey; use crate::utils; use crate::utils::{async_file_writer, file_exists_async, FileReadGuard, IO_BUFFER_SIZE}; use indexmap::IndexMap; use log::error; -use shared::concat_string; +use shared::{concat_string, notify_err_res}; use shared::error::{notify_err, str_to_io_error, string_to_io_error, TuliproxError}; use shared::model::{M3uPlaylistItem, PlaylistGroup}; use shared::model::{PlaylistItem, PlaylistItemType, XtreamCluster}; @@ -29,15 +29,6 @@ macro_rules! cant_write_result { } } -pub fn m3u_get_file_path_for_db(target_path: &Path) -> PathBuf { - target_path.join(PathBuf::from(concat_string!(storage_const::FILE_M3U, ".", storage_const::FILE_SUFFIX_DB))) -} - -pub fn m3u_get_epg_file_path_for_target(target_path: &Path) -> PathBuf { - let path = target_path.join(PathBuf::from(concat_string!(storage_const::FILE_M3U, ".", storage_const::FILE_SUFFIX_DB))); - utils::add_prefix_to_filename(&path, "epg_", Some(storage_const::FILE_SUFFIX_DB)) -} - macro_rules! await_playlist_write { ($expr:expr, $fmt:literal $(, $args:expr)* ) => {{ $expr.await.map_err(|err| { @@ -46,6 +37,35 @@ macro_rules! await_playlist_write { }}; } +pub fn m3u_get_file_path_for_db(target_path: &Path) -> PathBuf { + target_path.join(storage_const::PATH_M3U).join(concat_string!(storage_const::FILE_M3U, ".", storage_const::FILE_SUFFIX_DB)) +} + +pub fn m3u_get_epg_file_path_for_target(target_path: &Path) -> PathBuf { + let path = target_path.join(storage_const::PATH_M3U).join(concat_string!(storage_const::FILE_M3U, ".", storage_const::FILE_SUFFIX_DB)); + utils::add_prefix_to_filename(&path, "epg_", Some(storage_const::FILE_SUFFIX_DB)) +} + +pub fn m3u_get_storage_path(cfg: &Config, target_name: &str) -> Option { + get_target_storage_path(cfg, target_name).map(|target_path| target_path.join(PathBuf::from(storage_const::PATH_M3U))) +} + +pub async fn ensure_m3u_storage_path(cfg: &Config, target_name: &str) -> Result { + if let Some(path) = m3u_get_storage_path(cfg, target_name) { + if tokio::fs::create_dir_all(&path).await.is_err() { + let msg = format!( + "Failed to save m3u data, can't create directory {}", + &path.display() + ); + return notify_err_res!("{msg}"); + } + Ok(path) + } else { + let msg = format!("Failed to save m3u data, can't create directory for target {target_name}"); + notify_err_res!("{msg}") + } +} + async fn persist_m3u_playlist_as_text( cfg: &Config, target: &ConfigTarget, @@ -90,6 +110,9 @@ pub async fn m3u_write_playlist( return Ok(()); } + let config = cfg.config.load(); + let _m3u_path = ensure_m3u_storage_path(&config, target.name.as_str()).await?; + let m3u_path = m3u_get_file_path_for_db(target_path); let m3u_playlist = Arc::new( new_playlist @@ -169,7 +192,7 @@ pub async fn iter_raw_m3u_target_playlist(config: &AppConfig, target: &ConfigTar pub async fn iter_raw_m3u_input_playlist(app_config: &AppConfig, input: &ConfigInput, cluster: Option) -> Option<(FileReadGuard, Box + Send>)> { let working_dir = &app_config.config.load().working_dir; - let storage_path = get_input_storage_path(&input.name, working_dir).ok()?; + let storage_path = get_input_storage_path(&input.name, working_dir).await.ok()?; let m3u_path = get_input_m3u_playlist_file_path(&storage_path, &input.name); iter_raw_m3u_playlist::>(app_config, &m3u_path, cluster).await diff --git a/backend/src/repository/playlist_repository.rs b/backend/src/repository/playlist_repository.rs index 15be36a98..5246186e4 100644 --- a/backend/src/repository/playlist_repository.rs +++ b/backend/src/repository/playlist_repository.rs @@ -37,7 +37,7 @@ pub async fn persist_playlist(app_config: &Arc, playlist: &mut [Playl target: &ConfigTarget, playlist_state: Option<&Arc>) -> Result<(), Vec> { let mut errors = vec![]; let config = &app_config.config.load(); - let target_path = match ensure_target_storage_path(config, &target.name) { + let target_path = match ensure_target_storage_path(config, &target.name).await { Ok(path) => path, Err(err) => return Err(vec![err]), }; @@ -355,7 +355,7 @@ pub async fn persist_input_playlist(app_config: &Arc, input: &ConfigI match input.input_type { InputType::Xtream | InputType::XtreamBatch => { let working_dir = &app_config.config.load().working_dir; - let storage_path = match get_input_storage_path(&input.name, working_dir) { + let storage_path = match get_input_storage_path(&input.name, working_dir).await { Ok(storage_path) => storage_path, Err(err) => { return (playlist, Some(info_err!("Error creating input storage directory for input '{}' failed: {err}", input.name))); @@ -367,7 +367,7 @@ pub async fn persist_input_playlist(app_config: &Arc, input: &ConfigI InputType::M3u | InputType::M3uBatch => { // Persist M3U let working_dir = &app_config.config.load().working_dir; - let storage_path = match get_input_storage_path(&input.name, working_dir) { + let storage_path = match get_input_storage_path(&input.name, working_dir).await { Ok(storage_path) => storage_path, Err(err) => { return (playlist, Some(info_err!("Error creating input storage directory for input '{}' failed: {err}", input.name))); @@ -382,7 +382,7 @@ pub async fn persist_input_playlist(app_config: &Arc, input: &ConfigI InputType::Library => { // Persist local library playlist let working_dir = &app_config.config.load().working_dir; - let storage_path = match get_input_storage_path(&input.name, working_dir) { + let storage_path = match get_input_storage_path(&input.name, working_dir).await { Ok(storage_path) => storage_path, Err(err) => { return (playlist, Some(info_err!("Error creating input storage directory for input '{}' failed: {err}", input.name))); @@ -400,7 +400,7 @@ pub async fn persist_input_playlist(app_config: &Arc, input: &ConfigI pub async fn load_input_playlist(ctx: &PlaylistProcessingContext, input: &ConfigInput, clusters: Option<&[XtreamCluster]>) -> Result, TuliproxError> { let app_config = &ctx.config; let working_dir = &app_config.config.load().working_dir; - let storage_path = get_input_storage_path(&input.name, working_dir) + let storage_path = get_input_storage_path(&input.name, working_dir).await .map_err(|e| info_err!("Error getting input path: {e}"))?; let disk_based_processing = app_config.config.load().disk_based_processing; diff --git a/backend/src/repository/storage.rs b/backend/src/repository/storage.rs index 189fd3f26..79cd67768 100644 --- a/backend/src/repository/storage.rs +++ b/backend/src/repository/storage.rs @@ -11,9 +11,9 @@ pub(in crate::repository) fn get_target_id_mapping_file(target_path: &Path) -> P target_path.join(storage_const::FILE_ID_MAPPING) } -pub fn ensure_target_storage_path(cfg: &Config, target_name: &str) -> Result { +pub async fn ensure_target_storage_path(cfg: &Config, target_name: &str) -> Result { if let Some(path) = get_target_storage_path(cfg, target_name) { - if std::fs::create_dir_all(&path).is_err() { + if tokio::fs::create_dir_all(&path).await.is_err() { let msg = format!("Failed to save target data, can't create directory {}", path.display()); return notify_err_res!("{msg}"); } @@ -40,14 +40,14 @@ pub fn build_input_storage_path(input_name: &str, working_dir: &str) -> PathBuf Path::new(working_dir).join(name) } -pub fn get_input_storage_path(input_name: &str, working_dir: &str) -> std::io::Result { +pub async fn get_input_storage_path(input_name: &str, working_dir: &str) -> std::io::Result { let path = build_input_storage_path(input_name, working_dir); // Create the directory and return the path or propagate the error - std::fs::create_dir_all(&path).map(|()| path) + tokio::fs::create_dir_all(&path).await.map(|()| path) } -pub fn ensure_input_storage_path(cfg: &Config, input_name: &str) -> Result { - get_input_storage_path(input_name, &cfg.working_dir) +pub async fn ensure_input_storage_path(cfg: &Config, input_name: &str) -> Result { + get_input_storage_path(input_name, &cfg.working_dir).await .map_err(|err| { notify_err!("Failed to save input data, can't create directory for input {input_name}: {err}") }) diff --git a/backend/src/repository/storage_const.rs b/backend/src/repository/storage_const.rs index ba143fd49..4693aaeb3 100644 --- a/backend/src/repository/storage_const.rs +++ b/backend/src/repository/storage_const.rs @@ -4,6 +4,8 @@ pub(in crate::repository) const FILE_SUFFIX_INDEX: &str = "idx"; pub(in crate::repository) const FILE_ID_MAPPING: &str = "id_mapping.db"; pub(in crate::repository) const FILE_STRM: &str = "strm"; pub(in crate::repository) const FILE_M3U: &str = "m3u"; +pub(in crate::repository) const PATH_M3U: &str = "m3u"; + pub const M3U_STREAM_PATH: &str = "m3u-stream"; pub const M3U_RESOURCE_PATH: &str = "resource/m3u"; pub const EPG_RESOURCE_PATH: &str = "resource/epg"; diff --git a/backend/src/repository/strm_repository.rs b/backend/src/repository/strm_repository.rs index b3c890c34..13e2e28b0 100644 --- a/backend/src/repository/strm_repository.rs +++ b/backend/src/repository/strm_repository.rs @@ -722,7 +722,7 @@ pub async fn write_strm_playlist( let normalized_dir = normalize_string_path(&target_output.directory); let strm_file_prefix = hash_string_as_hex(&normalized_dir); let strm_index_path = - strm_get_file_paths(&strm_file_prefix, &ensure_target_storage_path(&config, target.name.as_str())?); + strm_get_file_paths(&strm_file_prefix, &ensure_target_storage_path(&config, target.name.as_str()).await?); let existing_strm = { let _file_lock = app_config .file_locks diff --git a/backend/src/repository/xtream_repository.rs b/backend/src/repository/xtream_repository.rs index 34c1e05c5..931d8596a 100644 --- a/backend/src/repository/xtream_repository.rs +++ b/backend/src/repository/xtream_repository.rs @@ -58,9 +58,9 @@ pub fn get_series_cat_collection_path(path: &Path) -> PathBuf { get_collection_path(path, storage_const::COL_CAT_SERIES) } -pub fn ensure_xtream_storage_path(cfg: &Config, target_name: &str) -> Result { +pub async fn ensure_xtream_storage_path(cfg: &Config, target_name: &str) -> Result { if let Some(path) = xtream_get_storage_path(cfg, target_name) { - if std::fs::create_dir_all(&path).is_err() { + if tokio::fs::create_dir_all(&path).await.is_err() { let msg = format!( "Failed to save xtream data, can't create directory {}", &path.display() @@ -118,7 +118,7 @@ pub async fn write_playlist_item_update( ) -> Result<(), TuliproxError> { let storage_path = { let config = app_config.config.load(); - ensure_xtream_storage_path(&config, target_name)? + ensure_xtream_storage_path(&config, target_name).await? }; let xtream_path = xtream_get_file_path(&storage_path, pli.xtream_cluster); { @@ -143,7 +143,7 @@ pub async fn write_playlist_batch_item_upsert( ) -> Result<(), TuliproxError> { let storage_path = { let config = app_config.config.load(); - ensure_xtream_storage_path(&config, target_name)? + ensure_xtream_storage_path(&config, target_name).await? }; let xtream_path = xtream_get_file_path(&storage_path, xtream_cluster); { @@ -245,7 +245,7 @@ pub async fn xtream_write_playlist( ) -> Result<(), TuliproxError> { let path = { let config = app_cfg.config.load(); - ensure_xtream_storage_path(&config, target.name.as_str())? + ensure_xtream_storage_path(&config, target.name.as_str()).await? }; let mut errors = Vec::new(); let mut cat_live_col = Vec::with_capacity(1_000); @@ -578,7 +578,7 @@ pub async fn iter_raw_xtream_target_playlist(app_config: &AppConfig, target: &Co pub async fn iter_raw_xtream_input_playlist(app_config: &AppConfig, input: &ConfigInput, cluster: XtreamCluster) -> Option<(FileReadGuard, Box + Send>)> { let config = app_config.config.load(); let working_dir = &config.working_dir; - let storage_path = get_input_storage_path(&input.name, working_dir).ok()?; + let storage_path = get_input_storage_path(&input.name, working_dir).await.ok()?; let xtream_path = xtream_get_file_path(&storage_path, cluster); iter_raw_xtream_playlist::(app_config, &xtream_path).await diff --git a/backend/src/utils/network/epg.rs b/backend/src/utils/network/epg.rs index 601439b7e..6e9b39f24 100644 --- a/backend/src/utils/network/epg.rs +++ b/backend/src/utils/network/epg.rs @@ -11,7 +11,7 @@ use shared::error::{info_err, TuliproxError}; use shared::utils::{sanitize_sensitive_info, short_hash}; use std::path::PathBuf; -pub fn get_input_raw_epg_file_path(url: &str, input: &ConfigInput, working_dir: &str) -> std::io::Result { +pub async fn get_input_raw_epg_file_path(url: &str, input: &ConfigInput, working_dir: &str) -> std::io::Result { let file_prefix = short_hash(url); if let Some(persist_path) = input.persist.as_deref() { @@ -23,7 +23,7 @@ pub fn get_input_raw_epg_file_path(url: &str, input: &ConfigInput, working_dir: } } - let download_path = get_input_storage_path(&input.name, working_dir)?; + let download_path = get_input_storage_path(&input.name, working_dir).await?; Ok(download_path.join(format!("{}_{}", file_prefix, storage_const::FILE_EPG))) } @@ -32,7 +32,7 @@ async fn download_epg_file(url: &str, ctx: &PlaylistProcessingContext, headers: Option<&reqwest::header::HeaderMap>, working_dir: &str) -> Result { debug!("Getting epg file path for url: {}", sanitize_sensitive_info(url)); - let persist_file_path = get_input_raw_epg_file_path(url, input, working_dir).map_err(|e| info_err!("Could not access epg file download directory: {}", e))?; + let persist_file_path = get_input_raw_epg_file_path(url, input, working_dir).await.map_err(|e| info_err!("Could not access epg file download directory: {}", e))?; if input.cache_duration_seconds > 0 { if let Ok(metadata) = tokio::fs::metadata(&persist_file_path).await { diff --git a/backend/src/utils/network/xtream.rs b/backend/src/utils/network/xtream.rs index da22720cb..415149c96 100644 --- a/backend/src/utils/network/xtream.rs +++ b/backend/src/utils/network/xtream.rs @@ -100,7 +100,7 @@ pub async fn get_xtream_stream_info(client: &reqwest::Client, XtreamCluster::Live => {} XtreamCluster::Video => { let working_dir = &app_config.config.load().working_dir; - if let Ok(storage_path) = get_input_storage_path(&input.name, working_dir) { + if let Ok(storage_path) = get_input_storage_path(&input.name, working_dir).await { match serde_json::from_str::(&content) { Ok(info) => { // parse downloaded info into StreamProperties @@ -141,7 +141,7 @@ pub async fn get_xtream_stream_info(client: &reqwest::Client, // parse series info let series_stream_props = SeriesStreamProperties::from_info(&info, pli); - if let Ok(storage_path) = get_input_storage_path(&input.name, working_dir) { + if let Ok(storage_path) = get_input_storage_path(&input.name, working_dir).await { // update input db if let Err(err) = persists_input_series_info(app_config, &storage_path, cluster, &input.name, provider_id, &series_stream_props).await { error!("Failed to persist series info for input {}: {err}", &input.name); @@ -511,7 +511,7 @@ async fn process_xtream_cluster_to_disk( let cfg = app_config.config.load(); // trace!("Starting process_xtream_cluster_to_disk for cluster {}", cluster); let storage_path = { - ensure_input_storage_path(&cfg, &input.name)? + ensure_input_storage_path(&cfg, &input.name).await? }; let xtream_path = xtream_get_file_path(&storage_path, cluster); diff --git a/shared/src/model/playlist.rs b/shared/src/model/playlist.rs index 622d1523b..194b345ee 100644 --- a/shared/src/model/playlist.rs +++ b/shared/src/model/playlist.rs @@ -452,7 +452,7 @@ pub struct M3uPlaylistItem { #[serde(with = "arc_str_serde")] pub input_name: Arc, pub item_type: PlaylistItemType, - #[serde(skip_serializing, default)] + #[serde(skip)] pub t_stream_url: Arc, #[serde(skip)] pub t_resource_url: Option,