From 2fa7a01c53d47e1be0676f45796bb0c5d752e2a7 Mon Sep 17 00:00:00 2001 From: euzu Date: Sat, 3 Jan 2026 19:43:43 +0100 Subject: [PATCH] To align input definitions with the SourceEditor, inputs are now defined globally in the inputs section of the config file. Each source can reference one or more inputs by their name in the inputs attribute. --- Cargo.lock | 23 +++----- backend/Cargo.toml | 2 +- backend/src/api/config_file.rs | 4 +- backend/src/repository/m3u_repository.rs | 54 ++++++++++++++++--- backend/src/repository/playlist_repository.rs | 42 +++++++-------- frontend/Cargo.toml | 4 +- shared/Cargo.toml | 4 +- shared/src/model/playlist.rs | 30 +++++++++++ 8 files changed, 112 insertions(+), 51 deletions(-) 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 {