diff --git a/Cargo.lock b/Cargo.lock index 5e14f4d7f..60bccee6a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1156,7 +1156,7 @@ dependencies = [ [[package]] name = "frontend" -version = "3.2.29" +version = "3.2.30" dependencies = [ "anyhow", "base64", @@ -1168,9 +1168,9 @@ dependencies = [ "futures-signals", "gloo-render 0.2.0", "gloo-storage 0.3.0", - "gloo-timers 0.3.0", + "gloo-timers 0.2.6", "gloo-utils 0.2.0", - "implicit-clone 0.6.0", + "implicit-clone", "js-sys", "log", "prost", @@ -1691,6 +1691,8 @@ version = "0.2.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9b995a66bb87bebce9a0f4a95aed01daca4872c050bfcb21653361c03bc35e5c" dependencies = [ + "futures-channel", + "futures-core", "js-sys", "wasm-bindgen", ] @@ -2234,15 +2236,6 @@ dependencies = [ "indexmap", ] -[[package]] -name = "implicit-clone" -version = "0.6.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1689b939ee35e3a075b0834b5672efd43aec8a6e81a1c6002b76a5ca2f211ae0" -dependencies = [ - "implicit-clone-derive", -] - [[package]] name = "implicit-clone-derive" version = "0.1.2" @@ -3968,7 +3961,7 @@ dependencies = [ [[package]] name = "shared" -version = "3.2.29" +version = "3.2.30" dependencies = [ "base64", "bitflags 2.10.0", @@ -4542,7 +4535,7 @@ checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b" [[package]] name = "tuliprox" -version = "3.2.29" +version = "3.2.30" dependencies = [ "arc-swap", "async-compression", @@ -5326,7 +5319,7 @@ dependencies = [ "console_error_panic_hook", "futures", "gloo 0.10.0", - "implicit-clone 0.4.9", + "implicit-clone", "indexmap", "js-sys", "prokio", diff --git a/backend/Cargo.toml b/backend/Cargo.toml index fb1550222..b6f197655 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "tuliprox" -version = "3.2.29" +version = "3.2.30" edition = "2021" rust-version = "1.87.0" diff --git a/backend/src/api/config_file.rs b/backend/src/api/config_file.rs index 1c0e530fa..2717b9a7a 100644 --- a/backend/src/api/config_file.rs +++ b/backend/src/api/config_file.rs @@ -65,8 +65,8 @@ impl ConfigFile { let mapping_changed = paths.mapping_file_path.as_ref() != config_dto.mapping_path.as_ref(); let mut config: Config = Config::from(config_dto); config.prepare(paths.config_path.as_str())?; - info!("Loaded config file {config_file}"); update_app_state_config(app_state, config).await?; + info!("Loaded config file {config_file}"); if mapping_changed { Self::load_mapping(app_state)?; } @@ -82,8 +82,8 @@ impl ConfigFile { }; prepare_sources_batch(&mut sources_dto, true).await?; let sources: SourcesConfig = SourcesConfig::try_from(sources_dto)?; - info!("Loaded sources file {sources_file}"); update_app_state_sources(app_state, sources).await?; + info!("Loaded sources file {sources_file}"); // mappings are not stored, so we need to reload and apply them if sources change. Self::load_mapping(app_state) } diff --git a/backend/src/repository/m3u_repository.rs b/backend/src/repository/m3u_repository.rs index 4340630e1..cd82c064c 100644 --- a/backend/src/repository/m3u_repository.rs +++ b/backend/src/repository/m3u_repository.rs @@ -6,18 +6,19 @@ use crate::repository::m3u_playlist_iterator::M3uPlaylistM3uTextIterator; use crate::repository::storage::get_target_storage_path; use crate::repository::storage_const; use crate::utils; +use crate::utils::{async_file_writer, IO_BUFFER_SIZE}; +use indexmap::IndexMap; +use log::error; use shared::error::{create_tuliprox_error, info_err, string_to_io_error}; use shared::error::{str_to_io_error, TuliproxError, TuliproxErrorKind}; -use shared::model::PlaylistItemType; use shared::model::{M3uPlaylistItem, PlaylistGroup}; +use shared::model::{PlaylistItem, PlaylistItemType, XtreamCluster}; use std::io::Error; use std::path::{Path, PathBuf}; use std::sync::Arc; -use log::error; use tokio::fs; -use tokio::io::{AsyncWriteExt}; +use tokio::io::AsyncWriteExt; use tokio::task; -use crate::utils::{async_file_writer, IO_BUFFER_SIZE}; macro_rules! cant_write_result { ($path:expr, $err:expr) => { @@ -146,7 +147,7 @@ pub async fn m3u_get_item_for_stream_id(stream_id: u32, app_state: &AppState, ta let target_path = get_target_storage_path(&cfg.config.load(), target.name.as_str()).ok_or_else(|| string_to_io_error(format!("Could not find path for target {}", &target.name)))?; let m3u_path = m3u_get_file_path(&target_path); let _file_lock = cfg.file_locks.read_lock(&m3u_path).await; - + let mut query = BPlusTreeQuery::::try_new(&m3u_path)?; match query.query(&stream_id) { Ok(Some(item)) => Ok(item), @@ -172,4 +173,45 @@ pub async fn iter_raw_m3u_playlist(config: &AppConfig, target: &ConfigTarget) -> } Err(_) => None } -} \ No newline at end of file +} + +pub async fn persist_input_m3u_playlist(app_config: &Arc, m3u_path: &Path, playlist: &[PlaylistGroup]) -> Result<(), TuliproxError> { + let _file_lock = app_config.file_locks.read_lock(m3u_path).await; + let mut tree = BPlusTree::new(); + for pg in playlist { + for item in &pg.channels { + let m3u = M3uPlaylistItem::from(item); + tree.insert(m3u.provider_id.clone(), m3u); + } + } + tree.store(m3u_path).map_err(|err| cant_write_result!(&m3u_path, err))?; + Ok(()) +} + +pub async fn load_input_m3u_playlist(app_config: &Arc, m3u_path: &Path) -> Result, TuliproxError> { + let mut groups: IndexMap = IndexMap::new(); + + if m3u_path.exists() { + // Load Items + let _file_lock = app_config.file_locks.read_lock(m3u_path).await; + if let Ok(mut query) = BPlusTreeQuery::::try_new(m3u_path) { + let mut group_cnt = 0; + for (_, ref item) in query.iter() { + let cat_id = item.group.clone(); + groups.entry(cat_id) + .or_insert_with(|| { + group_cnt += 1; + PlaylistGroup { + id: group_cnt, + title: item.group.clone(), + channels: Vec::new(), + xtream_cluster: XtreamCluster::Live, + } + }) + .channels.push(PlaylistItem::from(item)); + } + } + } + + Ok(groups.into_values().collect()) +} diff --git a/backend/src/repository/playlist_repository.rs b/backend/src/repository/playlist_repository.rs index 48943ae4e..4adde19d4 100644 --- a/backend/src/repository/playlist_repository.rs +++ b/backend/src/repository/playlist_repository.rs @@ -4,13 +4,13 @@ use crate::model::{AppConfig, ConfigInput, ConfigTarget, TargetOutput}; use crate::processing::processor::playlist::apply_filter_to_playlist; use crate::repository::bplustree::{BPlusTree, BPlusTreeQuery}; use crate::repository::epg_repository::epg_write; -use crate::repository::m3u_repository::{m3u_get_file_path, m3u_write_playlist}; +use crate::repository::m3u_repository::{m3u_get_file_path, m3u_write_playlist, persist_input_m3u_playlist}; use crate::repository::storage::{ensure_target_storage_path, get_input_storage_path, get_target_id_mapping_file, get_target_storage_path}; +use crate::repository::storage_const::FILE_SUFFIX_DB; use crate::repository::strm_repository::write_strm_playlist; use crate::repository::target_id_mapping::{TargetIdMapping, VirtualIdRecord}; use crate::repository::xtream_repository::{load_input_xtream_playlist, persist_input_xtream_playlist, xtream_get_file_path, xtream_get_storage_path, xtream_write_playlist}; use crate::utils; -use crate::utils::json_write_documents_to_file; use log::info; use shared::create_tuliprox_error; use shared::error::TuliproxError; @@ -19,7 +19,7 @@ use shared::model::xtream_const::XTREAM_CLUSTER; use shared::model::{InputType, M3uPlaylistItem, PlaylistEntry, PlaylistGroup, PlaylistItem, PlaylistItemHeader, PlaylistItemType, StreamProperties, XtreamCluster, XtreamPlaylistItem}; use shared::utils::{is_dash_url, is_hls_url}; use std::collections::HashMap; -use std::path::Path; +use std::path::{Path, PathBuf}; use std::sync::Arc; struct LocalEpisodeKey { @@ -337,7 +337,7 @@ pub async fn load_target_into_memory_cache(app_state: &AppState, target: &Arc, input: &ConfigInput, mut playlist: Vec) -> (Vec, Option) { playlist.iter_mut().for_each(PlaylistGroup::on_load); - let (result, err) = match input.input_type { + 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) { @@ -350,8 +350,6 @@ pub async fn persist_input_playlist(app_config: &Arc, input: &ConfigI } _ => { - // TODO what is written and why not as BPlusTree? - // Persist M3U/Other types let working_dir = &app_config.config.load().working_dir; let storage_path = match get_input_storage_path(&input.name, working_dir) { @@ -360,23 +358,13 @@ pub async fn persist_input_playlist(app_config: &Arc, input: &ConfigI return (playlist, Some(create_tuliprox_error!( TuliproxErrorKind::Info, "Error creating input storage directory for input '{}' failed: {err}", input.name))); } }; - let sanitized_input_name: String = input.name.chars() - .map(|c| if c.is_alphanumeric() { c } else { '_' }) - .collect(); - let file_path = storage_path.join(format!("{sanitized_input_name}_playlist.json")); - let file_path_clone = file_path.clone(); - let playlist_clone = playlist.clone(); - match tokio::task::spawn_blocking(move || { - json_write_documents_to_file(&file_path_clone, &playlist_clone) - }).await { - Ok(Ok(())) => (playlist, None), - Ok(Err(e)) => (playlist, Some(info_err!(format!("Failed to persist input playlist: {e}")))), - Err(e) => (playlist, Some(info_err!(format!("Failed to persist input playlist: {e}")))), + let file_path = get_input_playlist_file_path(&storage_path, input.name.as_str()); + if let Err(err) = persist_input_m3u_playlist(app_config, &file_path, &playlist).await { + return (playlist, Some(err)); } + (playlist, None) } - }; - - (result, err) + } } pub async fn load_input_playlist(app_config: &Arc, input: &ConfigInput, clusters: Option<&[XtreamCluster]>) -> Result, TuliproxError> { @@ -395,8 +383,8 @@ pub async fn load_input_playlist(app_config: &Arc, input: &ConfigInpu load_input_xtream_playlist(app_config, &storage_path, clusters_to_load).await } _ => { - // Load JSON for M3U - let file_path = storage_path.join("playlist.json"); + // Load M3U + let file_path = get_input_playlist_file_path(&storage_path, input.name.as_str()); if file_path.exists() { let content = tokio::fs::read_to_string(&file_path).await .map_err(|e| info_err!(format!("Failed to read input playlist cache: {e}")))?; @@ -408,4 +396,12 @@ pub async fn load_input_playlist(app_config: &Arc, input: &ConfigInpu } } } +} + + +fn get_input_playlist_file_path(storage_path: &Path, input_name: &str) -> PathBuf { + let sanitized_input_name: String = input_name.chars() + .map(|c| if c.is_alphanumeric() { c } else { '_' }) + .collect(); + storage_path.join(format!("{sanitized_input_name}_playlist.{FILE_SUFFIX_DB}")) } \ No newline at end of file diff --git a/frontend/Cargo.toml b/frontend/Cargo.toml index 91a094d18..7f622cdff 100644 --- a/frontend/Cargo.toml +++ b/frontend/Cargo.toml @@ -1,11 +1,11 @@ [package] name = "frontend" -version = "3.2.29" +version = "3.2.30" edition = "2021" rust-version = "1.87.0" [dependencies] -shared = { version = "3.2.29", path = "../shared" } +shared = { version = "3.2.30", path = "../shared" } chrono = "0" yew = "0.21" yew-router = "0.18" diff --git a/shared/Cargo.toml b/shared/Cargo.toml index d4ca5e701..008d6fc27 100644 --- a/shared/Cargo.toml +++ b/shared/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "shared" -version = "3.2.29" +version = "3.2.30" edition = "2021" rust-version = "1.87.0" @@ -30,4 +30,4 @@ serde-saphyr = "0.0.12" [target.'cfg(target_arch = "wasm32")'.dependencies] js-sys = "0.3.83" -getrandom = { version = "0.3", features = ["wasm_js"] } \ No newline at end of file +getrandom = { version = "0.3", features = ["wasm_js"] } diff --git a/shared/src/model/playlist.rs b/shared/src/model/playlist.rs index 5137d34d0..360e56ee6 100644 --- a/shared/src/model/playlist.rs +++ b/shared/src/model/playlist.rs @@ -1337,6 +1337,36 @@ impl From<&XtreamPlaylistItem> for PlaylistItem { } } +impl From<&M3uPlaylistItem> for PlaylistItem { + fn from(item: &M3uPlaylistItem) -> Self { + let header = PlaylistItemHeader { + uuid: item.get_uuid(), + virtual_id: item.virtual_id, + id: item.provider_id.clone(), + name: item.name.clone(), + title: item.title.clone(), + logo: item.logo.clone(), + logo_small: item.logo_small.clone(), + group: item.group.clone(), + parent_code: item.parent_code.clone(), + rec: item.rec.clone(), + url: item.url.clone(), + epg_channel_id: item.epg_channel_id.clone(), + xtream_cluster: XtreamCluster::Live, // TODO based on file ending + item_type: item.item_type, + category_id: 0, + input_name: item.input_name.clone(), + chno: item.chno, + audio_track: String::new(), + time_shift: String::new(), + additional_properties: None, + }; + + PlaylistItem { + header + } + } +} impl PlaylistEntry for PlaylistItem {