From 873235d8f32115d025f2e69918cdd3d3ae012e0b Mon Sep 17 00:00:00 2001 From: euzu Date: Tue, 13 Jan 2026 16:37:27 +0100 Subject: [PATCH] - Fixed Sort - Added cluster to favourites - Added Regex cache to avoid multiple compilations of same regex --- Cargo.lock | 11 +- README.md | 4 +- backend/Cargo.toml | 4 +- backend/src/api/model/app_state.rs | 4 +- backend/src/model/config/epg_smart_match.rs | 6 +- backend/src/model/config/favourites.rs | 8 +- backend/src/model/config/ip_check.rs | 19 +- backend/src/model/config/rename.rs | 5 +- backend/src/model/config/sort.rs | 5 +- backend/src/model/config/target.rs | 34 +- backend/src/model/config/video_download.rs | 7 +- backend/src/processing/processor/playlist.rs | 25 +- backend/src/processing/processor/sort.rs | 177 +++-- backend/src/processing/processor/trakt.rs | 42 +- backend/src/repository/bplustree.rs | 165 ++++- backend/src/repository/library_repository.rs | 12 +- .../src/repository/m3u_playlist_iterator.rs | 69 +- backend/src/repository/m3u_repository.rs | 43 +- backend/src/repository/mod.rs | 1 + backend/src/repository/playlist_repository.rs | 7 +- backend/src/repository/playlist_source.rs | 12 +- backend/src/repository/sorted_index.rs | 605 ++++++++++++++++++ backend/src/repository/storage.rs | 5 +- backend/src/repository/storage_const.rs | 1 + .../repository/xtream_playlist_iterator.rs | 64 +- backend/src/repository/xtream_repository.rs | 43 +- backend/src/utils/network/ip_checker.rs | 3 +- frontend/scss/app/_component.scss | 5 +- ...scss => _proxy_user_credentials_form.scss} | 0 .../app/components/userlist/_user_table.scss | 24 + .../components/userlist/_userlist_view.scss | 4 +- frontend/src/app/components/search.rs | 4 +- .../src/app/components/userlist/user_table.rs | 21 +- frontend/src/app/context.rs | 4 +- shared/Cargo.toml | 3 +- shared/src/foundation/filter.rs | 7 +- shared/src/foundation/mapper.rs | 5 +- shared/src/model/config/base.rs | 2 +- shared/src/model/config/epg_smart_match.rs | 2 +- shared/src/model/config/favourites.rs | 6 +- shared/src/model/config/ipcheck.rs | 5 +- shared/src/model/config/mod.rs | 1 + shared/src/model/config/rename.rs | 2 +- shared/src/model/config/sort.rs | 7 +- shared/src/model/config/target.rs | 2 +- shared/src/model/config/video_download.rs | 2 +- shared/src/model/mod.rs | 4 +- shared/src/model/playlist.rs | 1 - shared/src/model/playlist_request.rs | 2 +- shared/src/model/regex_cache.rs | 50 ++ shared/src/utils/constants.rs | 6 +- shared/src/utils/serde_utils.rs | 25 + 52 files changed, 1287 insertions(+), 288 deletions(-) create mode 100644 backend/src/repository/sorted_index.rs rename frontend/scss/app/components/userlist/{_proxy-user-credentials-form.scss => _proxy_user_credentials_form.scss} (100%) create mode 100644 frontend/scss/app/components/userlist/_user_table.scss create mode 100644 shared/src/model/regex_cache.rs diff --git a/Cargo.lock b/Cargo.lock index 871ed35b9..1dad16a88 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3284,9 +3284,9 @@ dependencies = [ [[package]] name = "quick-xml" -version = "0.38.4" +version = "0.39.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b66c2058c55a409d601666cffe35f04333cf1013010882cec174a7467cd4e21c" +checksum = "f2e3bf4aa9d243beeb01a7b3bc30b77cfe2c44e24ec02d751a7104a53c2c49a1" dependencies = [ "memchr", "tokio", @@ -3854,9 +3854,9 @@ dependencies = [ [[package]] name = "serde-saphyr" -version = "0.0.13" +version = "0.0.14" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "be7f28bb35aab6f9cdc5f464f9bbda0b1e3908766e819557111cc13e58ce7915" +checksum = "45400dbf0a2c4c2af106c08eb1028b3e45bea3cd45517d3a0e8170b86597122f" dependencies = [ "ahash", "annotate-snippets", @@ -3864,11 +3864,11 @@ dependencies = [ "encoding_rs_io", "nohash-hasher", "num-traits", - "ryu", "saphyr-parser", "serde", "serde_json", "smallvec 2.0.0-alpha.12", + "zmij", ] [[package]] @@ -3991,6 +3991,7 @@ dependencies = [ "bytes", "chrono", "ciborium", + "dashmap", "deunicode", "enum-iterator", "fastrand", diff --git a/README.md b/README.md index 675c343a1..1f9f077bd 100644 --- a/README.md +++ b/README.md @@ -1448,6 +1448,7 @@ sources: ### 2.2.2.8 `favourites` Allows you to explicitly add items to a favorite group based on a filter. This is processed after mapping and resolution. +- `cluster`: can be Series, Movie or Live. - `group`: The name of the group to add the favorite items to. - `filter`: A filter statement to select the original items. - `match_as_ascii`: _optional_ (default `false`). If `true`, the filter matching will be case-insensitive and normalized (e.g., "Cinema" matches "Cinéma"). @@ -1455,7 +1456,8 @@ Allows you to explicitly add items to a favorite group based on a filter. This i Example: ```yaml favourites: - - group: "My Favourites" + - cluster: series + group: "My Favourites" filter: 'Name ~ "Cinema"' match_as_ascii: true ``` diff --git a/backend/Cargo.toml b/backend/Cargo.toml index 4384cb5f6..07e881109 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -11,9 +11,9 @@ lol_html = "2.7" handlebars = "6" shared = { version = "3", path = "../shared" } serde = { version = "1.0", features = ["derive", "rc"] } -serde-saphyr = "0.0.13" +serde-saphyr = "0.0.14" serde_json = { version = "1", features = ["raw_value"] } -quick-xml = { version = "0.38", features = ["async-tokio"] } +quick-xml = { version = "0.39", features = ["async-tokio"] } regex = "1.12" clap = { version = "4", features = ["derive"] } reqwest = { version = "0", features = ["blocking", "json", "stream", "rustls-tls", "socks"] } diff --git a/backend/src/api/model/app_state.rs b/backend/src/api/model/app_state.rs index 78d091cb6..b4f48100a 100644 --- a/backend/src/api/model/app_state.rs +++ b/backend/src/api/model/app_state.rs @@ -10,7 +10,7 @@ use arc_swap::{ArcSwap, ArcSwapOption}; use log::{error, info}; use reqwest::Client; use shared::error::TuliproxError; -use shared::model::UserConnectionPermission; +use shared::model::{UserConnectionPermission}; use shared::utils::small_vecs_equal_unordered; use std::collections::HashMap; use std::sync::atomic::AtomicI8; @@ -293,6 +293,7 @@ impl AppState { self.geoip.store(new_geoip); } + shared::model::REGEX_CACHE.sweep(); Ok(changes) } @@ -326,6 +327,7 @@ impl AppState { self.app_config.set_sources(sources)?; self.active_provider.update_config(&self.app_config).await; + shared::model::REGEX_CACHE.sweep(); Ok(changes) } diff --git a/backend/src/model/config/epg_smart_match.rs b/backend/src/model/config/epg_smart_match.rs index 298dc11de..e01b24823 100644 --- a/backend/src/model/config/epg_smart_match.rs +++ b/backend/src/model/config/epg_smart_match.rs @@ -1,3 +1,4 @@ +use std::sync::Arc; use shared::utils::CONSTANTS; use regex::Regex; use shared::model::{EpgNamePrefix, EpgSmartMatchConfigDto}; @@ -6,14 +7,13 @@ use crate::model::macros; #[derive(Debug, Clone)] pub struct EpgSmartMatchConfig { pub enabled: bool, - pub normalize_regex: Regex, + pub normalize_regex: Arc, pub strip: Vec, pub name_prefix: EpgNamePrefix, pub name_prefix_separator: Vec, pub fuzzy_matching: bool, pub match_threshold: u16, pub best_match_threshold: u16, - } macros::from_impl!(EpgSmartMatchConfig); @@ -22,7 +22,7 @@ impl From<&EpgSmartMatchConfigDto> for EpgSmartMatchConfig { Self { enabled: dto.enabled, normalize_regex: match &dto.normalize_regex { - Some(regex_str) => regex::Regex::new(regex_str).unwrap_or_else(|_| CONSTANTS.re_epg_normalize.clone()), + Some(regex_str) => shared::model::REGEX_CACHE.get_or_compile(regex_str).unwrap_or_else(|_| CONSTANTS.re_epg_normalize.clone()), None => CONSTANTS.re_epg_normalize.clone(), }, strip: match &dto.strip { diff --git a/backend/src/model/config/favourites.rs b/backend/src/model/config/favourites.rs index 09000ab45..cb0fa52ce 100644 --- a/backend/src/model/config/favourites.rs +++ b/backend/src/model/config/favourites.rs @@ -1,11 +1,11 @@ use std::sync::Arc; use crate::model::macros; -use regex::Regex; use shared::foundation::filter::{CompiledRegex, Filter}; -use shared::model::{ConfigFavouritesDto, ItemField}; +use shared::model::{ConfigFavouritesDto, ItemField, XtreamCluster}; #[derive(Debug, Clone)] pub struct ConfigFavourites { + pub cluster: XtreamCluster, pub group: Arc, pub filter: Filter, pub match_as_ascii: bool, @@ -17,7 +17,7 @@ impl ConfigFavourites { ItemField::Group, CompiledRegex { restr: String::new(), - re: Regex::new(".*").unwrap(), + re: shared::model::REGEX_CACHE.get_or_compile(".*").expect("default regex '.*' must compile"), }, ) } @@ -28,6 +28,7 @@ macros::from_impl!(ConfigFavourites); impl From<&ConfigFavouritesDto> for ConfigFavourites { fn from(dto: &ConfigFavouritesDto) -> Self { Self { + cluster: dto.cluster, group: dto.group.clone(), filter: dto.t_filter.as_ref().map_or_else(Self::default_filter, Clone::clone), match_as_ascii: dto.match_as_ascii, @@ -38,6 +39,7 @@ impl From<&ConfigFavouritesDto> for ConfigFavourites { impl From<&ConfigFavourites> for ConfigFavouritesDto { fn from(instance: &ConfigFavourites) -> Self { Self { + cluster: instance.cluster, group: instance.group.clone(), filter: instance.filter.to_string(), match_as_ascii: instance.match_as_ascii, diff --git a/backend/src/model/config/ip_check.rs b/backend/src/model/config/ip_check.rs index ad15794de..7098285c1 100644 --- a/backend/src/model/config/ip_check.rs +++ b/backend/src/model/config/ip_check.rs @@ -1,5 +1,6 @@ +use std::sync::Arc; use regex::Regex; -use shared::model::IpCheckConfigDto; +use shared::model::{IpCheckConfigDto}; use crate::model::macros; #[derive(Debug, Clone)] @@ -7,8 +8,8 @@ pub struct IpCheckConfig { pub url: Option, pub url_ipv4: Option, pub url_ipv6: Option, - pub pattern_ipv4: Option, - pub pattern_ipv6: Option, + pub pattern_ipv4: Option>, + pub pattern_ipv6: Option>, } macros::from_impl!(IpCheckConfig); @@ -18,8 +19,16 @@ impl From<&IpCheckConfigDto> for IpCheckConfig { url: dto.url.clone(), url_ipv4: dto.url_ipv4.clone(), url_ipv6: dto.url_ipv6.clone(), - pattern_ipv4: dto.pattern_ipv4.as_ref().and_then(|s| Regex::new(s).ok()), - pattern_ipv6: dto.pattern_ipv6.as_ref().and_then(|s| Regex::new(s).ok()), + pattern_ipv4: dto.pattern_ipv4.as_ref().and_then(|s| { + shared::model::REGEX_CACHE.get_or_compile(s) + .map_err(|e| log::warn!("Invalid pattern_ipv4 regex '{s}': {e}")) + .ok() + }), + pattern_ipv6: dto.pattern_ipv6.as_ref().and_then(|s| { + shared::model::REGEX_CACHE.get_or_compile(s) + .map_err(|e| log::warn!("Invalid pattern_ipv6 regex '{s}': {e}")) + .ok() + }), } } } diff --git a/backend/src/model/config/rename.rs b/backend/src/model/config/rename.rs index 2068eb08d..1449a9470 100644 --- a/backend/src/model/config/rename.rs +++ b/backend/src/model/config/rename.rs @@ -1,3 +1,4 @@ +use std::sync::Arc; use shared::model::{ConfigRenameDto, ItemField}; use crate::model::macros; @@ -5,7 +6,7 @@ use crate::model::macros; pub struct ConfigRename { pub field: ItemField, pub new_name: String, - pub pattern: regex::Regex, + pub pattern: Arc, } macros::from_impl!(ConfigRename); @@ -14,7 +15,7 @@ impl From<&ConfigRenameDto> for ConfigRename { Self { field: dto.field, new_name: dto.new_name.clone(), - pattern: regex::Regex::new(&dto.pattern).unwrap() + pattern: shared::model::REGEX_CACHE.get_or_compile(&dto.pattern).unwrap_or_else(|_| panic!("Invalid regex pattern {}", dto.pattern)), } } } diff --git a/backend/src/model/config/sort.rs b/backend/src/model/config/sort.rs index d433ea0ad..93d84ee00 100644 --- a/backend/src/model/config/sort.rs +++ b/backend/src/model/config/sort.rs @@ -1,3 +1,4 @@ +use std::sync::Arc; use regex::Regex; use shared::model::{ConfigSortRuleDto, ConfigSortDto, ItemField, SortOrder, SortTarget}; use shared::foundation::filter::Filter; @@ -10,7 +11,7 @@ pub struct ConfigSortRule { pub target: SortTarget, pub order: SortOrder, pub field: ItemField, - pub sequence: Option>, + pub sequence: Option>>, pub filter: Filter, } @@ -33,7 +34,7 @@ impl From<&ConfigSortRule> for ConfigSortRuleDto { target: instance.target, order: instance.order, field: instance.field, - sequence: instance.sequence.as_ref().map(|l: &Vec| l.iter().map(ToString::to_string).collect()), + sequence: instance.sequence.as_ref().map(|l: &Vec>| l.iter().map(ToString::to_string).collect()), filter: instance.filter.to_string(), t_sequence: None, t_filter: Some(instance.filter.clone()), diff --git a/backend/src/model/config/target.rs b/backend/src/model/config/target.rs index 3adb242e3..a522b9206 100644 --- a/backend/src/model/config/target.rs +++ b/backend/src/model/config/target.rs @@ -1,15 +1,14 @@ -use crate::model::mapping::Mapping; +use crate::model::config::favourites::ConfigFavourites; use crate::model::config::trakt::TraktConfig; +use crate::model::mapping::Mapping; +use crate::model::{macros, ConfigRename, ConfigSort}; use arc_swap::ArcSwapOption; use shared::model::{ConfigTargetDto, ConfigTargetOptions, HdHomeRunTargetOutputDto, M3uTargetOutputDto, ProcessingOrder, StrmExportStyle, StrmTargetOutputDto, TargetOutputDto, TargetType, TraktConfigDto, XtreamTargetOutputDto}; use shared::model::PlaylistItemType; use std::sync::Arc; -use regex::Regex; use shared::foundation::filter::Filter; use shared::foundation::filter::ValueProvider; -use crate::model::{macros, ConfigRename, ConfigSort}; -use crate::model::config::favourites::ConfigFavourites; #[derive(Clone, Debug)] pub struct ProcessTargets { @@ -197,9 +196,9 @@ impl From<&TargetOutputDto> for TargetOutput { fn from(dto: &TargetOutputDto) -> Self { match dto { TargetOutputDto::Xtream(o) => TargetOutput::Xtream(XtreamTargetOutput::from(o)), - TargetOutputDto::M3u(o) => TargetOutput::M3u(M3uTargetOutput::from(o)), - TargetOutputDto::Strm(o) => TargetOutput::Strm(StrmTargetOutput::from(o)), - TargetOutputDto::HdHomeRun(o) => TargetOutput::HdHomeRun(HdHomeRunTargetOutput::from(o)), + TargetOutputDto::M3u(o) => TargetOutput::M3u(M3uTargetOutput::from(o)), + TargetOutputDto::Strm(o) => TargetOutput::Strm(StrmTargetOutput::from(o)), + TargetOutputDto::HdHomeRun(o) => TargetOutput::HdHomeRun(HdHomeRunTargetOutput::from(o)), } } } @@ -208,9 +207,9 @@ impl From<&TargetOutput> for TargetOutputDto { fn from(instance: &TargetOutput) -> Self { match instance { TargetOutput::Xtream(o) => TargetOutputDto::Xtream(XtreamTargetOutputDto::from(o)), - TargetOutput::M3u(o) => TargetOutputDto::M3u(M3uTargetOutputDto::from(o)), - TargetOutput::Strm(o) => TargetOutputDto::Strm(StrmTargetOutputDto::from(o)), - TargetOutput::HdHomeRun(o) => TargetOutputDto::HdHomeRun(HdHomeRunTargetOutputDto::from(o)), + TargetOutput::M3u(o) => TargetOutputDto::M3u(M3uTargetOutputDto::from(o)), + TargetOutput::Strm(o) => TargetOutputDto::Strm(StrmTargetOutputDto::from(o)), + TargetOutput::HdHomeRun(o) => TargetOutputDto::HdHomeRun(HdHomeRunTargetOutputDto::from(o)), } } } @@ -229,14 +228,13 @@ pub struct ConfigTarget { pub mapping: Arc>>, pub favourites: Option>, pub processing_order: ProcessingOrder, - pub watch: Option>, + pub watch: Option>>, pub use_memory_cache: bool, } impl ConfigTarget { - pub fn filter(&self, provider: &ValueProvider) -> bool { - self.filter.filter(provider) + self.filter.filter(provider) } pub(crate) fn get_xtream_output(&self) -> Option<&XtreamTargetOutput> { @@ -297,7 +295,6 @@ impl ConfigTarget { macros::from_impl!(ConfigTarget); impl From<&ConfigTargetDto> for ConfigTarget { fn from(dto: &ConfigTargetDto) -> Self { - Self { id: dto.id, enabled: dto.enabled, @@ -311,7 +308,14 @@ impl From<&ConfigTargetDto> for ConfigTarget { mapping: Arc::new(ArcSwapOption::new(None)), favourites: dto.favourites.as_ref().map(|f| f.iter().map(Into::into).collect()), processing_order: dto.processing_order, - watch: dto.watch.as_ref().map(|list| list.iter().filter_map(|s| Regex::new(s).ok()).collect()), + watch: dto.watch.as_ref().map(|list| list.iter().filter_map(|s| + match shared::model::REGEX_CACHE.get_or_compile(s) { + Ok(re) => Some(re), + Err(e) => { + log::warn!("Invalid watch regex pattern '{s}': {e}"); + None + } + }).collect()), use_memory_cache: dto.use_memory_cache, } } diff --git a/backend/src/model/config/video_download.rs b/backend/src/model/config/video_download.rs index 0ca0493c8..bcafcff8d 100644 --- a/backend/src/model/config/video_download.rs +++ b/backend/src/model/config/video_download.rs @@ -1,5 +1,6 @@ use regex::Regex; use std::collections::HashMap; +use std::sync::Arc; use shared::model::{VideoConfigDto, VideoDownloadConfigDto}; use crate::model::macros; @@ -8,7 +9,7 @@ pub struct VideoDownloadConfig { pub headers: HashMap, pub directory: String, pub organize_into_directories: bool, - pub episode_pattern: Option, + pub episode_pattern: Option>, } macros::from_impl!(VideoDownloadConfig); @@ -18,7 +19,9 @@ impl From<&VideoDownloadConfigDto> for VideoDownloadConfig { headers: dto.headers.clone(), directory: dto.directory.as_ref().map_or_else(|| "downloads".to_string(), ToString::to_string), organize_into_directories: dto.organize_into_directories, - episode_pattern: dto.episode_pattern.as_ref().and_then(|s| Regex::new(s).ok()), + episode_pattern: dto.episode_pattern.as_ref().and_then(|s| shared::model::REGEX_CACHE.get_or_compile(s) + .map_err(|e| log::warn!("Invalid episode_pattern regex '{s}': {e}")) + .ok()), } } } diff --git a/backend/src/processing/processor/playlist.rs b/backend/src/processing/processor/playlist.rs index 1bf3abfae..987fb7980 100644 --- a/backend/src/processing/processor/playlist.rs +++ b/backend/src/processing/processor/playlist.rs @@ -40,7 +40,7 @@ use shared::foundation::filter::{get_field_value, set_field_value, Filter, Value use shared::model::xtream_const::XTREAM_CLUSTER; use shared::model::{CounterModifier, FieldGetAccessor, FieldSetAccessor, InputType, ItemField, MsgKind, PlaylistGroup, PlaylistItem, PlaylistItemType, PlaylistUpdateState, - ProcessingOrder, UUIDType}; + ProcessingOrder, UUIDType, XtreamCluster}; use shared::utils::{create_alias_uuid, default_as_default, StringInterner}; use std::time::Instant; @@ -50,13 +50,14 @@ fn is_valid(pli: &PlaylistItem, filter: &Filter, match_as_ascii: bool) -> bool { } pub fn apply_filter_to_source(source: &mut dyn PlaylistSource, filter: &Filter) -> Option> { - let mut groups: IndexMap, PlaylistGroup> = IndexMap::new(); + let mut groups: IndexMap = IndexMap::new(); for pli in source.into_items() { if is_valid(&pli, filter, false) { let group_title = pli.header.group.clone(); let cluster = pli.header.xtream_cluster; let cat_id = pli.header.category_id; - groups.entry(group_title.clone()) + let key = (cluster, group_title.clone()); + groups.entry(key) .or_insert_with(|| PlaylistGroup { id: cat_id, title: group_title, @@ -130,7 +131,7 @@ fn exec_rename(pli: &mut PlaylistItem, rename: Option<&Vec>, inter fn rename_playlist(source: &mut dyn PlaylistSource, target: &ConfigTarget, interner: &mut StringInterner) -> Option> { match &target.rename { Some(renames) if !renames.is_empty() => { - let mut groups: IndexMap, PlaylistGroup> = IndexMap::new(); + let mut groups: IndexMap<(XtreamCluster, Arc), PlaylistGroup> = IndexMap::new(); for mut pli in source.into_items() { // Handle group rename first if it's in the renames for r in renames { @@ -146,7 +147,7 @@ fn rename_playlist(source: &mut dyn PlaylistSource, target: &ConfigTarget, inter let group_title = pli.header.group.clone(); let cluster = pli.header.xtream_cluster; let cat_id = pli.header.category_id; - groups.entry(group_title.clone()) + groups.entry((cluster, group_title.clone())) .or_insert_with(|| PlaylistGroup { id: cat_id, title: group_title, @@ -205,12 +206,12 @@ fn map_playlist(source: &mut dyn PlaylistSource, target: &ConfigTarget, _interne Box::new(iter.flat_map(move |chan| map_channel_and_flatten(chan, mapping))) as Box> }); - let mut next_groups: IndexMap, PlaylistGroup> = IndexMap::new(); + let mut next_groups: IndexMap<(XtreamCluster, Arc), PlaylistGroup> = IndexMap::new(); let mut grp_id: u32 = 0; for channel in mapped_iter { let group_title = channel.header.group.clone(); let cluster = channel.header.xtream_cluster; - next_groups.entry(group_title.clone()) + next_groups.entry((cluster, group_title.clone())) .or_insert_with(|| { grp_id += 1; PlaylistGroup { @@ -774,7 +775,7 @@ async fn process_playlist_for_target(ctx: &PlaylistProcessingContext, pub fn process_favourites(playlist: &mut Vec, favourites_cfg: Option<&[ConfigFavourites]>) { if let Some(favourites) = favourites_cfg { - let mut fav_groups: IndexMap, Vec> = IndexMap::new(); + let mut fav_groups: IndexMap> = IndexMap::new(); for pg in playlist.iter() { for pli in &pg.channels { // series episodes can't be included in favourites @@ -782,13 +783,13 @@ pub fn process_favourites(playlist: &mut Vec, favourites_cfg: Opt continue; } for fav in favourites { - if is_valid(pli, &fav.filter, fav.match_as_ascii) { + if pli.header.xtream_cluster == fav.cluster && is_valid(pli, &fav.filter, fav.match_as_ascii) { let mut channel = pli.clone(); channel.header.group.clone_from(&fav.group); // Update UUID to be an alias of the original channel.header.uuid = create_alias_uuid(&pli.header.uuid, &fav.group); fav_groups - .entry(fav.group.clone()) + .entry((fav.cluster, fav.group.clone())) .or_default() .push(channel); } @@ -796,9 +797,9 @@ pub fn process_favourites(playlist: &mut Vec, favourites_cfg: Opt } } - for (group_name, channels) in fav_groups { + for (fav_group, channels) in fav_groups { if !channels.is_empty() { - let xtream_cluster = channels[0].header.xtream_cluster; + let (xtream_cluster, group_name) = fav_group; playlist.push(PlaylistGroup { id: 0, title: group_name, diff --git a/backend/src/processing/processor/sort.rs b/backend/src/processing/processor/sort.rs index e5c84e517..5abd958ad 100644 --- a/backend/src/processing/processor/sort.rs +++ b/backend/src/processing/processor/sort.rs @@ -2,6 +2,7 @@ use crate::model::{ConfigSortRule, ConfigTarget}; use shared::foundation::filter::ValueProvider; use shared::model::{PlaylistGroup, SortOrder, SortTarget}; use std::cmp::Ordering; +use std::sync::Arc; fn direction(order: SortOrder, ordering: Ordering) -> Ordering { match (order, ordering) { @@ -12,7 +13,7 @@ fn direction(order: SortOrder, ordering: Ordering) -> Ordering { } fn playlist_comparator( - sequence: Option<&Vec>, + sequence: Option<&Vec>>, order: SortOrder, value_a: &str, value_b: &str, @@ -77,24 +78,11 @@ fn playlist_comparator( } } - let o = value_a.cmp(value_b); - match order { - SortOrder::Asc => o, - SortOrder::Desc => o.reverse(), - SortOrder::None => Ordering::Equal, - } + Ordering::Equal } (Some(_), None) => direction(order, Ordering::Less), (None, Some(_)) => direction(order, Ordering::Greater), - (None, None) => { - // NP match → fallback - let o = value_a.cmp(value_b); - match order { - SortOrder::Asc => o, - SortOrder::Desc => o.reverse(), - SortOrder::None => Ordering::Equal, - } - } + (None, None) => Ordering::Equal, } } else { // No Regex-Sequence defined → fallback @@ -123,65 +111,64 @@ pub(in crate::processing::processor) fn sort_playlist( true } -fn sort_groups(groups: &mut [PlaylistGroup], rules: &[ConfigSortRule], match_as_ascii: bool) { +fn sort_groups( + groups: &mut [PlaylistGroup], + rules: &[ConfigSortRule], + match_as_ascii: bool, +) { let group_rules: Vec<_> = rules .iter() .filter(|r| matches!(r.target, SortTarget::Group)) + .filter(|r| r.order != SortOrder::None) .collect(); if group_rules.is_empty() { return; } - for rule in &group_rules { - if rule.order == SortOrder::None { - continue; - } - groups.sort_by(|a_grp, b_grp| { - let (a_chan, b_chan) = match (a_grp.channels.first(), b_grp.channels.first()) { - (None, None) => return Ordering::Equal, + groups.sort_by(|a_grp, b_grp| { + let a_chan = a_grp.channels.first(); + let b_chan = b_grp.channels.first(); + + for rule in &group_rules { + let (vp_a, vp_b) = match (a_chan, b_chan) { + (Some(a), Some(b)) => ( + ValueProvider { pli: a, match_as_ascii }, + ValueProvider { pli: b, match_as_ascii }, + ), (Some(_), None) => return direction(rule.order, Ordering::Less), (None, Some(_)) => return direction(rule.order, Ordering::Greater), - (Some(a), Some(b)) => (a, b), + (None, None) => continue, }; - let provider_a = ValueProvider { - pli: a_chan, - match_as_ascii, - }; - let provider_b = ValueProvider { - pli: b_chan, - match_as_ascii, - }; + let fa = rule.filter.filter(&vp_a); + let fb = rule.filter.filter(&vp_b); - match ( - rule.filter.filter(&provider_a), - rule.filter.filter(&provider_b), - ) { - // (false, false) => return Ordering::Equal, - // (true, false) => return direction(rule.order, Ordering::Less), - // (false, true) => return direction(rule.order, Ordering::Greater), - (true, true) => { /* fallthrough */ } - _ => return Ordering::Equal, + match (fa, fb) { + (false, false) => continue, + (true, false) => return direction(rule.order, Ordering::Less), + (false, true) => return direction(rule.order, Ordering::Greater), + _ => {} } - let va = provider_a.get(rule.field.as_str()); - let vb = provider_b.get(rule.field.as_str()); - match (va, vb) { - (None, None) => return Ordering::Equal, - (Some(_), None) => return direction(rule.order, Ordering::Less), - (None, Some(_)) => return direction(rule.order, Ordering::Greater), + let va = vp_a.get(rule.field.as_str()); + let vb = vp_b.get(rule.field.as_str()); + let ord = match (va, vb) { + (None, None) => Ordering::Equal, + (Some(_), None) => direction(rule.order, Ordering::Less), + (None, Some(_)) => direction(rule.order, Ordering::Greater), (Some(va), Some(vb)) => { - let ord = playlist_comparator(rule.sequence.as_ref(), rule.order, &va, &vb); - - if ord != Ordering::Equal { - return ord; - } + playlist_comparator(rule.sequence.as_ref(), rule.order, &va, &vb) } + }; + + if ord != Ordering::Equal { + return ord; } - Ordering::Equal - }); - } + } + + Ordering::Equal + }); } fn sort_channels_in_groups( @@ -192,50 +179,48 @@ fn sort_channels_in_groups( let channel_rules: Vec<_> = rules .iter() .filter(|r| matches!(r.target, SortTarget::Channel)) + .filter(|r| r.order != SortOrder::None) .collect(); if channel_rules.is_empty() { return; } - for rule in channel_rules { - for group in &mut *groups { - if rule.order == SortOrder::None { - continue; - } - group.channels.sort_by(|a, b| { - let provider_a = ValueProvider { - pli: a, - match_as_ascii, - }; - let provider_b = ValueProvider { - pli: b, - match_as_ascii, - }; + for group in groups { + group.channels.sort_by(|a, b| { + let vp_a = ValueProvider { pli: a, match_as_ascii }; + let vp_b = ValueProvider { pli: b, match_as_ascii }; - match ( - rule.filter.filter(&provider_a), - rule.filter.filter(&provider_b), - ) { - // (false, false) => return Ordering::Equal, - // (true, false) => return direction(rule.order, Ordering::Less), - // (false, true) => return direction(rule.order, Ordering::Greater), - (true, true) => { /* fallthrough */ } - _ => return Ordering::Equal, + for rule in &channel_rules { + let fa = rule.filter.filter(&vp_a); + let fb = rule.filter.filter(&vp_b); + + match (fa, fb) { + (false, false) => continue, + (true, false) => return direction(rule.order, Ordering::Less), + (false, true) => return direction(rule.order, Ordering::Greater), + _ => {} } - let va = provider_a.get(rule.field.as_str()); - let vb = provider_b.get(rule.field.as_str()); - match (va, vb) { + let va = vp_a.get(rule.field.as_str()); + let vb = vp_b.get(rule.field.as_str()); + + let ord = match (va, vb) { (None, None) => Ordering::Equal, (Some(_), None) => direction(rule.order, Ordering::Less), (None, Some(_)) => direction(rule.order, Ordering::Greater), (Some(va), Some(vb)) => { playlist_comparator(rule.sequence.as_ref(), rule.order, &va, &vb) } + }; + + if ord != Ordering::Equal { + return ord; } - }); - } + } + + Ordering::Equal + }); } } @@ -243,7 +228,6 @@ fn sort_channels_in_groups( mod tests { use crate::model::ConfigSortRule; use crate::processing::processor::sort::playlist_comparator; - use regex::Regex; use shared::foundation::filter::Filter; use shared::model::{ItemField, PlaylistItem, PlaylistItemHeader, SortOrder, SortTarget}; use std::cmp::Ordering; @@ -280,9 +264,9 @@ mod tests { field: ItemField::Caption, order: SortOrder::Asc, sequence: Some(vec![ - Regex::new(r"(?P.*?)\bUHD\b").unwrap(), - Regex::new(r"(?P.*?)\bFHD\b").unwrap(), - Regex::new(r"(?P.*?)\bHD\b").unwrap(), + shared::model::REGEX_CACHE.get_or_compile(r"(?P.*?)\bUHD\b").unwrap(), + shared::model::REGEX_CACHE.get_or_compile(r"(?P.*?)\bFHD\b").unwrap(), + shared::model::REGEX_CACHE.get_or_compile(r"(?P.*?)\bHD\b").unwrap(), ]), filter: Filter::default(), }; @@ -368,14 +352,14 @@ mod tests { field: ItemField::Caption, order: SortOrder::Asc, sequence: Some(vec![ - Regex::new(r"^US\| EAST.*?\[\bUHD\b\](?P.*)").unwrap(), - Regex::new(r"^US\| EAST.*?\[\bFHD\b\](?P.*)").unwrap(), - Regex::new(r"^US\| EAST.*?\[\bHD\b\](?P.*)").unwrap(), - Regex::new(r"^US\| EAST.*?\[\bSD\b\](?P.*)").unwrap(), - Regex::new(r"^US\| WEST.*?\[\bUHD\b\](?P.*)").unwrap(), - Regex::new(r"^US\| WEST.*?\[\bFHD\b\](?P.*)").unwrap(), - Regex::new(r"^US\| WEST.*?\[\bHD\b\](?P.*)").unwrap(), - Regex::new(r"^US\| WEST.*?\[\bSD\b\](?P.*)").unwrap(), + shared::model::REGEX_CACHE.get_or_compile(r"^US\| EAST.*?\[\bUHD\b\](?P.*)").unwrap(), + shared::model::REGEX_CACHE.get_or_compile(r"^US\| EAST.*?\[\bFHD\b\](?P.*)").unwrap(), + shared::model::REGEX_CACHE.get_or_compile(r"^US\| EAST.*?\[\bHD\b\](?P.*)").unwrap(), + shared::model::REGEX_CACHE.get_or_compile(r"^US\| EAST.*?\[\bSD\b\](?P.*)").unwrap(), + shared::model::REGEX_CACHE.get_or_compile(r"^US\| WEST.*?\[\bUHD\b\](?P.*)").unwrap(), + shared::model::REGEX_CACHE.get_or_compile(r"^US\| WEST.*?\[\bFHD\b\](?P.*)").unwrap(), + shared::model::REGEX_CACHE.get_or_compile(r"^US\| WEST.*?\[\bHD\b\](?P.*)").unwrap(), + shared::model::REGEX_CACHE.get_or_compile(r"^US\| WEST.*?\[\bSD\b\](?P.*)").unwrap(), ]), filter: Filter::default(), }; @@ -429,4 +413,5 @@ mod tests { assert_eq!(expected, sorted); } + } diff --git a/backend/src/processing/processor/trakt.rs b/backend/src/processing/processor/trakt.rs index f4353338f..20db6845d 100644 --- a/backend/src/processing/processor/trakt.rs +++ b/backend/src/processing/processor/trakt.rs @@ -7,6 +7,7 @@ use shared::error::TuliproxError; use shared::utils::StringInterner; use shared::model::{FieldGetAccessor, FieldSetAccessor, PlaylistGroup, PlaylistItem, TraktContentType, XtreamCluster}; use shared::utils::CONSTANTS; +use indexmap::IndexMap; use std::borrow::Cow; use strsim::normalized_levenshtein; @@ -130,10 +131,10 @@ fn create_category_from_matches<'a>( matches: Vec>, list_config: &'a TraktListConfig, interner: &mut StringInterner, -) -> Option { - if matches.is_empty() { return None; } +) -> Vec { + if matches.is_empty() { return vec![]; } - let mut matched_items = Vec::new(); + let mut matched_items_by_cluster: IndexMap> = IndexMap::new(); let mut sorted_matches = matches; sorted_matches.sort_by(|a, b| { @@ -150,11 +151,7 @@ fn create_category_from_matches<'a>( for match_result in sorted_matches { let mut modified_item = match_result.playlist_item.clone(); - // Use the (possibly numbered) title from the match result (which now contains the original playlist title) with!(mut modified_item.header => header { - // Synchronize name with title so both fields show the same value - // header.title.clone_from(&match_result.trakt_item.title.to_string()); - // header.name.clone_from(&match_result.trakt_item.title.to_string()); let title = header.get_field("caption").unwrap_or_else(|| Cow::Borrowed(&header.title)); if extract_quality(&title).is_none() { if let Some(quality) = extract_quality(&header.group) { @@ -168,28 +165,18 @@ fn create_category_from_matches<'a>( } header.group = interner.intern(group_title); header.gen_uuid(); + matched_items_by_cluster.entry(header.xtream_cluster).or_default().push(modified_item); }); - matched_items.push(modified_item); } - if matched_items.is_empty() { return None; } - - - let cluster = match list_config.content_type { - TraktContentType::Vod => XtreamCluster::Video, - TraktContentType::Series => XtreamCluster::Series, - TraktContentType::Both => { - matched_items.first() - .map_or(XtreamCluster::Video, |item| item.header.xtream_cluster) + matched_items_by_cluster.into_iter().map(|(cluster, channels)| { + PlaylistGroup { + id: 0, + title: interner.intern(group_title), + channels, + xtream_cluster: cluster, } - }; - - Some(PlaylistGroup { - id: 0, - title: interner.intern(group_title), - channels: matched_items, - xtream_cluster: cluster, - }) + }).collect() } fn match_trakt_items_with_playlist<'a>( @@ -197,7 +184,7 @@ fn match_trakt_items_with_playlist<'a>( playlist: &'a [PlaylistGroup], list_config: &'a TraktListConfig, interner: &mut StringInterner, -) -> Option { +) -> Vec { let trakt_match_items: Vec> = trakt_items .iter() .filter(|item| should_include_item(item, list_config.content_type)) @@ -256,7 +243,8 @@ impl TraktCategoriesProcessor { Ok(trakt_items) => { debug!("Processing Trakt list {cache_key} with {} items", trakt_items.len()); - if let Some(category) = match_trakt_items_with_playlist(&trakt_items, playlist, list_config, &mut interner) { + let categories = match_trakt_items_with_playlist(&trakt_items, playlist, list_config, &mut interner); + for category in categories { if !category.channels.is_empty() { total_matches += category.channels.len(); let category_len = category.channels.len(); diff --git a/backend/src/repository/bplustree.rs b/backend/src/repository/bplustree.rs index 179fec81e..cf01b39c9 100644 --- a/backend/src/repository/bplustree.rs +++ b/backend/src/repository/bplustree.rs @@ -12,10 +12,11 @@ use std::marker::PhantomData; use std::mem::size_of; use std::path::{Path, PathBuf}; use tempfile::NamedTempFile; +use crate::repository::storage::get_file_path_for_db_index; // Constants (Restored) const PAGE_SIZE: u16 = 4096; -const PAGE_SIZE_USIZE: usize = PAGE_SIZE as usize; +pub const PAGE_SIZE_USIZE: usize = PAGE_SIZE as usize; const LEN_SIZE: usize = 4; const FLAG_SIZE: usize = 1; const MAGIC: &[u8; 4] = b"BTRE"; @@ -41,7 +42,7 @@ const PACK_VALUE_HEADER_SIZE: usize = 4; const COMPRESSION_MIN_SIZE: usize = 64; const COMPRESSION_THRESHOLD_PERCENT: usize = 85; const COMPRESSION_FLAG_NONE: u8 = 0x00; -const COMPRESSION_FLAG_LZ4: u8 = 0x01; +pub const COMPRESSION_FLAG_LZ4: u8 = 0x01; // Page Configuration const PAGE_HEADER_SIZE: u16 = 16; @@ -1444,6 +1445,64 @@ where } } + /// Store the tree and build a sorted index file. + /// + /// # Arguments + /// * `filepath` - Path to store the `BPlusTree` + /// * `sort_key_extractor` - Closure that extracts the sort key from a value + /// + /// # Example + /// ```ignore + /// tree.store_with_index(&db_path, |v| v.name.clone())?; + /// ``` + pub fn store_with_index( + &mut self, + filepath: &Path, + sort_key_extractor: F, + ) -> io::Result + where + SortKey: Ord + Serialize, + F: Fn(&V) -> SortKey, + { + // Store the tree first + let result = self.store(filepath)?; + if result > 0 { + Self::store_index(filepath, sort_key_extractor)?; + } + + Ok(result) + } + + pub fn store_index(filepath: &Path, sort_key_extractor: F) -> io::Result<()> + where + SortKey: Ord + Serialize, + F: Fn(&V) -> SortKey + { + let index_path = get_file_path_for_db_index(filepath); + + // Re-open the stored tree to get value locations + let mut query = BPlusTreeQuery::::try_new(filepath)?; + let entries_with_locations = query.collect_with_locations()?; + + // Collect (sort_key, primary_key, location) and sort + let mut sorted_entries: Vec<(SortKey, K, super::sorted_index::ValueLocation)> = + entries_with_locations + .into_iter() + .map(|(k, v, loc)| (sort_key_extractor(&v), k, loc)) + .collect(); + + // Sort by sort key + sorted_entries.sort_by(|a, b| a.0.cmp(&b.0)); + + // Write index file + let mut writer = super::sorted_index::SortedIndexWriter::new(&index_path)?; + for (sort_key, primary_key, location) in &sorted_entries { + writer.push(sort_key, primary_key, *location)?; + } + writer.finish()?; + Ok(()) + } + /// Internal store without locking, used for compaction or initial save. fn store_internal(&mut self, filepath: &Path) -> io::Result { let tempfile = NamedTempFile::new()?; @@ -1798,6 +1857,7 @@ where /// at the cost of higher memory usage. pub struct BPlusTreeQuery { file: BufReader, + filepath: PathBuf, buffer: Vec, cache: IndexMap>, root_offset: u64, @@ -1827,6 +1887,7 @@ where Ok(Self { file: utils::file_reader(file), + filepath: PathBuf::new(), // Unknown when created from file directly buffer: vec![0u8; PAGE_SIZE_USIZE], cache: IndexMap::with_capacity(CACHE_CAPACITY), root_offset, @@ -1837,7 +1898,14 @@ where pub fn try_new(filepath: &Path) -> io::Result { let file = File::open(filepath)?; - Self::try_from_file(file) + let mut query = Self::try_from_file(file)?; + query.filepath = filepath.to_path_buf(); + Ok(query) + } + + /// Returns the filepath this query was opened from. + pub fn filepath(&self) -> &Path { + &self.filepath } pub fn query(&mut self, key: &K) -> Result, BPlusTreeError> { @@ -1867,6 +1935,29 @@ where pub fn iter(&mut self) -> BPlusTreeDiskIterator<'_, K, V> { BPlusTreeDiskIterator::new(self) } + + /// Owned iterator that traverses the tree in order defined by a secondary sorted index. + /// + /// The index path is automatically derived from the tree filepath by changing + /// the extension to `.idx`. For example, if the tree is at `/data/items.bin`, + /// the index is expected at `/data/items.idx`. + /// + /// This iterator reads values directly from stored offsets in O(1) time, + /// avoiding O(log n) tree traversal per item. + pub fn disk_iter_sorted(self) -> io::Result> + where + SortKey: for<'de> Deserialize<'de>, + { + super::sorted_index::BPlusTreeSortedIteratorOwned::new(self.filepath.clone(), self.file) + } + + /// Owned iterator with explicit index path. + pub fn disk_iter_sorted_with_path(self, index_path: &Path) -> io::Result> + where + SortKey: for<'de> Deserialize<'de>, + { + super::sorted_index::BPlusTreeSortedIteratorOwned::with_index_path(self.filepath.clone(), self.file, index_path) + } /// Traverses the tree and calls the provided closure for each leaf's keys and values. pub fn traverse(&mut self, mut f: F) -> io::Result<()> @@ -1882,12 +1973,59 @@ where Ok(()) } + /// Collects all entries with their value locations. + /// Used for building sorted indexes that need direct value access. + pub fn collect_with_locations(&mut self) -> io::Result> { + let mut result = Vec::new(); + let mut stack = vec![self.root_offset]; + + while let Some(offset) = stack.pop() { + let (node, pointers) = BPlusTreeNode::::deserialize_from_block( + &mut self.file, + &mut self.buffer, + offset, + false, + )?; + + if node.is_leaf { + for (key, info) in node.keys.into_iter().zip(node.value_info.iter()) { + let value = BPlusTreeNode::::load_value_from_info(&mut self.file, info)?; + let location = match info.mode { + ValueStorageMode::Single(offset) => { + super::sorted_index::ValueLocation::Single { + offset, + length: info.length, + } + } + ValueStorageMode::Packed(block_offset, index) => { + super::sorted_index::ValueLocation::Packed { + block_offset, + index, + length: info.length, + } + } + }; + result.push((key, value, location)); + } + } else if let Some(ptrs) = pointers { + // Process children in reverse order so we pop them in correct order + for ptr in ptrs.into_iter().rev() { + stack.push(ptr); + } + } + } + + Ok(result) + } + /// Provides an owned disk-backed iterator. pub fn disk_iter(self) -> BPlusTreeDiskIteratorOwned { BPlusTreeDiskIteratorOwned::new(self) } + } + pub struct BPlusTreeDiskIteratorOwned { query: BPlusTreeQuery, stack: Vec<(u64, usize)>, @@ -1973,6 +2111,27 @@ where } } +/// Generic reader that can be either a sorted index iterator or a regular disk iterator. +/// Used for fallback logic (Sorted -> Unsorted). +pub enum PlaylistIteratorReader { + Sorted(super::sorted_index::BPlusTreeSortedIteratorOwned), + Unsorted(BPlusTreeDiskIteratorOwned), +} + +impl Iterator for PlaylistIteratorReader +where + V: Serialize + for<'de> Deserialize<'de> + Clone, +{ + type Item = io::Result<(u32, V)>; + + fn next(&mut self) -> Option { + match self { + PlaylistIteratorReader::Sorted(iter) => iter.next(), + PlaylistIteratorReader::Unsorted(iter) => iter.next().map(Ok), + } + } +} + pub struct BPlusTreeDiskIterator<'a, K, V> { query: &'a mut BPlusTreeQuery, stack: Vec<(u64, usize)>, // (node_offset, next_child_index) diff --git a/backend/src/repository/library_repository.rs b/backend/src/repository/library_repository.rs index 71f233213..47b95453d 100644 --- a/backend/src/repository/library_repository.rs +++ b/backend/src/repository/library_repository.rs @@ -1,10 +1,11 @@ use crate::model::AppConfig; use crate::repository::bplustree::{BPlusTree, BPlusTreeQuery}; use shared::error::{TuliproxError, notify_err_res}; -use shared::model::{PlaylistGroup, PlaylistItem, UUIDType, XtreamPlaylistItem}; +use shared::model::{PlaylistGroup, PlaylistItem, UUIDType, XtreamCluster, XtreamPlaylistItem}; use std::path::Path; use std::sync::Arc; use indexmap::IndexMap; +use crate::repository::xtream_repository::CategoryKey; pub async fn persist_input_library_playlist(app_config: &Arc, library_path: &Path, playlist: &[PlaylistGroup]) -> Result<(), TuliproxError> { if playlist.is_empty() { @@ -26,7 +27,7 @@ pub async fn persist_input_library_playlist(app_config: &Arc, library pub async fn load_input_local_library_playlist(app_config: &Arc, lib_path: &Path) -> Result, TuliproxError> { - let mut groups: IndexMap, PlaylistGroup> = IndexMap::new(); + let mut groups: IndexMap = IndexMap::new(); if tokio::fs::try_exists(lib_path).await.unwrap_or(false) { // Load Items @@ -34,15 +35,16 @@ pub async fn load_input_local_library_playlist(app_config: &Arc, lib_ if let Ok(mut query) = BPlusTreeQuery::::try_new(lib_path) { let mut group_cnt = 0; for (_, ref item) in query.iter() { - let cat_id = item.group.clone(); - groups.entry(cat_id) + let cluster = XtreamCluster::try_from(item.item_type).unwrap_or(XtreamCluster::Live); + let key = (cluster, item.group.clone()); + groups.entry(key) .or_insert_with(|| { group_cnt += 1; PlaylistGroup { id: group_cnt, title: item.group.clone(), channels: Vec::new(), - xtream_cluster: item.xtream_cluster, + xtream_cluster: cluster, } }) .channels.push(PlaylistItem::from(item)); diff --git a/backend/src/repository/m3u_playlist_iterator.rs b/backend/src/repository/m3u_playlist_iterator.rs index 8fbadf756..f572971aa 100644 --- a/backend/src/repository/m3u_playlist_iterator.rs +++ b/backend/src/repository/m3u_playlist_iterator.rs @@ -3,18 +3,20 @@ use shared::error::{TuliproxError}; use crate::model::{AppConfig, ProxyUserCredentials}; use crate::model::{ConfigTarget}; use shared::model::{ConfigTargetOptions, M3uPlaylistItem, PlaylistItemType, ProxyType, TargetType, XtreamCluster}; -use crate::repository::bplustree::{BPlusTreeDiskIteratorOwned, BPlusTreeQuery}; -use crate::repository::m3u_repository::m3u_get_file_path; -use crate::repository::storage::ensure_target_storage_path; +use crate::repository::bplustree::{BPlusTreeQuery, PlaylistIteratorReader}; +use crate::repository::m3u_repository::m3u_get_file_path_for_db; +use crate::repository::storage::{ensure_target_storage_path, get_file_path_for_db_index}; use crate::repository::storage_const; use crate::repository::user_repository::user_get_bouquet_filter; use crate::utils::FileReadGuard; use std::collections::HashSet; +use std::iter::Peekable; +use log::error; // concat_string! macro from shared utils is used for efficient String building #[allow(clippy::struct_excessive_bools)] pub struct M3uPlaylistIterator { - reader: BPlusTreeDiskIteratorOwned, + reader: Peekable>, base_url: String, username: String, password: String, @@ -40,13 +42,28 @@ 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 m3u_path = m3u_get_file_path(&target_path); + let m3u_path = m3u_get_file_path_for_db(&target_path); let file_lock = cfg.file_locks.read_lock(&m3u_path).await; - let query = BPlusTreeQuery::::try_new(&m3u_path) - .map_err(|err| info_err!("Could not open BPlusTreeQuery {m3u_path:?} - {err}"))?; - let reader = query.disk_iter(); + let index_path = get_file_path_for_db_index(&m3u_path); + let reader = if index_path.exists() { + let query = BPlusTreeQuery::::try_new(&m3u_path) + .map_err(|err| info_err!("Could not open BPlusTreeQuery {m3u_path:?} - {err}"))?; + match query.disk_iter_sorted() { + Ok(reader) => PlaylistIteratorReader::Sorted(reader), + Err(err) => { + error!("Sorted index error for m3u, fallback: {err}"); + let query = BPlusTreeQuery::::try_new(&m3u_path) + .map_err(|err| info_err!("Could not open BPlusTreeQuery {m3u_path:?} - {err}"))?; + PlaylistIteratorReader::Unsorted(query.disk_iter()) + } + } + } else { + let query = BPlusTreeQuery::::try_new(&m3u_path) + .map_err(|err| info_err!("Could not open BPlusTreeQuery {m3u_path:?} - {err}"))?; + PlaylistIteratorReader::Unsorted(query.disk_iter()) + }.peekable(); let filter = user_get_bouquet_filter(&config, &user.username, None, TargetType::M3u, XtreamCluster::Live).await; @@ -112,18 +129,35 @@ impl M3uPlaylistIterator { self.get_rewritten_url(m3u_pli, false, storage_const::M3U_RESOURCE_PATH) } + fn find_next_matching(reader: &mut Peekable>, set: &HashSet) -> Option { + loop { + match reader.next() { + Some(Ok((_, item))) => { + if set.contains(&*item.group) { + return Some(item); + } + } + Some(Err(e)) => { + error!("Iterator error: {e}"); + return None; + } + None => return None, + } + } + } + fn get_next(&mut self) -> Option<(M3uPlaylistItem, bool)> { let entry = if let Some(set) = &self.filter { if let Some((current_item, _)) = self.lookup_item.take() { // Avoid cloning strings while filtering - let next_valid = self.reader.find(|(_, pli)| set.contains(&*pli.group)); - self.lookup_item = next_valid.map(|(_k, v)| (v, true)); // has_next handled by iterator usually, but here we just need the item + let next_valid = Self::find_next_matching(&mut self.reader, set); + self.lookup_item = next_valid.map(|v| (v, true)); // has_next handled by iterator usually, but here we just need the item let has_next = self.lookup_item.is_some(); Some((current_item, has_next)) } else { - let current_item = self.reader.find(|(_, item)| set.contains(&*item.group)); - if let Some((_, item)) = current_item { - self.lookup_item = self.reader.find(|(_, item)| set.contains(&*item.group)).map(|(_, item)| (item, true)); + let current_item = Self::find_next_matching(&mut self.reader, set); + if let Some(item) = current_item { + self.lookup_item = Self::find_next_matching(&mut self.reader, set).map(|item| (item, true)); let has_next = self.lookup_item.is_some(); Some((item, has_next)) } else { @@ -131,7 +165,14 @@ impl M3uPlaylistIterator { } } } else { - self.reader.next().map(|(_, v)| (v, !self.reader.is_empty())) + match self.reader.next() { + Some(Ok((_, v))) => Some((v, self.reader.peek().is_some())), + Some(Err(e)) => { + error!("Iterator error: {e}"); + None + } + None => None, + } }; // TODO hls and unknown reverse proxy diff --git a/backend/src/repository/m3u_repository.rs b/backend/src/repository/m3u_repository.rs index d3fc925d9..6a75b0f20 100644 --- a/backend/src/repository/m3u_repository.rs +++ b/backend/src/repository/m3u_repository.rs @@ -3,7 +3,7 @@ use crate::model::{AppConfig, ProxyUserCredentials}; use crate::model::{Config, ConfigTarget, M3uTargetOutput}; use crate::repository::bplustree::{BPlusTree, BPlusTreeQuery}; use crate::repository::m3u_playlist_iterator::M3uPlaylistM3uTextIterator; -use crate::repository::storage::get_target_storage_path; +use crate::repository::storage::{get_file_path_for_db_index, get_target_storage_path}; use crate::repository::storage_const; use crate::utils; use crate::utils::{async_file_writer, IO_BUFFER_SIZE}; @@ -12,12 +12,14 @@ use log::error; use shared::error::{notify_err, info_err, string_to_io_error, str_to_io_error, TuliproxError}; use shared::model::{M3uPlaylistItem, PlaylistGroup}; use shared::model::{PlaylistItem, PlaylistItemType, XtreamCluster}; +use crate::repository::xtream_repository::CategoryKey; use std::io::Error; use std::path::{Path, PathBuf}; use std::sync::Arc; use tokio::fs; use tokio::io::AsyncWriteExt; use tokio::task; +use shared::concat_string; macro_rules! cant_write_result { ($path:expr, $err:expr) => { @@ -25,8 +27,8 @@ macro_rules! cant_write_result { } } -pub fn m3u_get_file_path(target_path: &Path) -> PathBuf { - target_path.join(PathBuf::from(format!("{}.{}", storage_const::FILE_M3U, storage_const::FILE_SUFFIX_DB))) +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(target_path: &Path) -> PathBuf { @@ -86,7 +88,7 @@ pub async fn m3u_write_playlist( return Ok(()); } - let m3u_path = m3u_get_file_path(target_path); + let m3u_path = m3u_get_file_path_for_db(target_path); let m3u_playlist = Arc::new( new_playlist .iter() @@ -111,7 +113,7 @@ pub async fn m3u_write_playlist( for m3u in playlist.iter() { tree.insert(m3u.virtual_id, m3u.clone()); } - tree.store(&m3u_path_clone).map_err(|err| cant_write_result!(&m3u_path_clone, err))?; + tree.store_with_index(&m3u_path_clone, |pli| pli.source_ordinal).map_err(|err| cant_write_result!(&m3u_path_clone, err))?; Ok(()) }) .await @@ -144,7 +146,7 @@ pub async fn m3u_get_item_for_stream_id(stream_id: u32, app_state: &AppState, ta let cfg: &AppConfig = &app_state.app_config; 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 m3u_path = m3u_get_file_path_for_db(&target_path); let _file_lock = cfg.file_locks.read_lock(&m3u_path).await; let mut query = BPlusTreeQuery::::try_new(&m3u_path)?; @@ -158,7 +160,7 @@ pub async fn m3u_get_item_for_stream_id(stream_id: u32, app_state: &AppState, ta pub async fn iter_raw_m3u_playlist(config: &AppConfig, target: &ConfigTarget) -> Option<(utils::FileReadGuard, impl Iterator)> { let target_path = get_target_storage_path(&config.config.load(), target.name.as_str())?; - let m3u_path = m3u_get_file_path(&target_path); + let m3u_path = m3u_get_file_path_for_db(&target_path); if let Ok(false) = tokio::fs::try_exists(&m3u_path).await { return None; } @@ -166,7 +168,23 @@ pub async fn iter_raw_m3u_playlist(config: &AppConfig, target: &ConfigTarget) -> match BPlusTreeQuery::::try_new(&m3u_path) .map_err(|err| info_err!("Could not open BPlusTreeQuery {m3u_path:?} - {err}")) { Ok(mut query) => { - let items: Vec = query.iter().map(|(_, v)| v).collect(); + let index_path = get_file_path_for_db_index(&m3u_path); + let items: Vec = if index_path.exists() { + match query.disk_iter_sorted::() { + Ok(iter) => iter.filter_map(Result::ok).map(|(_, v)| v).collect(), + Err(err) => { + error!("Sorted index error {}: {err}", m3u_path.display()); + // Re-open query for fallback + match BPlusTreeQuery::::try_new(&m3u_path) { + Ok(mut query) => query.iter().map(|(_, v)| v).collect(), + Err(_) => Vec::new(), + } + } + } + } else { + query.iter().map(|(_, v)| v).collect() + }; + let len = items.len(); Some((file_lock, items.into_iter().enumerate().map(move |(i, v)| (v, i < len - 1)))) } @@ -201,7 +219,7 @@ pub async fn persist_input_m3u_playlist(app_config: &Arc, m3u_path: & } pub async fn load_input_m3u_playlist(app_config: &Arc, m3u_path: &Path) -> Result, TuliproxError> { - let mut groups: IndexMap, PlaylistGroup> = IndexMap::new(); + let mut groups: IndexMap = IndexMap::new(); if tokio::fs::try_exists(m3u_path).await.unwrap_or(false) { // Load Items @@ -209,15 +227,16 @@ pub async fn load_input_m3u_playlist(app_config: &Arc, m3u_path: &Pat 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) + let cluster = XtreamCluster::try_from(item.item_type).unwrap_or(XtreamCluster::Live); + let key = (cluster, item.group.clone()); + groups.entry(key) .or_insert_with(|| { group_cnt += 1; PlaylistGroup { id: group_cnt, title: item.group.clone(), channels: Vec::new(), - xtream_cluster: XtreamCluster::try_from(item.item_type).unwrap_or(XtreamCluster::Live), + xtream_cluster: cluster, } }) .channels.push(PlaylistItem::from(item)); diff --git a/backend/src/repository/mod.rs b/backend/src/repository/mod.rs index b1f80a820..614a171fb 100644 --- a/backend/src/repository/mod.rs +++ b/backend/src/repository/mod.rs @@ -13,6 +13,7 @@ pub mod storage_const; mod playlist_scratch; mod playlist_source; mod library_repository; +pub mod sorted_index; pub use playlist_source::*; diff --git a/backend/src/repository/playlist_repository.rs b/backend/src/repository/playlist_repository.rs index 8a7ae5f33..43ca14308 100644 --- a/backend/src/repository/playlist_repository.rs +++ b/backend/src/repository/playlist_repository.rs @@ -4,7 +4,7 @@ use crate::model::{AppConfig, ConfigInput, ConfigTarget, TargetOutput}; use crate::processing::processor::playlist::{apply_filter_to_playlist, PlaylistProcessingContext}; use crate::repository::bplustree::{BPlusTree, BPlusTreeQuery}; use crate::repository::epg_repository::epg_write; -use crate::repository::m3u_repository::{load_input_m3u_playlist, m3u_get_file_path, m3u_write_playlist, persist_input_m3u_playlist}; +use crate::repository::m3u_repository::{load_input_m3u_playlist, m3u_get_file_path_for_db, 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; @@ -47,10 +47,13 @@ pub async fn persist_playlist(app_config: &Arc, playlist: &mut [Playl let mut local_library_series = HashMap::>::new(); let mut provider_series = HashMap::>::new(); + let mut source_ordinal: u32 = 0; // Virtual IDs assignment for group in playlist.iter_mut() { for channel in &mut group.channels { let header = &mut channel.header; + source_ordinal += 1; + header.source_ordinal = source_ordinal; let provider_id = header.get_provider_id().unwrap_or_default(); if provider_id == 0 { header.item_type = match (is_hls_url(&header.url), header.item_type) { @@ -286,7 +289,7 @@ async fn load_m3u_target_storage(app_config: &AppConfig, target: &ConfigTarget) let target_path = get_target_storage_path(&config, target.name.as_str()).ok_or_else(|| info_err!("Could not find path for target {}", &target.name))?; - let m3u_path = m3u_get_file_path(&target_path); + let m3u_path = m3u_get_file_path_for_db(&target_path); let _file_lock = app_config.file_locks.read_lock(&m3u_path).await; let mut tree = BPlusTree::::new(); if let Ok(mut query) = BPlusTreeQuery::::try_new(&m3u_path) { diff --git a/backend/src/repository/playlist_source.rs b/backend/src/repository/playlist_source.rs index fb37e3183..37135a057 100644 --- a/backend/src/repository/playlist_source.rs +++ b/backend/src/repository/playlist_source.rs @@ -188,7 +188,7 @@ impl PlaylistSource for XtreamDiskPlaylistSource { fn take_groups(&mut self) -> Vec { // Build groups on-the-fly using disk iterator (streams one leaf at a time) - let mut groups_map: IndexMap = IndexMap::new(); + let mut groups_map: IndexMap<(XtreamCluster, u32), PlaylistGroup> = IndexMap::new(); let mut iters: Vec<(XtreamCluster, Box + Send>)> = vec![]; if let Some((q, _)) = self.live.as_mut() { iters.push((XtreamCluster::Live, Box::new(q.iter().map(|(_, item)| item)))); @@ -202,7 +202,7 @@ impl PlaylistSource for XtreamDiskPlaylistSource { for (cluster, iter) in iters { for item in iter { - groups_map.entry(item.category_id) + groups_map.entry((cluster, item.category_id)) .or_insert_with(|| PlaylistGroup { id: item.category_id, title: item.group.clone(), @@ -351,14 +351,16 @@ macro_rules! impl_single_file_disk_source { fn take_groups(&mut self) -> Vec { // Build groups on-the-fly using disk iterator (streams one leaf at a time) if let Some(q) = self.playlist.as_mut() { - let mut groups_map: IndexMap, PlaylistGroup> = IndexMap::new(); + let mut groups_map: IndexMap<(XtreamCluster, Arc), PlaylistGroup> = IndexMap::new(); for (_, item) in q.iter() { - groups_map.entry(item.group.clone()) + let cluster = XtreamCluster::try_from(item.item_type).unwrap_or(XtreamCluster::Live); + let key = (cluster, item.group.clone()); + groups_map.entry(key) .or_insert_with(|| PlaylistGroup { id: 0, title: item.group.clone(), channels: vec![], - xtream_cluster: XtreamCluster::try_from(item.item_type).unwrap_or(XtreamCluster::Live), + xtream_cluster: cluster, }) .channels.push(PlaylistItem::from(&item)); } diff --git a/backend/src/repository/sorted_index.rs b/backend/src/repository/sorted_index.rs new file mode 100644 index 000000000..caf796c19 --- /dev/null +++ b/backend/src/repository/sorted_index.rs @@ -0,0 +1,605 @@ +//! Sorted Index for `BPlusTree` +//! +//! This module provides a compact sorted index file that enables iteration +//! over a `BPlusTree` in order of a secondary sort key rather than the primary key. +//! +//! # File Format (v3) +//! ```text +//! [magic: 4 bytes]["SIDX"] +//! [version: u32] +//! [count: u64] +//! [entry0][entry1]... +//! +//! Entry format (with value location for O(1) access): +//! [sort_key_len: u32][sort_key_bytes][primary_key_len: u32][primary_key_bytes][value_location] +//! +//! Value location format: +//! - Single mode: [mode: u8 = 0][offset: u64][length: u32] +//! - Packed mode: [mode: u8 = 1][block_offset: u64][index: u16][length: u32] +//! ``` +//! +//! # Optimization +//! By storing the value location directly in the index, iteration can +//! read values in O(1) time by seeking directly to the offset, avoiding O(log n) +//! tree traversal for each lookup. +//! +//! > **Important**: The index is tightly coupled to the B+Tree file structure. +//! > It must be rebuilt after any operation that changes value offsets (e.g., `compact()`). + +use crate::utils::{binary_deserialize, binary_serialize}; +use serde::{Deserialize, Serialize}; +use std::fs::{File, OpenOptions}; +use std::io::{self, BufReader, BufWriter, Read, Seek, SeekFrom, Write}; +use std::marker::PhantomData; +use std::path::{Path, PathBuf}; +use crate::repository::bplustree::{COMPRESSION_FLAG_LZ4, PAGE_SIZE_USIZE}; +use indexmap::IndexMap; +use crate::repository::storage::get_file_path_for_db_index; + +const MAGIC: &[u8; 4] = b"SIDX"; +const VERSION: u32 = 3; // Bumped for new format with flexible value location +const HEADER_SIZE: usize = 16; // 4 (magic) + 4 (version) + 8 (count) + +const MODE_SINGLE: u8 = 0; +const MODE_PACKED: u8 = 1; + + + +/// Represents how a value can be located and read from the tree file. +#[derive(Debug, Clone, Copy)] +pub enum ValueLocation { + /// Single value stored at a specific offset with a given length. + Single { offset: u64, length: u32 }, + /// Value packed in a block with other values, identified by block offset and index. + Packed { block_offset: u64, index: u16, length: u32 }, +} + +impl ValueLocation { + /// Serialize the value location to bytes. + fn to_bytes(self) -> Vec { + match self { + ValueLocation::Single { offset, length } => { + let mut bytes = vec![MODE_SINGLE]; + bytes.extend_from_slice(&offset.to_le_bytes()); + bytes.extend_from_slice(&length.to_le_bytes()); + bytes + } + ValueLocation::Packed { block_offset, index, length } => { + let mut bytes = vec![MODE_PACKED]; + bytes.extend_from_slice(&block_offset.to_le_bytes()); + bytes.extend_from_slice(&index.to_le_bytes()); + bytes.extend_from_slice(&length.to_le_bytes()); + bytes + } + } + } + + /// Deserialize a value location from a reader. + fn from_reader(reader: &mut R) -> io::Result { + let mut mode = [0u8; 1]; + reader.read_exact(&mut mode)?; + + match mode[0] { + MODE_SINGLE => { + let mut offset_buf = [0u8; 8]; + let mut length_buf = [0u8; 4]; + reader.read_exact(&mut offset_buf)?; + reader.read_exact(&mut length_buf)?; + Ok(ValueLocation::Single { + offset: u64::from_le_bytes(offset_buf), + length: u32::from_le_bytes(length_buf), + }) + } + MODE_PACKED => { + let mut block_offset_buf = [0u8; 8]; + let mut index_buf = [0u8; 2]; + let mut length_buf = [0u8; 4]; + reader.read_exact(&mut block_offset_buf)?; + reader.read_exact(&mut index_buf)?; + reader.read_exact(&mut length_buf)?; + Ok(ValueLocation::Packed { + block_offset: u64::from_le_bytes(block_offset_buf), + index: u16::from_le_bytes(index_buf), + length: u32::from_le_bytes(length_buf), + }) + } + _ => Err(io::Error::new( + io::ErrorKind::InvalidData, + format!("Unknown value location mode: {}", mode[0]), + )), + } + } +} + +/// Entry containing value location information for direct access. +#[derive(Debug, Clone)] +pub struct IndexEntry { + pub sort_key: SortKey, + pub primary_key: K, + pub location: ValueLocation, +} + +/// Writer for building a sorted index file. +/// +/// Entries must be pushed in sorted order. The writer buffers writes +/// and flushes on `finish()`. +pub struct SortedIndexWriter { + writer: BufWriter, + count: u64, + _marker: PhantomData<(SortKey, K)>, +} + +impl SortedIndexWriter +where + SortKey: Serialize, + K: Serialize, +{ + /// Create a new index writer at the given path. + /// Overwrites any existing file. + pub fn new(path: &Path) -> io::Result { + let file = OpenOptions::new() + .write(true) + .create(true) + .truncate(true) + .open(path)?; + + let mut writer = BufWriter::new(file); + + // Write header (count will be updated on finish) + writer.write_all(MAGIC)?; + writer.write_all(&VERSION.to_le_bytes())?; + writer.write_all(&0u64.to_le_bytes())?; // placeholder count + + Ok(Self { + writer, + count: 0, + _marker: PhantomData, + }) + } + + /// Append an entry to the index with value location for O(1) access. + /// Caller must ensure entries are pushed in sorted order. + pub fn push( + &mut self, + sort_key: &SortKey, + primary_key: &K, + location: ValueLocation, + ) -> io::Result<()> { + let sk_bytes = binary_serialize(sort_key)?; + let pk_bytes = binary_serialize(primary_key)?; + + // Write sort key + let sk_len = u32::try_from(sk_bytes.len()) + .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?; + self.writer.write_all(&sk_len.to_le_bytes())?; + self.writer.write_all(&sk_bytes)?; + + // Write primary key + let pk_len = u32::try_from(pk_bytes.len()) + .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?; + self.writer.write_all(&pk_len.to_le_bytes())?; + self.writer.write_all(&pk_bytes)?; + + // Write value location + self.writer.write_all(&location.to_bytes())?; + + self.count += 1; + Ok(()) + } + + /// Finalize the index file by writing the count to the header. + pub fn finish(mut self) -> io::Result { + self.writer.flush()?; + + // Seek back to count position and write final count + let file = self.writer.into_inner()?; + let mut file = file; + file.seek(SeekFrom::Start(8))?; // After magic + version + file.write_all(&self.count.to_le_bytes())?; + file.sync_all()?; + + Ok(self.count) + } +} + +/// Reader for iterating over a sorted index file. +pub struct SortedIndexReader { + reader: BufReader, + remaining: u64, + _marker: PhantomData<(SortKey, K)>, +} + +impl SortedIndexReader +where + SortKey: for<'de> Deserialize<'de>, + K: for<'de> Deserialize<'de>, +{ + /// Open an existing index file for reading. + pub fn open(path: &Path) -> io::Result { + let file = File::open(path)?; + let mut reader = BufReader::new(file); + + // Verify header + let mut header = [0u8; HEADER_SIZE]; + reader.read_exact(&mut header)?; + + if &header[0..4] != MAGIC { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "Invalid sorted index magic", + )); + } + + let version = u32::from_le_bytes(header[4..8].try_into().unwrap_or_default()); + if version != VERSION { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + format!("Unsupported sorted index version: {version} (expected {VERSION})"), + )); + } + + let count = match header[8..16].try_into() { + Ok(count) => u64::from_le_bytes(count), + Err(e) => { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + format!("Invalid count {e}"), + )); + } + }; + + Ok(Self { + reader, + remaining: count, + _marker: PhantomData, + }) + } + + /// Returns the number of remaining entries to read. + pub fn remaining(&self) -> u64 { + self.remaining + } + + /// Returns true if the index is empty. + pub fn is_empty(&self) -> bool { + self.remaining == 0 + } + + /// Read the next entry from the index. + pub fn read_next(&mut self) -> io::Result>> { + if self.remaining == 0 { + return Ok(None); + } + + // Read sort key + let mut len_buf = [0u8; 4]; + self.reader.read_exact(&mut len_buf)?; + let sk_len = u32::from_le_bytes(len_buf) as usize; + + let mut sk_bytes = vec![0u8; sk_len]; + self.reader.read_exact(&mut sk_bytes)?; + let sort_key: SortKey = binary_deserialize(&sk_bytes)?; + + // Read primary key + self.reader.read_exact(&mut len_buf)?; + let pk_len = u32::from_le_bytes(len_buf) as usize; + + let mut pk_bytes = vec![0u8; pk_len]; + self.reader.read_exact(&mut pk_bytes)?; + let primary_key: K = binary_deserialize(&pk_bytes)?; + + // Read value location + let location = ValueLocation::from_reader(&mut self.reader)?; + + self.remaining -= 1; + Ok(Some(IndexEntry { + sort_key, + primary_key, + location, + })) + } +} + +impl Iterator for SortedIndexReader +where + SortKey: for<'de> Deserialize<'de>, + K: for<'de> Deserialize<'de>, +{ + type Item = io::Result>; + + fn next(&mut self) -> Option { + match self.read_next() { + Ok(Some(entry)) => Some(Ok(entry)), + Ok(None) => None, + Err(e) => Some(Err(e)), + } + } +} + +/// Owned iterator that combines a sorted index with a primary `BPlusTree` file. +/// +/// Iterates through the index in order and reads values directly from the +/// tree file using stored offsets - no tree traversal needed (O(1) per item). +pub struct BPlusTreeSortedIteratorOwned { + index_reader: SortedIndexReader, + tree_file: BufReader, + filepath: PathBuf, + block_cache: IndexMap>, + _marker: PhantomData, +} + +const CACHE_CAPACITY: usize = 8; + +impl BPlusTreeSortedIteratorOwned +where + K: for<'de> Deserialize<'de>, + V: for<'de> Deserialize<'de>, + SortKey: for<'de> Deserialize<'de>, +{ + /// Create a new owned sorted iterator. + /// + /// The index path is automatically derived from the tree filepath + /// by changing the extension to `.idx`. + pub fn new(filepath: PathBuf, tree_file: BufReader) -> io::Result { + let index_path = get_file_path_for_db_index(&filepath); + let index_reader = SortedIndexReader::open(&index_path)?; + Ok(Self { + index_reader, + tree_file, + filepath, + block_cache: IndexMap::with_capacity(CACHE_CAPACITY), + _marker: PhantomData, + }) + } + + /// Create from an explicit index path. + pub fn with_index_path( + filepath: PathBuf, + tree_file: BufReader, + index_path: &Path, + ) -> io::Result { + let index_reader = SortedIndexReader::open(index_path)?; + Ok(Self { + index_reader, + tree_file, + filepath, + block_cache: IndexMap::with_capacity(CACHE_CAPACITY), + _marker: PhantomData, + }) + } + + /// Returns the number of remaining items. + pub fn remaining(&self) -> u64 { + self.index_reader.remaining + } + + /// Returns the path to the tree file. + pub fn filepath(&self) -> &Path { + &self.filepath + } + + /// Read a value from the tree file using the given `ValueLocation`. + fn read_value(&mut self, location: ValueLocation) -> io::Result { + match location { + ValueLocation::Single { offset, length } => { + self.read_value_single(offset, length) + } + ValueLocation::Packed { block_offset, index, length: _ } => { + self.read_value_packed(block_offset, index) + } + } + } + + /// Read a single value directly from the tree file at the given offset. + fn read_value_single(&mut self, offset: u64, length: u32) -> io::Result { + self.tree_file.seek(SeekFrom::Start(offset))?; + + // Read compression flag + let mut flag = [0u8; 1]; + self.tree_file.read_exact(&mut flag)?; + + let data = if flag[0] == COMPRESSION_FLAG_LZ4 { + // Compressed: [flag:1][lz4_payload_with_prepended_size] + let compressed_len = length as usize - 1; + let mut compressed = vec![0u8; compressed_len]; + self.tree_file.read_exact(&mut compressed)?; + + lz4_flex::decompress_size_prepended(&compressed).map_err(|e| { + io::Error::new( + io::ErrorKind::InvalidData, + format!("LZ4 decompression failed: {e}"), + ) + })? + } else { + // Uncompressed: [flag:1][payload] + let payload_len = length as usize - 1; + let mut data = vec![0u8; payload_len]; + self.tree_file.read_exact(&mut data)?; + data + }; + + crate::utils::binary_deserialize(&data) + } + + /// Read a value from a packed block at the given index. + fn read_value_packed(&mut self, block_offset: u64, value_index: u16) -> io::Result { + + // Try cache first + if !self.block_cache.contains_key(&block_offset) { + // Miss - read from disk + self.tree_file.seek(SeekFrom::Start(block_offset))?; + let mut buf = vec![0u8; PAGE_SIZE_USIZE]; + self.tree_file.read_exact(&mut buf)?; + + // Update cache + if self.block_cache.len() >= CACHE_CAPACITY { + self.block_cache.shift_remove_index(0); // Remove LRU + } + self.block_cache.insert(block_offset, buf.clone()); + } + + let block_buffer = self.block_cache.get(&block_offset).unwrap(); + + // Read count (first 4 bytes) + let count = u32::from_le_bytes(block_buffer[0..4].try_into().map_err(|e| { + io::Error::new(io::ErrorKind::InvalidData, format!("Invalid count: {e}")) + })?); + let mut pos = 4; + + // Skip to target value + for i in 0..=value_index { + if pos + 4 > PAGE_SIZE_USIZE { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + format!("Packed block corrupted: position {pos} exceeds block size"), + )); + } + + let len = u32::from_le_bytes(block_buffer[pos..pos + 4].try_into().map_err(|e| { + io::Error::new(io::ErrorKind::InvalidData, format!("Invalid length: {e}")) + })?) as usize; + pos += 4; + + if i == value_index { + // Found target value + if pos + len > PAGE_SIZE_USIZE { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + format!("Packed value corrupted: length {len} at position {pos} exceeds block size"), + )); + } + let value_data = &block_buffer[pos..pos + len]; + return crate::utils::binary_deserialize(value_data); + } + + pos += len; + } + + Err(io::Error::new( + io::ErrorKind::InvalidData, + format!("Value index {value_index} not found in packed block (count: {count})"), + )) + } +} + +impl Iterator for BPlusTreeSortedIteratorOwned +where + K: for<'de> Deserialize<'de> + Clone, + V: for<'de> Deserialize<'de>, + SortKey: for<'de> Deserialize<'de>, +{ + type Item = io::Result<(K, V)>; + + fn next(&mut self) -> Option { + // Get next entry from index + let entry = match self.index_reader.read_next() { + Ok(Some(entry)) => entry, + Ok(None) => return None, + Err(e) => return Some(Err(e)), + }; + + // Read value from tree file using the stored location + match self.read_value(entry.location) { + Ok(value) => Some(Ok((entry.primary_key, value))), + Err(e) => Some(Err(e)), + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use tempfile::tempdir; + + #[test] + fn test_sorted_index_write_read() { + let dir = tempdir().unwrap(); + let path = dir.path().join("test.idx"); + + // Write entries with value locations + let mut writer = SortedIndexWriter::::new(&path).unwrap(); + writer.push( + &"apple".to_string(), + &1u32, + ValueLocation::Single { offset: 100, length: 50 } + ).unwrap(); + writer.push( + &"banana".to_string(), + &2u32, + ValueLocation::Packed { block_offset: 200, index: 3, length: 75 } + ).unwrap(); + writer.push( + &"cherry".to_string(), + &3u32, + ValueLocation::Single { offset: 300, length: 100 } + ).unwrap(); + let count = writer.finish().unwrap(); + assert_eq!(count, 3); + + // Read entries + let reader = SortedIndexReader::::open(&path).unwrap(); + let entries: Vec<_> = reader.map(|r| r.unwrap()).collect(); + + assert_eq!(entries.len(), 3); + assert_eq!(entries[0].sort_key, "apple".to_string()); + assert_eq!(entries[0].primary_key, 1u32); + match entries[0].location { + ValueLocation::Single { offset, length } => { + assert_eq!(offset, 100); + assert_eq!(length, 50); + } + _ => panic!("Expected Single location"), + } + + assert_eq!(entries[1].sort_key, "banana".to_string()); + assert_eq!(entries[1].primary_key, 2u32); + match entries[1].location { + ValueLocation::Packed { block_offset, index, length } => { + assert_eq!(block_offset, 200); + assert_eq!(index, 3); + assert_eq!(length, 75); + } + _ => panic!("Expected Packed location"), + } + + assert_eq!(entries[2].sort_key, "cherry".to_string()); + assert_eq!(entries[2].primary_key, 3u32); + match entries[2].location { + ValueLocation::Single { offset, length } => { + assert_eq!(offset, 300); + assert_eq!(length, 100); + } + _ => panic!("Expected Single location"), + } + } + + #[test] + fn test_empty_index() { + let dir = tempdir().unwrap(); + let path = dir.path().join("empty.idx"); + + let writer = SortedIndexWriter::::new(&path).unwrap(); + let count = writer.finish().unwrap(); + assert_eq!(count, 0); + + let reader = SortedIndexReader::::open(&path).unwrap(); + assert!(reader.is_empty()); + // Note: collect() consumes the reader, so we check is_empty() and remaining first + let remaining = reader.remaining; + assert_eq!(remaining, 0); + + let entries: Vec<_> = reader.collect(); + assert!(entries.is_empty()); + } + + #[test] + fn test_index_path_derivation() { + let tree_path = Path::new("/data/my_tree.bin"); + let idx_path = get_file_path_for_db_index(tree_path); + assert_eq!(idx_path, PathBuf::from("/data/my_tree.idx")); + + let tree_path2 = Path::new("/data/items"); + let idx_path2 = get_file_path_for_db_index(tree_path2); + assert_eq!(idx_path2, PathBuf::from("/data/items.idx")); + } +} diff --git a/backend/src/repository/storage.rs b/backend/src/repository/storage.rs index 6d8e115c2..189fd3f26 100644 --- a/backend/src/repository/storage.rs +++ b/backend/src/repository/storage.rs @@ -53,7 +53,10 @@ pub fn ensure_input_storage_path(cfg: &Config, input_name: &str) -> Result PathBuf { Path::new(working_dir).join("geoip.db") +} + +pub fn get_file_path_for_db_index(db_path: &Path) -> PathBuf { + db_path.with_extension(storage_const::FILE_SUFFIX_INDEX) } \ No newline at end of file diff --git a/backend/src/repository/storage_const.rs b/backend/src/repository/storage_const.rs index 46ac5cab9..ba143fd49 100644 --- a/backend/src/repository/storage_const.rs +++ b/backend/src/repository/storage_const.rs @@ -1,5 +1,6 @@ pub const FILE_EPG: &str = "epg.xml"; pub(in crate::repository) const FILE_SUFFIX_DB: &str = "db"; +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"; diff --git a/backend/src/repository/xtream_playlist_iterator.rs b/backend/src/repository/xtream_playlist_iterator.rs index 18fc1db99..8cd7841ef 100644 --- a/backend/src/repository/xtream_playlist_iterator.rs +++ b/backend/src/repository/xtream_playlist_iterator.rs @@ -1,6 +1,7 @@ use crate::model::ConfigTarget; use crate::model::{xtream_mapping_option_from_target_options, AppConfig, ProxyUserCredentials}; -use crate::repository::bplustree::{BPlusTreeDiskIteratorOwned, BPlusTreeQuery}; +use crate::repository::bplustree::{BPlusTreeQuery, PlaylistIteratorReader}; + use crate::repository::user_repository::user_get_bouquet_filter; use crate::repository::xtream_repository::{xtream_get_file_path, xtream_get_storage_path}; use crate::utils::FileReadGuard; @@ -8,9 +9,10 @@ use log::error; use shared::error::{TuliproxError, info_err, info_err_res}; use shared::model::{PlaylistItemType, TargetType, XtreamCluster, XtreamMappingOptions, XtreamPlaylistItem}; use std::collections::HashSet; +use crate::repository::storage::get_file_path_for_db_index; pub struct XtreamPlaylistIterator { - reader: BPlusTreeDiskIteratorOwned, + reader: PlaylistIteratorReader, options: XtreamMappingOptions, cluster: XtreamCluster, // Use parsed numeric filter to avoid per-item String allocations (no to_string per check) @@ -39,9 +41,25 @@ impl XtreamPlaylistIterator { } let file_lock = app_config.file_locks.read_lock(&xtream_path).await; - let query = BPlusTreeQuery::::try_new(&xtream_path) - .map_err(|err| info_err!("Could not open BPlusTreeQuery {xtream_path:?} - {err}"))?; - let reader = query.disk_iter(); + let index_path = get_file_path_for_db_index(&xtream_path); + let reader = if index_path.exists() { + let query = BPlusTreeQuery::::try_new(&xtream_path) + .map_err(|err| info_err!("Could not open BPlusTreeQuery {xtream_path:?} - {err}"))?; + match query.disk_iter_sorted() { + Ok(reader) => PlaylistIteratorReader::Sorted(reader), + Err(err) => { + error!("Sorted index error, falling back to unsorted: {err}"); + // Query was consumed, re-open for fallback + let query = BPlusTreeQuery::::try_new(&xtream_path) + .map_err(|err| info_err!("Could not open BPlusTreeQuery {xtream_path:?} - {err}"))?; + PlaylistIteratorReader::Unsorted(query.disk_iter()) + } + } + } else { + let query = BPlusTreeQuery::::try_new(&xtream_path) + .map_err(|err| info_err!("Could not open BPlusTreeQuery {xtream_path:?} - {err}"))?; + PlaylistIteratorReader::Unsorted(query.disk_iter()) + }; let server_info = app_config.get_user_server_info(user); let options = xtream_mapping_option_from_target_options(target, xtream_output, app_config, user, Some(server_info.get_base_url().as_str())); @@ -87,31 +105,37 @@ impl XtreamPlaylistIterator { true } - fn get_next(&mut self) -> Option<(XtreamPlaylistItem, bool)> { + /// Helper to find the next matching item from the iterator, handling `io::Result`. + fn find_next_matching(&mut self) -> Option<(XtreamPlaylistItem, bool)> { let filter_ids = self.filter_ids.as_ref(); let cluster = self.cluster; - let predicate = |(_, item): &(u32, XtreamPlaylistItem)| { - Self::matches_filters(cluster, filter_ids, item) - }; + loop { + match self.reader.next() { + Some(Ok((_, item))) => { + if Self::matches_filters(cluster, filter_ids, &item) { + return Some((item, false)); + } + // Continue to next item if filter doesn't match + } + Some(Err(err)) => { + error!("Error reading sorted index: {err}"); + return None; + } + None => return None, + } + } + } + fn get_next(&mut self) -> Option<(XtreamPlaylistItem, bool)> { if self.lookup_item.is_none() { - self.lookup_item = if self.cluster == XtreamCluster::Series || self.filter_ids.is_some() { - self.reader.find(predicate).map(|(_, item)| (item, false)) - } else { - self.reader.next().map(|(_, item)| (item, false)) - }; + self.lookup_item = self.find_next_matching(); } let (current_item, _) = self.lookup_item.take()?; // prefetch next - self.lookup_item = if self.cluster == XtreamCluster::Series || self.filter_ids.is_some() { - self.reader.find(predicate).map(|(_, item)| (item, false)) - } else { - self.reader.next().map(|(_, item)| (item, false)) - }; - + self.lookup_item = self.find_next_matching(); let has_next = self.lookup_item.is_some(); Some((current_item, has_next)) diff --git a/backend/src/repository/xtream_repository.rs b/backend/src/repository/xtream_repository.rs index 1d5083d17..faff4c68a 100644 --- a/backend/src/repository/xtream_repository.rs +++ b/backend/src/repository/xtream_repository.rs @@ -4,7 +4,7 @@ use crate::model::{AppConfig, ProxyUserCredentials}; use crate::model::{Config, ConfigTarget}; use crate::repository::bplustree::{BPlusTree, BPlusTreeQuery, BPlusTreeUpdate}; use crate::repository::playlist_scratch::PlaylistScratch; -use crate::repository::storage::{get_target_id_mapping_file, get_target_storage_path}; +use crate::repository::storage::{get_file_path_for_db_index, get_target_id_mapping_file, get_target_storage_path}; use crate::repository::storage_const; use crate::repository::target_id_mapping::VirtualIdRecord; use crate::repository::xtream_playlist_iterator::XtreamPlaylistJsonIterator; @@ -83,6 +83,7 @@ enum StorageKey { async fn write_playlists_to_file( app_config: &Arc, storage_path: &Path, + with_index: bool, storage_key: StorageKey, collections: Vec<(XtreamCluster, Vec)>, ) -> Result<(), TuliproxError> { @@ -100,7 +101,11 @@ async fn write_playlists_to_file( StorageKey::ProviderId => item.provider_id, }, item); } - tree.store(&xtream_path).map_err(|err| cant_write_result!(&xtream_path, err))?; + if with_index { + tree.store_with_index(&xtream_path, |pli| pli.source_ordinal).map_err(|err| cant_write_result!(&xtream_path, err))?; + } else { + tree.store(&xtream_path).map_err(|err| cant_write_result!(&xtream_path, err))?; + } } } Ok(()) @@ -297,6 +302,7 @@ pub async fn xtream_write_playlist( if let Err(err) = write_playlists_to_file( app_cfg, &path, + true, StorageKey::VirtualId, vec![ (XtreamCluster::Live, live_col.iter().map(|item| XtreamPlaylistItem::from(&**item)).collect::>()), @@ -582,11 +588,27 @@ pub async fn iter_raw_xtream_playlist(app_config: &AppConfig, target: &ConfigTar let file_lock = app_config.file_locks.read_lock(&xtream_path).await; match BPlusTreeQuery::::try_new(&xtream_path) .map_err(|err| info_err!("Could not open BPlusTreeQuery {xtream_path:?} - {err}")) { - Ok(mut query) => { - let items: Vec = query.iter().map(|(_, v)| v).collect(); - let len = items.len(); - Some((file_lock, items.into_iter().enumerate().map(move |(i, v)| (v, i < len - 1)))) - } + Ok(mut query) => { + let index_path = get_file_path_for_db_index(&xtream_path); + let items: Vec = if index_path.exists() { + match query.disk_iter_sorted::() { + Ok(iter) => iter.filter_map(Result::ok).map(|(_, v)| v).collect(), + Err(err) => { + error!("Sorted index error {}: {err}", xtream_path.display()); + // Re-open query for fallback + match BPlusTreeQuery::::try_new(&xtream_path) { + Ok(mut query) => query.iter().map(|(_, v)| v).collect(), + Err(_) => Vec::new(), + } + } + } + } else { + query.iter().map(|(_, v)| v).collect() + }; + + let len = items.len(); + Some((file_lock, items.into_iter().enumerate().map(move |(i, v)| (v, i < len - 1)))) + } Err(_) => None } } else { @@ -732,6 +754,7 @@ pub async fn persist_input_xtream_playlist(app_config: &Arc, storage_ if let Err(err) = write_playlists_to_file( app_config, storage_path, + false, StorageKey::ProviderId, vec![(cluster, col.iter().map(Into::into).collect::>())], ).await { @@ -869,7 +892,7 @@ pub async fn persist_input_series_info_batch(app_config: &Arc, storag } pub async fn load_input_xtream_playlist(app_config: &Arc, storage_path: &Path, clusters: &[XtreamCluster]) -> Result, TuliproxError> { - let mut groups: IndexMap = IndexMap::new(); + let mut groups: IndexMap<(XtreamCluster, u32), PlaylistGroup> = IndexMap::new(); for &cluster in clusters { let xtream_path = xtream_get_file_path(storage_path, cluster); @@ -885,7 +908,7 @@ pub async fn load_input_xtream_playlist(app_config: &Arc, storage_pat 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 { + groups.insert((cluster, cat.category_id), PlaylistGroup { id: cat.category_id, title: cat.category_name, channels: Vec::new(), @@ -901,7 +924,7 @@ pub async fn load_input_xtream_playlist(app_config: &Arc, storage_pat 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) + groups.entry((cluster, cat_id)) .or_insert_with(|| PlaylistGroup { id: cat_id, title: intern("Unknown"), diff --git a/backend/src/utils/network/ip_checker.rs b/backend/src/utils/network/ip_checker.rs index 1ee5dcdf4..ae6312bd9 100644 --- a/backend/src/utils/network/ip_checker.rs +++ b/backend/src/utils/network/ip_checker.rs @@ -1,9 +1,10 @@ +use std::sync::Arc; use shared::error::{TuliproxError, TuliproxErrorKind}; use crate::model::IpCheckConfig; use regex::Regex; use shared::utils::sanitize_sensitive_info; -async fn fetch_ip(client: &reqwest::Client, url: &str, regex: Option<&Regex>) -> Result { +async fn fetch_ip(client: &reqwest::Client, url: &str, regex: Option<&Arc>) -> Result { let response = client.get(url).send().await .map_err(|e| TuliproxError::new(TuliproxErrorKind::Info, format!("Failed to request {}: {e}", sanitize_sensitive_info(url))))?; diff --git a/frontend/scss/app/_component.scss b/frontend/scss/app/_component.scss index e21003fcc..da9d850e9 100644 --- a/frontend/scss/app/_component.scss +++ b/frontend/scss/app/_component.scss @@ -21,7 +21,6 @@ @forward "components/list_view"; @forward "components/reveal_content"; @forward "components/hide_content"; -@forward "components/userlist/userlist_view"; @forward "components/playlist/playlist_editor_view"; @forward "components/playlist/playlist_explorer_view"; @forward "components/playlist/playlist_explorer"; @@ -53,10 +52,12 @@ @forward "components/filter"; @forward "components/mapper_script"; @forward "components/mapper_counter"; +@forward "components/userlist/userlist_view"; @forward "components/userlist/user_state"; @forward "components/userlist/proxy_type"; @forward "components/userlist/max_connections"; -@forward "components/userlist/proxy-user-credentials-form"; +@forward "components/userlist/proxy_user_credentials_form"; +@forward "components/userlist/user_table"; @forward "components/radio_button_group"; @forward "components/no_content"; @forward "components/search"; diff --git a/frontend/scss/app/components/userlist/_proxy-user-credentials-form.scss b/frontend/scss/app/components/userlist/_proxy_user_credentials_form.scss similarity index 100% rename from frontend/scss/app/components/userlist/_proxy-user-credentials-form.scss rename to frontend/scss/app/components/userlist/_proxy_user_credentials_form.scss diff --git a/frontend/scss/app/components/userlist/_user_table.scss b/frontend/scss/app/components/userlist/_user_table.scss new file mode 100644 index 000000000..bc86c94cc --- /dev/null +++ b/frontend/scss/app/components/userlist/_user_table.scss @@ -0,0 +1,24 @@ +.tp__user-table { + display: flex; + flex-flow: column; + box-sizing: border-box; + overflow: hidden; + + .tp__icon-button { + width: 2rem; + } + + .tp__filter { + max-width: var(--max-table-cell-width); + overflow: hidden !important; + + .tp__filter__code { + overflow: hidden !important; + } + } + + &__invalid-target { + font-weight: bold; + background-color: var(--attention-color); + } +} diff --git a/frontend/scss/app/components/userlist/_userlist_view.scss b/frontend/scss/app/components/userlist/_userlist_view.scss index e5d6bf0d6..5b53efeee 100644 --- a/frontend/scss/app/components/userlist/_userlist_view.scss +++ b/frontend/scss/app/components/userlist/_userlist_view.scss @@ -30,8 +30,6 @@ } .tp__userlist-edit { - max-width: 800px; + max-width: 800px; } - - } diff --git a/frontend/src/app/components/search.rs b/frontend/src/app/components/search.rs index f7371d397..1f6dc276a 100644 --- a/frontend/src/app/components/search.rs +++ b/frontend/src/app/components/search.rs @@ -45,7 +45,7 @@ pub fn Search(props: &SearchProps) -> Html { RegexState::Inactive => { if let Some(input) = input.cast::() { let text = input.value(); - if regex::Regex::new(&text).is_ok() { + if shared::model::REGEX_CACHE.get_or_compile(&text).is_ok() { regex_active.set(RegexState::Active); } else { regex_active.set(RegexState::Invalid); @@ -83,7 +83,7 @@ pub fn Search(props: &SearchProps) -> Html { let text = input.value(); if text.len() >= min_length { if !matches!(*regex, RegexState::Inactive) { - if regex::Regex::new(&text).is_ok() { + if shared::model::REGEX_CACHE.get_or_compile(&text).is_ok() { regex.set(RegexState::Active); cb_search.emit(SearchRequest::Regexp(text, (*search_fields).clone())); } else { diff --git a/frontend/src/app/components/userlist/user_table.rs b/frontend/src/app/components/userlist/user_table.rs index de84d3642..479fcb85c 100644 --- a/frontend/src/app/components/userlist/user_table.rs +++ b/frontend/src/app/components/userlist/user_table.rs @@ -1,7 +1,9 @@ use std::cmp::Ordering; use crate::app::components::menu_item::MenuItem; use crate::app::components::popup_menu::PopupMenu; -use crate::app::components::{convert_bool_to_chip_style, AppIcon, CellValue, Chip, HideContent, MaxConnections, ProxyTypeView, RevealContent, Table, TableDefinition, UserStatus, UserlistContext, UserlistPage}; +use crate::app::components::{convert_bool_to_chip_style, AppIcon, CellValue, Chip, HideContent, + MaxConnections, ProxyTypeView, RevealContent, Table, TableDefinition, + UserStatus, UserlistContext, UserlistPage}; use crate::app::context::TargetUser; use crate::model::DialogResult; use crate::services::DialogService; @@ -11,11 +13,12 @@ use shared::utils::{unix_ts_to_str, Substring}; use std::fmt::Display; use std::rc::Rc; use std::str::FromStr; +use std::collections::HashSet; use yew::platform::spawn_local; use yew::prelude::*; use yew_hooks::use_clipboard; use yew_i18n::use_translation; -use crate::app::TargetUserList; +use crate::app::{ConfigContext, TargetUserList}; use crate::hooks::use_service_context; const HEADERS: [&str; 15] = [ @@ -102,12 +105,21 @@ pub fn UserTable(props: &UserTableProps) -> Html { let translate = use_translation(); let clipboard = use_clipboard(); let service_ctx = use_service_context(); + let config_ctx = use_context::().expect("Config context not found"); let dialog = use_context::().expect("Dialog service not found"); let userlist_context = use_context::().expect("Userlist context not found"); let popup_anchor_ref = use_state(|| None::); let popup_is_open = use_state(|| false); let selected_dto = use_state(|| None::>); let user_list = use_state(|| props.users.clone()); + let target_names = use_memo(config_ctx.clone(), |cfg| + cfg.config.as_ref().map(|c| c.sources.sources.iter().flat_map(|s| s.targets.iter()) + .map(|t| t.name.clone()) + .collect::>() + ) + .unwrap_or_default() + ); + { let user_list = user_list.clone(); @@ -156,6 +168,7 @@ pub fn UserTable(props: &UserTableProps) -> Html { let render_data_cell = { let translator = translate.clone(); let popup_onclick = handle_popup_onclick.clone(); + let target_names = target_names.clone(); Callback::<(usize, usize, Rc), Html>::from( move |(row, col, dto): (usize, usize, Rc)| { let user_active = dto.credentials.is_active(); @@ -174,7 +187,7 @@ pub fn UserTable(props: &UserTableProps) -> Html { label={if user_active {translator.t("LABEL.ENABLED")} else { translator.t("LABEL.DISABLED")} } /> }, 2 => html! { }, - 3 => html! { dto.target.as_str() }, + 3 => html! { {dto.target.as_str()} }, 4 => html! { dto.credentials.username.as_str() }, 5 => html! { }, 6 => html! { dto.credentials.token.as_ref().map_or_else(|| html!{}, |token| html! { }) }, @@ -317,7 +330,7 @@ pub fn UserTable(props: &UserTableProps) -> Html { }; html! { -
+
{ html! { <> diff --git a/frontend/src/app/context.rs b/frontend/src/app/context.rs index 6f23f4e6b..792c2115b 100644 --- a/frontend/src/app/context.rs +++ b/frontend/src/app/context.rs @@ -1,5 +1,4 @@ use std::rc::Rc; -use regex::Regex; use yew::UseStateHandle; use shared::model::{AppConfigDto, ConfigTargetDto, PlaylistRequest, ProxyUserCredentialsDto, SearchRequest, StatusCheck, SystemInfo, UiPlaylistCategories}; use crate::app::components::{InputRow, PlaylistEditorPage, PlaylistExplorerPage, UserlistPage}; @@ -23,7 +22,6 @@ pub struct PlaylistExplorerContext { pub playlist_request: UseStateHandle>, } - #[derive(Clone, PartialEq)] pub struct TargetUser { pub target: String, @@ -79,7 +77,7 @@ impl UserlistContext { self.filtered_users.set(Some(Rc::new(filtered))); } SearchRequest::Regexp(text, search_fields) => { - if let Ok(regex) = Regex::new(text) { + if let Ok(regex) = shared::model::REGEX_CACHE.get_or_compile(text) { let filter_username = search_fields.as_ref().is_none_or(|f| f.iter().any(|s| s == "username")); let filter_server = search_fields.as_ref().is_none_or(|f| f.iter().any(|s| s == "server")); let filter_playlist = search_fields.as_ref().is_none_or(|f| f.iter().any(|s| s == "playlist")); diff --git a/shared/Cargo.toml b/shared/Cargo.toml index e1c55c9f0..46e175e6f 100644 --- a/shared/Cargo.toml +++ b/shared/Cargo.toml @@ -26,7 +26,8 @@ ciborium = "0.2.2" hex = "0.4.3" lz4_flex = "0.12.0" deunicode = "1.6.2" -serde-saphyr = "0.0.13" +serde-saphyr = "0.0.14" +dashmap = "6.1.0" [target.'cfg(target_arch = "wasm32")'.dependencies] js-sys = "0.3.83" diff --git a/shared/src/foundation/filter.rs b/shared/src/foundation/filter.rs index 1f366a52c..435ad1d60 100644 --- a/shared/src/foundation/filter.rs +++ b/shared/src/foundation/filter.rs @@ -9,6 +9,7 @@ use pest::iterators::Pair; use pest::Parser; use std::cmp::Ordering; use std::collections::HashMap; +use std::sync::Arc; use crate::error::{info_err_res, TuliproxError}; use crate::info_err; pub use crate::model::{ItemField, PatternTemplate, PlaylistItem, PlaylistItemType, FieldGetAccessor, FieldSetAccessor, TemplateValue}; @@ -86,7 +87,7 @@ impl ValueAccessor<'_> { #[derive(Debug, Clone)] pub struct CompiledRegex { pub restr: String, - pub re: regex::Regex, + pub re: Arc, } impl PartialEq for CompiledRegex { @@ -157,7 +158,7 @@ pub enum Filter { impl Default for Filter { fn default() -> Self { - Self::Group(Box::new(Filter::FieldComparison(ItemField::Group, CompiledRegex { restr: ".*".to_string(), re: regex::Regex::new(".*").unwrap() }))) + Self::Group(Box::new(Filter::FieldComparison(ItemField::Group, CompiledRegex { restr: ".*".to_string(), re: crate::model::REGEX_CACHE.get_or_compile(".*").unwrap() }))) } } @@ -296,7 +297,7 @@ fn get_parser_regexp( parsed_text.pop(); parsed_text.remove(0); let regstr = apply_templates_to_pattern_single(&parsed_text, templates)?; - let re = regex::Regex::new(regstr.as_str()); + let re = crate::model::REGEX_CACHE.get_or_compile(regstr.as_str()); if re.is_err() { return info_err_res!("can't parse regex: {}", regstr); } diff --git a/shared/src/foundation/mapper.rs b/shared/src/foundation/mapper.rs index 1aec08928..edecad895 100644 --- a/shared/src/foundation/mapper.rs +++ b/shared/src/foundation/mapper.rs @@ -17,6 +17,7 @@ use std::fmt::Display; use std::fmt::Write; use std::ops::Deref; use std::str::FromStr; +use std::sync::Arc; #[derive(Parser)] #[grammar_inline = r##" @@ -181,7 +182,7 @@ pub enum Expression { NumberLiteral(f64), FieldAccess(String), VarAccess(String, String), - RegexExpr { field: RegexSource, pattern: String, re_pattern: Regex }, + RegexExpr { field: RegexSource, pattern: String, re_pattern: Arc }, FunctionCall { name: BuiltInFunction, args: Vec }, Assignment { target: AssignmentTarget, expr: ExprId }, MatchBlock(Vec), @@ -506,7 +507,7 @@ impl MapperScript { }; let pattern_raw = inner.next().unwrap().as_str(); let pattern = &pattern_raw[1..pattern_raw.len() - 1]; // Strip quotes - match Regex::new(pattern) { + match crate::model::REGEX_CACHE.get_or_compile(pattern) { Ok(re) => Ok(Some(Expression::RegexExpr { field, pattern: pattern.to_string(), re_pattern: re })), Err(_) => info_err_res!("Invalid regex {}", pattern), } diff --git a/shared/src/model/config/base.rs b/shared/src/model/config/base.rs index a3dff4a01..c463788ca 100644 --- a/shared/src/model/config/base.rs +++ b/shared/src/model/config/base.rs @@ -227,7 +227,7 @@ impl ConfigDto { if let Some(download) = &video.download { if let Some(episode_pattern) = &download.episode_pattern { if !episode_pattern.is_empty() { - let re = regex::Regex::new(episode_pattern); + let re = crate::model::REGEX_CACHE.get_or_compile(episode_pattern); if re.is_err() { return false; } diff --git a/shared/src/model/config/epg_smart_match.rs b/shared/src/model/config/epg_smart_match.rs index d97033773..f9e65ae70 100644 --- a/shared/src/model/config/epg_smart_match.rs +++ b/shared/src/model/config/epg_smart_match.rs @@ -96,7 +96,7 @@ impl EpgSmartMatchConfigDto { } if let Some(regstr) = self.normalize_regex.as_ref() { - let re = regex::Regex::new(regstr.as_str()); + let re = crate::model::REGEX_CACHE.get_or_compile(regstr.as_str()); if re.is_err() { return info_err_res!("can't parse regex: {}", regstr); } diff --git a/shared/src/model/config/favourites.rs b/shared/src/model/config/favourites.rs index 3a809252a..b7c39c957 100644 --- a/shared/src/model/config/favourites.rs +++ b/shared/src/model/config/favourites.rs @@ -1,12 +1,14 @@ use std::sync::Arc; use crate::error::{TuliproxError}; use crate::foundation::filter::{get_filter, Filter}; -use crate::model::{PatternTemplate}; -use crate::utils::arc_str_serde; +use crate::model::{PatternTemplate, XtreamCluster}; +use crate::utils::{arc_str_serde, xtream_cluster_serde}; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq)] #[serde(deny_unknown_fields)] pub struct ConfigFavouritesDto { + #[serde(with = "xtream_cluster_serde")] + pub cluster: XtreamCluster, #[serde(with = "arc_str_serde")] pub group: Arc, #[serde(default)] diff --git a/shared/src/model/config/ipcheck.rs b/shared/src/model/config/ipcheck.rs index d8c046822..67951bd96 100644 --- a/shared/src/model/config/ipcheck.rs +++ b/shared/src/model/config/ipcheck.rs @@ -1,6 +1,5 @@ use crate::error::{TuliproxError, TuliproxErrorKind}; use crate::utils::{is_blank_optional_str, is_blank_optional_string}; -use regex::Regex; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default, PartialEq)] #[serde(deny_unknown_fields)] @@ -68,7 +67,7 @@ impl IpCheckConfigDto { // } if let Some(p4) = &self.pattern_ipv4 { - Regex::new(p4).map_err(|err| { + crate::model::REGEX_CACHE.get_or_compile(p4).map_err(|err| { TuliproxError::new( TuliproxErrorKind::Info, format!("Invalid IPv4 regex: {p4} {err}"), @@ -76,7 +75,7 @@ impl IpCheckConfigDto { })?; } if let Some(p6) = &self.pattern_ipv6 { - Regex::new(p6).map_err(|err| { + crate::model::REGEX_CACHE.get_or_compile(p6).map_err(|err| { TuliproxError::new( TuliproxErrorKind::Info, format!("Invalid IPv6 regex: {p6} {err}"), diff --git a/shared/src/model/config/mod.rs b/shared/src/model/config/mod.rs index f8210060a..9a95e570f 100644 --- a/shared/src/model/config/mod.rs +++ b/shared/src/model/config/mod.rs @@ -36,6 +36,7 @@ mod proxy_user_status; mod favourites; mod geoip; mod library; + pub use proxy_type::*; pub use proxy_user_status::*; pub use base::*; diff --git a/shared/src/model/config/rename.rs b/shared/src/model/config/rename.rs index 65822a876..9bc6a5cf4 100644 --- a/shared/src/model/config/rename.rs +++ b/shared/src/model/config/rename.rs @@ -13,7 +13,7 @@ pub struct ConfigRenameDto { impl ConfigRenameDto { pub fn prepare(&mut self, templates: Option<&Vec>) -> Result<(), TuliproxError> { self.pattern = apply_templates_to_pattern_single(&self.pattern, templates)?; - if let Err(err) = regex::Regex::new(&self.pattern) { + if let Err(err) = crate::model::REGEX_CACHE.get_or_compile(&self.pattern) { return info_err_res!("can't parse regex: {} {err}", &self.pattern); } Ok(()) diff --git a/shared/src/model/config/sort.rs b/shared/src/model/config/sort.rs index b6a84ddab..c40f98a00 100644 --- a/shared/src/model/config/sort.rs +++ b/shared/src/model/config/sort.rs @@ -6,12 +6,13 @@ use regex::Regex; use serde::{Deserialize, Deserializer, Serialize}; use std::fmt::{Display, Formatter}; use std::str::FromStr; +use std::sync::Arc; -fn compile_regex_vec(patterns: Option<&Vec>) -> Result>, TuliproxError> { +fn compile_regex_vec(patterns: Option<&Vec>) -> Result>>, TuliproxError> { patterns.as_ref() .map(|seq| { seq.iter() - .map(|s| Regex::new(s).map_err(|err| { + .map(|s| crate::model::REGEX_CACHE.get_or_compile(s).map_err(|err| { info_err!("can't parse regex: {s} {err}") })) .collect::, _>>() @@ -107,7 +108,7 @@ pub struct ConfigSortRuleDto { pub sequence: Option>, pub filter: String, #[serde(skip)] - pub t_sequence: Option>, + pub t_sequence: Option>>, #[serde(skip)] pub t_filter: Option, } diff --git a/shared/src/model/config/target.rs b/shared/src/model/config/target.rs index 8d965dc74..253cbf150 100644 --- a/shared/src/model/config/target.rs +++ b/shared/src/model/config/target.rs @@ -379,7 +379,7 @@ impl ConfigTargetDto { if let Some(watch) = &self.watch { for pat in watch { - if let Err(err) = regex::Regex::new(pat) { + if let Err(err) = crate::model::REGEX_CACHE.get_or_compile(pat) { return info_err_res!("Invalid watch regular expression: {}", err); } } diff --git a/shared/src/model/config/video_download.rs b/shared/src/model/config/video_download.rs index 3b3af23c1..96e8936a3 100644 --- a/shared/src/model/config/video_download.rs +++ b/shared/src/model/config/video_download.rs @@ -67,7 +67,7 @@ impl VideoConfigDto { } if let Some(episode_pattern) = &downl.episode_pattern { - if let Err(err) = regex::Regex::new(episode_pattern) { + if let Err(err) = crate::model::REGEX_CACHE.get_or_compile(episode_pattern) { return info_err_res!("can't parse regex: {episode_pattern} {err}"); } } diff --git a/shared/src/model/mod.rs b/shared/src/model/mod.rs index 05caf6981..bc022be20 100644 --- a/shared/src/model/mod.rs +++ b/shared/src/model/mod.rs @@ -28,6 +28,7 @@ mod stream_properties; mod xtream; mod info_doc_utils; mod playlist_info_document; +mod regex_cache; pub use self::cluster_flags::*; pub use self::playlist::*; @@ -55,4 +56,5 @@ pub use self::system_info::*; pub use self::library_request::*; pub use self::stream_properties::*; pub use self::xtream::*; -pub use self::playlist_info_document::*; \ No newline at end of file +pub use self::playlist_info_document::*; +pub use self::regex_cache::*; \ No newline at end of file diff --git a/shared/src/model/playlist.rs b/shared/src/model/playlist.rs index d063cd42d..d9bf84947 100644 --- a/shared/src/model/playlist.rs +++ b/shared/src/model/playlist.rs @@ -1232,7 +1232,6 @@ pub struct PlaylistGroup { #[serde(with = "arc_str_serde")] pub title: Arc, pub channels: Vec, - #[serde(skip)] pub xtream_cluster: XtreamCluster, } diff --git a/shared/src/model/playlist_request.rs b/shared/src/model/playlist_request.rs index 36df924ab..000660df8 100644 --- a/shared/src/model/playlist_request.rs +++ b/shared/src/model/playlist_request.rs @@ -211,7 +211,7 @@ impl UiPlaylistCategories { build_result(live, video, series) } SearchRequest::Regexp(text, _search_fields) => { - if let Ok(regex) = Regex::new(text) { + if let Ok(regex) = crate::model::REGEX_CACHE.get_or_compile(text) { let live = filter_channels_re(self.live.as_ref(), ®ex); let video = filter_channels_re(self.vod.as_ref(), ®ex); let series = filter_channels_re(self.series.as_ref(), ®ex); diff --git a/shared/src/model/regex_cache.rs b/shared/src/model/regex_cache.rs new file mode 100644 index 000000000..712215f0a --- /dev/null +++ b/shared/src/model/regex_cache.rs @@ -0,0 +1,50 @@ +use std::sync::{Arc, LazyLock}; +use dashmap::DashMap; +use regex::Regex; +use crate::error::TuliproxError; +use crate::info_err; + +pub static REGEX_CACHE: LazyLock = LazyLock::new(RegexCache::new); + +pub struct RegexCache { + cache: DashMap>, +} + +impl Default for RegexCache { + fn default() -> Self { + Self::new() + } +} + +impl RegexCache { + pub fn new() -> Self { + Self { + cache: DashMap::new(), + } + } + + pub fn get_or_compile( + &self, + pattern: &str, + ) -> Result, TuliproxError> { + // Try to get existing entry first + if let Some(cached) = self.cache.get(pattern) { + return Ok(cached.clone()); + } + // Compile outside the lock + let regex = Regex::new(pattern).map_err(|e| { + info_err!("can't parse regex: {pattern} {e}") + })?; + let arc_regex = Arc::new(regex); + // Use entry API to avoid overwriting if another thread inserted + Ok(self.cache.entry(pattern.to_owned()) + .or_insert(arc_regex) + .clone()) + } + + + /// Removes regexes that are only held by the cache itself (strong_count == 1). + pub fn sweep(&self) { + self.cache.retain(|_k, v| Arc::strong_count(v) > 1); + } +} \ No newline at end of file diff --git a/shared/src/utils/constants.rs b/shared/src/utils/constants.rs index 9ebd9f9ff..843670f48 100644 --- a/shared/src/utils/constants.rs +++ b/shared/src/utils/constants.rs @@ -2,7 +2,7 @@ use regex::Regex; use std::collections::HashSet; use std::string::ToString; use std::sync::atomic::AtomicBool; -use std::sync::LazyLock; +use std::sync::{Arc, LazyLock}; pub const USER_FILE: &str = "user.txt"; @@ -83,7 +83,7 @@ pub struct Constants { pub re_base_href_wasm: Regex, pub re_env_var: Regex, pub re_memory_usage: Regex, - pub re_epg_normalize: Regex, + pub re_epg_normalize: Arc, pub re_template_var: Regex, pub re_template_tag: Regex, pub re_template_attribute: Regex, @@ -119,7 +119,7 @@ pub static CONSTANTS: LazyLock = LazyLock::new(|| re_base_href_wasm: Regex::new("'(/frontend\\-)").unwrap(), re_env_var: Regex::new(r"\$\{env:(?P[a-zA-Z_][a-zA-Z0-9_]*)}").unwrap(), re_memory_usage: Regex::new(r"VmRSS:\s+(\d+) kB").unwrap(), - re_epg_normalize: Regex::new(r"[^a-zA-Z0-9\-]").unwrap(), + re_epg_normalize: Arc::new(Regex::new(r"[^a-zA-Z0-9\-]").unwrap()), re_template_var: Regex::new("!(.*?)!").unwrap(), re_template_tag: Regex::new("").unwrap(), re_template_attribute: Regex::new("<(.*?)>").unwrap(), diff --git a/shared/src/utils/serde_utils.rs b/shared/src/utils/serde_utils.rs index 9b6ff0e22..427288b79 100644 --- a/shared/src/utils/serde_utils.rs +++ b/shared/src/utils/serde_utils.rs @@ -398,3 +398,28 @@ where { serializer.serialize_str(&value.to_string()) } + + +/// Serde support for `XtreamCluster` fields. +/// Serializes as string (e.g., "live", "video", "series") and deserializes via `FromStr`. +pub mod xtream_cluster_serde { + use std::str::FromStr; + use serde::{Deserialize, Deserializer, Serializer}; + use serde::de::Error; + use crate::model::XtreamCluster; + + pub fn serialize(value: &XtreamCluster, serializer: S) -> Result + where + S: Serializer, + { + serializer.serialize_str(value.as_str()) + } + + pub fn deserialize<'de, D>(deserializer: D) -> Result + where + D: Deserializer<'de>, + { + let raw = String::deserialize(deserializer)?; + XtreamCluster::from_str(&raw).map_err(D::Error::custom) + } +} \ No newline at end of file