From 5bb32ebc7238775bc9fe0ab0e6d5f9b1858fc003 Mon Sep 17 00:00:00 2001 From: euzu Date: Sun, 4 Jan 2026 14:55:19 +0100 Subject: [PATCH] Fixed category creation during xtream playlist save --- README.md | 1 + backend/src/processing/processor/playlist.rs | 4 +- backend/src/repository/xtream_repository.rs | 142 ++++++++++++------- backend/src/utils/file/mapping_reader.rs | 2 +- shared/src/utils/default_utils.rs | 2 +- 5 files changed, 99 insertions(+), 52 deletions(-) diff --git a/README.md b/README.md index 315f2788b..baed2387d 100644 --- a/README.md +++ b/README.md @@ -158,6 +158,7 @@ messaging: token: user: url: `optional`, default is `https://api.pushover.net/1/messages.json` +``` ### 1.4.1 Messaging Templating For `discord` and `rest` messaging, you can use [Handlebars](https://handlebarsjs.com/) templates to format the message body. diff --git a/backend/src/processing/processor/playlist.rs b/backend/src/processing/processor/playlist.rs index 9592e5eb6..4fe2768c3 100644 --- a/backend/src/processing/processor/playlist.rs +++ b/backend/src/processing/processor/playlist.rs @@ -567,7 +567,7 @@ fn execute_pipe<'a>(target: &ConfigTarget, pipe: &ProcessingPipe, fpl: &FetchedP }; if target.options.as_ref().is_some_and(|opt| opt.remove_duplicates) { for group in &mut new_fpl.playlist_groups { - // `HashSet::insert` returns true for first insert, otherweise false + // `HashSet::insert` returns true for first insert, otherwise false group.channels.retain(|item| duplicates.insert(duplicate_hash(item))); } } @@ -614,7 +614,6 @@ async fn process_playlist_for_target(app_config: &Arc, event_manager: Option>, playlist_state: Option<&Arc>, ) -> Result<(), Vec> { - let pipe = get_processing_pipe(target); debug_if_enabled!("Processing order is {}", &target.processing_order); let mut duplicates: HashSet = HashSet::new(); @@ -623,6 +622,7 @@ async fn process_playlist_for_target(app_config: &Arc, debug!("Executing processing pipes"); let broadcast_step = create_broadcast_callback(event_manager.as_ref()); + let pipe = get_processing_pipe(target); let mut step = StepMeasure::new(&target.name, broadcast_step); for provider_fpl in playlists.iter_mut() { step.broadcast("Executing transformations on '{}' playlist", &target.name); diff --git a/backend/src/repository/xtream_repository.rs b/backend/src/repository/xtream_repository.rs index 2a229c7c0..56a02fb5f 100644 --- a/backend/src/repository/xtream_repository.rs +++ b/backend/src/repository/xtream_repository.rs @@ -21,6 +21,7 @@ use shared::error::{create_tuliprox_error, create_tuliprox_error_result, info_er use shared::model::xtream_const::XTREAM_CLUSTER; use shared::model::{PlaylistGroup, PlaylistItem, PlaylistItemType, SeriesStreamProperties, StreamProperties, VideoStreamProperties, XtreamCluster, XtreamPlaylistItem}; use shared::utils::get_u32_from_serde_value; +use std::borrow::Cow; use std::collections::HashMap; use std::fs::File; use std::io::{Error, ErrorKind}; @@ -263,38 +264,29 @@ pub async fn xtream_write_playlist( let mut series_col = Vec::with_capacity(50_000); let mut vod_col = Vec::with_capacity(50_000); - // preserve category_ids - let (max_cat_id, existing_cat_ids) = load_old_category_ids(&path).await; - let mut cat_id_counter = max_cat_id; - for plg in playlist.iter_mut() { - if !plg.channels.is_empty() { - let cat_key = format!("{}{}", plg.xtream_cluster, &plg.title); - let cat_id = existing_cat_ids.get(&cat_key).unwrap_or_else(|| { - cat_id_counter += 1; - &cat_id_counter - }); - plg.id = *cat_id; - - match &plg.xtream_cluster { + let categories = create_categories(playlist, &path).await; + { + for (xtream_cluster, category) in categories { + match xtream_cluster { XtreamCluster::Live => &mut cat_live_col, XtreamCluster::Series => &mut cat_series_col, XtreamCluster::Video => &mut cat_vod_col, - }.push(json!(CategoryEntry { - category_id: *cat_id, - category_name: plg.title.clone(), - parent_id: 0 - })); + }.push(category); + } + } - for pli in &mut plg.channels { - let header = &mut pli.header; - header.category_id = *cat_id; - let col = match header.xtream_cluster { - XtreamCluster::Live => &mut live_col, - XtreamCluster::Series => &mut series_col, - XtreamCluster::Video => &mut vod_col, - }; - col.push(pli); - } + for plg in playlist.iter_mut() { + if plg.channels.is_empty() { + continue; + } + + for pli in &plg.channels { + let col = match pli.header.xtream_cluster { + XtreamCluster::Live => &mut live_col, + XtreamCluster::Series => &mut series_col, + XtreamCluster::Video => &mut vod_col, + }; + col.push(pli); } } @@ -344,6 +336,60 @@ pub async fn xtream_write_playlist( Ok(()) } +async fn create_categories(playlist: &mut [PlaylistGroup], path: &Path) -> Vec<(XtreamCluster, CategoryEntry)> { + // preserve category_ids + let (max_cat_id, existing_cat_ids) = load_old_category_ids(path).await; + let mut cat_id_counter = max_cat_id; + + let mut new_categories: IndexMap<(XtreamCluster, Cow<'_, str>), CategoryEntry> = IndexMap::new(); + + let mut last_cluster: Option = None; + let mut last_group = String::new(); + let mut last_category_id: u32 = 0; + + for plg in playlist.iter_mut() { + if plg.channels.is_empty() { + continue; + } + + for channel in &mut plg.channels { + let cluster = channel.header.xtream_cluster; + let group = channel.header.group.as_str(); + + // Fast path + if last_cluster == Some(cluster) && last_group == group { + channel.header.category_id = last_category_id; + continue; + } + + let key = (cluster, Cow::Borrowed(group)); + + let entry = new_categories.entry(key).or_insert_with(|| { + let cat_id = existing_cat_ids.get(group).copied().unwrap_or_else(|| { + cat_id_counter += 1; + cat_id_counter + }); + + CategoryEntry { + category_id: cat_id, + category_name: group.to_string(), + parent_id: 0, + } + }); + + last_cluster = Some(cluster); + last_group.clear(); + last_group.push_str(group); + last_category_id = entry.category_id; + + channel.header.category_id = last_category_id; + } + } + new_categories.into_iter() + .map(|((cluster, _group), value)| (cluster, value)) + .collect::>() +} + pub fn xtream_get_collection_path( cfg: &Config, target_name: &str, @@ -810,40 +856,40 @@ pub async fn load_input_xtream_playlist(app_config: &Arc, storage_pat XtreamCluster::Series => storage_const::COL_CAT_SERIES, }; let cat_path = get_collection_path(storage_path, cat_col_name); - + if cat_path.exists() { - if let Ok(content) = tokio::fs::read_to_string(&cat_path).await { - if let Ok(cats) = serde_json::from_str::>(&content) { - for cat in cats { - groups.insert(cat.category_id, PlaylistGroup { - id: cat.category_id, - title: cat.category_name, - channels: Vec::new(), - xtream_cluster: cluster, - }); - } - } - } + if let Ok(content) = tokio::fs::read_to_string(&cat_path).await { + if let Ok(cats) = serde_json::from_str::>(&content) { + for cat in cats { + groups.insert(cat.category_id, PlaylistGroup { + id: cat.category_id, + title: cat.category_name, + channels: Vec::new(), + xtream_cluster: cluster, + }); + } + } + } } // Load Items let _file_lock = app_config.file_locks.read_lock(&xtream_path).await; if let Ok(mut query) = BPlusTreeQuery::::try_new(&xtream_path) { for (_, ref item) in query.iter() { - let cat_id = item.category_id; - groups.entry(cat_id) + let cat_id = item.category_id; + groups.entry(cat_id) .or_insert_with(|| PlaylistGroup { - id: cat_id, - title: String::from("Unknown"), - channels: Vec::new(), - xtream_cluster: cluster, + id: cat_id, + title: String::from("Unknown"), + channels: Vec::new(), + xtream_cluster: cluster, }) .channels.push(PlaylistItem::from(item)); } } } } - + Ok(groups.into_values().collect()) } diff --git a/backend/src/utils/file/mapping_reader.rs b/backend/src/utils/file/mapping_reader.rs index f48bc6bd4..d3eccbf54 100644 --- a/backend/src/utils/file/mapping_reader.rs +++ b/backend/src/utils/file/mapping_reader.rs @@ -101,7 +101,7 @@ fn read_mappings_from_directory(path: &Path, resolve_env: bool) -> Result