2025-06-24 01:06:25 +02:00
use crate ::model ::{ AppConfig , ConfigInput , ConfigRename };
2025-05-04 10:59:08 +02:00
use crate ::utils ::epg ;
use crate ::utils ::m3u ;
use crate ::utils ::xtream ;
2025-04-30 17:26:07 +02:00
use crate ::Config ;
2023-10-27 19:03:29 +02:00
use std ::collections ::{ HashMap , HashSet };
2025-01-06 09:17:05 +01:00
use std ::path ::PathBuf ;
2025-04-30 17:26:07 +02:00
use std ::sync ::Arc ;
2023-02-13 12:19:58 +01:00
use std ::thread ;
2025-04-30 17:26:07 +02:00
use tokio ::sync ::Mutex ;
2023-10-19 18:56:08 +02:00
2025-06-24 01:06:25 +02:00
use crate ::messaging ::send_message ;
use crate ::model ::Epg ;
use crate ::model ::{ ConfigTarget , ProcessTargets };
use crate ::model ::{ Mapping };
use crate ::model ::FetchedPlaylist ;
2025-05-04 10:59:08 +02:00
use crate ::model ::{ InputStats , PlaylistStats , SourceStats , TargetStats };
2025-06-24 01:06:25 +02:00
use crate ::processing ::parser ::xmltv ::flatten_tvguide ;
2023-10-13 18:13:16 +02:00
use crate ::processing ::playlist_watch ::process_group_watch ;
2025-06-24 01:06:25 +02:00
use crate ::processing ::processor ::epg ::process_playlist_epg ;
use crate ::processing ::processor ::sort ::sort_playlist ;
2025-06-16 17:53:03 +02:00
use crate ::processing ::processor ::trakt ::process_trakt_categories_for_target ;
2025-06-24 01:06:25 +02:00
use crate ::processing ::processor ::xtream_series ::playlist_resolve_series ;
use crate ::processing ::processor ::xtream_vod ::playlist_resolve_vod ;
2024-05-04 20:30:20 +02:00
use crate ::repository ::playlist_repository ::persist_playlist ;
2025-04-30 17:26:07 +02:00
use crate ::utils ::debug_if_enabled ;
2025-06-24 01:06:25 +02:00
use crate ::utils ::StepMeasure ;
2025-04-30 17:26:07 +02:00
use deunicode ::deunicode ;
use log ::{ debug , error , info , log_enabled , trace , warn , Level };
2025-06-16 17:53:03 +02:00
use reqwest ::Client ;
2025-08-13 23:07:21 +02:00
use shared ::error ::{ get_errors_notify_message , notify_err , TuliproxError };
2025-06-24 01:06:25 +02:00
use shared ::foundation ::filter ::{ get_field_value , set_field_value , ValueAccessor , ValueProvider };
2025-08-04 14:18:46 +02:00
use shared ::model ::{ CounterModifier , FieldGetAccessor , FieldSetAccessor , InputType , ItemField , MsgKind , PlaylistEntry ,
PlaylistGroup , PlaylistItem , PlaylistUpdateState , ProcessingOrder , UUIDType , XtreamCluster };
2025-06-24 01:06:25 +02:00
use shared ::utils ::default_as_default ;
use std ::time ::Instant ;
2025-08-03 10:24:23 +02:00
use crate ::api ::model ::{ EventManager , EventMessage };
2025-04-03 18:03:59 +03:00
2024-05-04 20:30:20 +02:00
fn is_valid ( pli : & PlaylistItem , target : & ConfigTarget ) -> bool {
2025-03-11 14:56:10 +01:00
let provider = ValueProvider { pli };
2024-05-04 20:30:20 +02:00
target . filter ( & provider )
}
2021-10-15 16:59:20 +02:00
2024-05-10 12:01:48 +02:00
#[allow(clippy::unnecessary_wraps)]
2023-10-08 19:54:59 +02:00
fn filter_playlist ( playlist : & mut [ PlaylistGroup ], target : & ConfigTarget ) -> Option < Vec < PlaylistGroup >> {
2023-10-28 22:20:57 +02:00
debug! ( "Filtering {} groups" , playlist . len ());
2025-01-09 13:31:05 +01:00
let mut new_playlist = Vec ::with_capacity ( 128 );
2024-05-10 12:01:48 +02:00
for pg in playlist . iter_mut () {
2024-05-04 20:30:20 +02:00
let channels = pg . channels . iter ()
. filter ( |& pli | is_valid ( pli , target )). cloned (). collect ::< Vec < PlaylistItem >> ();
2024-11-03 14:52:56 +01:00
trace! ( "Filtered group {} has now {}/{} items" , pg . title , channels . len (), pg . channels . len ());
2023-05-03 17:46:39 +02:00
if ! channels . is_empty () {
2023-01-13 15:56:03 +01:00
new_playlist . push ( PlaylistGroup {
2023-10-06 19:24:06 +02:00
id : pg . id ,
2023-01-13 15:56:03 +01:00
title : pg . title . clone (),
channels ,
2024-09-10 20:05:29 +02:00
xtream_cluster : pg . xtream_cluster ,
2023-01-13 15:56:03 +01:00
});
}
2024-05-10 12:01:48 +02:00
}
2023-01-13 15:56:03 +01:00
Some ( new_playlist )
}
2024-05-05 19:15:54 +02:00
2025-05-04 10:59:08 +02:00
fn assign_channel_no_playlist ( new_playlist : & mut [ PlaylistGroup ]) {
2025-04-12 18:20:46 +02:00
let assigned_chnos : HashSet < u32 > = new_playlist . iter (). flat_map ( | g | & g . channels )
. filter ( | c | ! c . header . chno . is_empty ())
2025-04-30 17:26:07 +02:00
. map ( | c | c . header . chno . as_str ())
2025-04-12 20:03:03 +02:00
. flat_map ( str ::parse ::< u32 > ). collect ();
2025-01-29 17:33:06 +01:00
let mut chno = 1 ;
for group in new_playlist {
2025-03-11 14:56:10 +01:00
for chan in & mut group . channels {
2025-04-12 18:20:46 +02:00
if chan . header . chno . is_empty () {
while assigned_chnos . contains ( & chno ) {
chno += 1 ;
}
chan . header . chno = chno . to_string ();
chno += 1 ;
}
2025-01-29 17:33:06 +01:00
}
}
}
2025-03-11 14:56:10 +01:00
fn exec_rename ( pli : & mut PlaylistItem , rename : Option <& Vec < ConfigRename >> ) {
2023-05-03 17:46:39 +02:00
if let Some ( renames ) = rename {
if ! renames . is_empty () {
let result = pli ;
for r in renames {
2025-05-26 13:57:29 +02:00
let value = get_field_value ( result , r . field );
2025-06-24 01:06:25 +02:00
let cap = r . pattern . replace_all ( value . as_str (), & r . new_name );
2025-05-12 21:38:02 +02:00
if log_enabled! ( log ::Level ::Debug ) && * value != cap {
2025-01-01 22:15:49 +01:00
debug_if_enabled! ( "Renamed {}={} to {}" , & r . field , value , cap );
2025-01-01 10:59:51 +01:00
}
2023-05-03 17:46:39 +02:00
let value = cap . into_owned ();
2025-05-26 13:57:29 +02:00
set_field_value ( result , r . field , value );
2023-01-13 15:56:03 +01:00
}
}
}
}
2023-10-08 19:54:59 +02:00
fn rename_playlist ( playlist : & mut [ PlaylistGroup ], target : & ConfigTarget ) -> Option < Vec < PlaylistGroup >> {
2023-01-13 15:56:03 +01:00
match & target . rename {
Some ( renames ) => {
2023-05-03 17:46:39 +02:00
if ! renames . is_empty () {
2025-01-09 13:31:05 +01:00
let mut new_playlist : Vec < PlaylistGroup > = Vec ::with_capacity ( playlist . len ());
2023-01-13 15:56:03 +01:00
for g in playlist {
let mut grp = g . clone ();
2023-01-12 15:59:40 +01:00
for r in renames {
2024-11-04 18:46:56 +01:00
if matches! ( r . field , ItemField ::Group ) {
2025-06-24 01:06:25 +02:00
let cap = r . pattern . replace_all ( & grp . title , & r . new_name );
2024-12-09 19:26:39 +01:00
debug_if_enabled! ( "Renamed group {} to {} for {}" , & grp . title , cap , target . name );
2025-03-11 14:56:10 +01:00
grp . title = cap . into_owned ();
2023-01-12 15:59:40 +01:00
}
}
2023-01-13 15:56:03 +01:00
2024-12-02 17:24:06 +01:00
grp . channels . iter_mut (). for_each ( | pli | exec_rename ( pli , target . rename . as_ref ()));
2023-01-13 15:56:03 +01:00
new_playlist . push ( grp );
2023-01-12 15:59:40 +01:00
}
2023-02-13 12:19:58 +01:00
return Some ( new_playlist );
2023-01-12 15:59:40 +01:00
}
2023-01-13 15:56:03 +01:00
None
2023-01-12 15:59:40 +01:00
}
2023-01-13 15:56:03 +01:00
_ => None
2023-01-12 15:59:40 +01:00
}
}
2025-03-11 14:56:10 +01:00
fn map_channel ( mut channel : PlaylistItem , mapping : & Mapping ) -> PlaylistItem {
2025-05-14 18:23:22 +02:00
if let Some ( mapper ) = & mapping . mapper {
if ! mapper . is_empty () {
let header = & channel . header ;
2025-08-26 13:53:11 +02:00
let channel_name = if mapping . match_as_ascii { deunicode ( & header . name ) } else { header . name . clone () };
2025-05-14 18:23:22 +02:00
if mapping . match_as_ascii && log_enabled! ( Level ::Trace ) { trace! ( "Decoded {} for matching to {}" , & header . name , & channel_name ); }
let ref_chan = & mut channel ;
2025-06-12 16:20:15 +02:00
let templates = mapping . templates . as_ref ();
2025-05-14 18:23:22 +02:00
for m in mapper {
2025-05-22 13:34:45 +02:00
if let Some ( script ) = m . t_script . as_ref () {
if let Some ( filter ) = & m . t_filter {
2025-06-12 17:27:59 +02:00
let provider = ValueProvider { pli : ref_chan };
2025-05-22 13:34:45 +02:00
if filter . filter ( & provider ) {
2025-06-12 17:27:59 +02:00
let mut accessor = ValueAccessor { pli : ref_chan };
2025-06-12 16:20:15 +02:00
script . eval ( & mut accessor , templates );
2025-05-14 18:23:22 +02:00
}
}
2022-04-05 12:23:30 +02:00
}
2025-04-10 00:06:48 +02:00
}
2023-01-12 15:59:40 +01:00
}
}
2023-10-23 18:12:24 +02:00
channel
2023-01-12 15:59:40 +01:00
}
2023-10-08 19:54:59 +02:00
fn map_playlist ( playlist : & mut [ PlaylistGroup ], target : & ConfigTarget ) -> Option < Vec < PlaylistGroup >> {
2025-06-24 01:06:25 +02:00
if let Some ( mappings ) = target . mapping . load (). as_ref () {
2023-05-03 17:46:39 +02:00
let new_playlist : Vec < PlaylistGroup > = playlist . iter (). map ( | playlist_group | {
2023-02-14 18:18:19 +01:00
let mut grp = playlist_group . clone ();
2025-05-14 19:11:43 +02:00
mappings . iter (). filter ( |& mapping | mapping . mapper . as_ref (). is_some_and ( | v | ! v . is_empty ()))
2025-05-14 18:23:22 +02:00
. for_each ( | mapping |
2025-05-23 18:17:47 +02:00
grp . channels = grp . channels . drain ( .. ). map ( | chan | map_channel ( chan , mapping )). collect ());
2023-05-03 17:46:39 +02:00
grp
}). collect ();
2023-04-26 14:30:43 +02:00
// if the group names are changed, restructure channels to the right groups
// we use
2025-01-09 13:31:05 +01:00
let mut new_groups : Vec < PlaylistGroup > = Vec ::with_capacity ( 128 );
2024-03-28 16:27:39 +01:00
let mut grp_id : u32 = 0 ;
2023-04-26 14:30:43 +02:00
for playlist_group in new_playlist {
for channel in & playlist_group . channels {
2025-03-11 14:56:10 +01:00
let cluster = & channel . header . xtream_cluster ;
let title = & channel . header . group ;
2024-05-10 12:01:48 +02:00
if let Some ( grp ) = new_groups . iter_mut (). find ( | x | * x . title == ** title ) {
grp . channels . push ( channel . clone ());
} else {
grp_id += 1 ;
new_groups . push ( PlaylistGroup {
id : grp_id ,
2025-08-26 13:53:11 +02:00
title : title . clone (),
2024-05-10 12:01:48 +02:00
channels : vec ! [ channel . clone ()],
2024-09-10 20:05:29 +02:00
xtream_cluster : * cluster ,
2024-05-10 12:01:48 +02:00
});
2023-04-26 14:30:43 +02:00
}
}
}
Some ( new_groups )
2023-01-12 17:15:50 +01:00
} else {
None
2022-04-05 12:23:30 +02:00
}
}
2025-03-11 14:56:10 +01:00
fn map_playlist_counter ( target : & ConfigTarget , playlist : & mut [ PlaylistGroup ]) {
2025-06-24 01:06:25 +02:00
if let Some ( guard ) = &* target . mapping . load () {
let mappings = guard . as_ref ();
for mapping in mappings {
2024-09-11 19:41:50 +02:00
if let Some ( counter_list ) = & mapping . t_counter {
for counter in counter_list {
2025-03-11 14:56:10 +01:00
for plg in & mut * playlist {
for channel in & mut plg . channels {
let provider = ValueProvider { pli : channel };
2025-05-22 13:34:45 +02:00
if counter . filter . filter ( & provider ) {
2025-04-28 08:30:48 +02:00
let cntval = counter . value . fetch_add ( 1 , core ::sync ::atomic ::Ordering ::SeqCst );
let padded_cntval = if counter . padding > 0 {
format! ( " {:0width$} " , cntval , width = counter . padding as usize )
} else {
2024-09-11 19:41:50 +02:00
cntval . to_string ()
2025-04-28 08:30:48 +02:00
};
let new_value = if counter . modifier == CounterModifier ::Assign {
padded_cntval
2024-09-11 19:41:50 +02:00
} else {
2025-03-11 14:56:10 +01:00
let value = channel . header . get_field ( & counter . field ). map_or_else ( String ::new , | field_value | field_value . to_string ());
2024-09-11 19:41:50 +02:00
if counter . modifier == CounterModifier ::Suffix {
2025-04-28 08:30:48 +02:00
format! ( " {value}{}{padded_cntval} " , counter . concat )
2024-10-25 00:39:22 +02:00
} else {
2025-04-28 08:30:48 +02:00
format! ( " {padded_cntval}{}{value} " , counter . concat )
2024-09-11 19:41:50 +02:00
}
};
2025-03-11 14:56:10 +01:00
channel . header . set_field ( & counter . field , new_value . as_str ());
2024-09-11 19:41:50 +02:00
}
}
}
}
}
}
}
}
2023-10-12 11:08:37 +02:00
// If no input is enabled but the user set the target as command line argument,
// we force the input to be enabled.
// If there are enabled input, then only these are used.
2025-05-23 18:17:47 +02:00
fn is_input_enabled ( input : & ConfigInput , user_targets : & ProcessTargets ) -> bool {
2025-02-14 17:51:33 +01:00
let input_enabled = input . enabled ;
let input_id = input . id ;
2025-05-23 18:17:47 +02:00
( ! user_targets . enabled && input_enabled ) || user_targets . has_input ( input_id )
2023-10-12 11:08:37 +02:00
}
2023-12-08 19:08:04 +01:00
fn is_target_enabled ( target : & ConfigTarget , user_targets : & ProcessTargets ) -> bool {
( ! user_targets . enabled && target . enabled ) || ( user_targets . enabled && user_targets . has_target ( target . id ))
}
2025-08-13 13:40:41 +02:00
async fn playlist_download_from_input ( client : & Arc < reqwest ::Client > , config : & Arc < Config > , input : & ConfigInput ) -> ( Vec < PlaylistGroup > , Vec < TuliproxError > ) {
let working_dir = & config . working_dir ;
match input . input_type {
InputType ::M3u => m3u ::get_m3u_playlist ( Arc ::clone ( client ), config , input , working_dir ). await ,
InputType ::Xtream => xtream ::get_xtream_playlist ( config , Arc ::clone ( client ), input , working_dir ). await ,
InputType ::M3uBatch | InputType ::XtreamBatch => ( vec! [], vec! [])
}
}
2025-08-04 14:18:46 +02:00
async fn process_source ( client : Arc < reqwest ::Client > , cfg : Arc < AppConfig > , source_idx : usize ,
user_targets : Arc < ProcessTargets > , event_manager : Option < Arc < EventManager >> )
-> ( Vec < InputStats > , Vec < TargetStats > , Vec < TuliproxError > ) {
2025-06-24 01:06:25 +02:00
let sources = cfg . sources . load ();
2023-10-13 13:59:22 +02:00
let mut errors = vec! [];
2025-01-27 15:58:57 +01:00
let mut input_stats = HashMap ::< String , InputStats > ::new ();
2025-01-05 13:29:28 +01:00
let mut target_stats = Vec ::< TargetStats > ::new ();
2025-07-22 11:32:46 +02:00
if let Some ( source ) = sources . get_source_at ( source_idx ) {
let mut source_playlists = Vec ::with_capacity ( 128 );
// Download the sources
let mut source_downloaded = false ;
for input in & source . inputs {
if is_input_enabled ( input , & user_targets ) {
let config = cfg . config . load ();
let working_dir = & config . working_dir ;
source_downloaded = true ;
let start_time = Instant ::now ();
2025-08-13 13:40:41 +02:00
let ( mut playlistgroups , mut error_list ) = playlist_download_from_input ( & client , & config , input ). await ;
2025-07-22 11:32:46 +02:00
let ( tvguide , mut tvguide_errors ) = if error_list . is_empty () {
epg ::get_xmltv ( Arc ::clone ( & client ), input , working_dir ). await
} else {
( None , vec! [])
};
errors . append ( & mut error_list );
errors . append ( & mut tvguide_errors );
let group_count = playlistgroups . len ();
let channel_count = playlistgroups . iter ()
. map ( | group | group . channels . len ())
. sum ();
let input_name = & input . name ;
if playlistgroups . is_empty () {
info! ( "Source is empty {input_name}" );
errors . push ( notify_err! ( format! ( "Source is empty {input_name} " )));
} else {
playlistgroups . iter_mut (). for_each ( PlaylistGroup ::on_load );
source_playlists . push (
FetchedPlaylist {
input ,
playlistgroups ,
epg : tvguide ,
}
);
}
let elapsed = start_time . elapsed (). as_secs ();
2025-08-26 13:53:11 +02:00
input_stats . insert ( input_name . clone (), create_input_stat ( group_count , channel_count , error_list . len (),
2025-07-22 11:32:46 +02:00
input . input_type , input_name , elapsed ));
2022-03-24 14:08:25 +01:00
}
}
2025-07-22 11:32:46 +02:00
if source_downloaded {
if source_playlists . is_empty () {
debug! ( "Source at index {source_idx} is empty" );
2025-07-22 17:37:56 +02:00
errors . push ( notify_err! ( format! ( "Source at index {source_idx} is empty: {} " , source . inputs . iter (). map ( | i | i . name . as_str ()). collect ::< Vec < _ >> (). join ( ", " ))));
2025-07-22 11:32:46 +02:00
} else {
debug_if_enabled! ( "Source has {} groups" , source_playlists . iter (). map ( | fpl | fpl . playlistgroups . len ()). sum ::< usize > ());
2025-08-04 14:18:46 +02:00
let event_manager_clone = event_manager . clone ();
2025-07-22 11:32:46 +02:00
for target in & source . targets {
2025-08-04 14:18:46 +02:00
let event_manager_clone = event_manager_clone . clone ();
2025-07-22 11:32:46 +02:00
if is_target_enabled ( target , & user_targets ) {
2025-08-04 14:18:46 +02:00
match process_playlist_for_target ( & cfg , Arc ::clone ( & client ), & mut source_playlists , target , & mut input_stats , & mut errors , event_manager_clone ). await {
2025-07-22 11:32:46 +02:00
Ok (()) => {
target_stats . push ( TargetStats ::success ( & target . name ));
}
Err ( mut err ) => {
target_stats . push ( TargetStats ::failure ( & target . name ));
errors . append ( & mut err );
}
2025-05-23 18:17:47 +02:00
}
2025-01-05 13:29:28 +01:00
}
2023-10-12 11:08:37 +02:00
}
}
2023-12-08 19:08:04 +01:00
}
2023-10-12 11:08:37 +02:00
}
2025-01-05 13:29:28 +01:00
( input_stats . into_values (). collect (), target_stats , errors )
2023-02-13 12:19:58 +01:00
}
2024-11-03 14:52:56 +01:00
fn create_input_stat ( group_count : usize , channel_count : usize , error_count : usize , input_type : InputType , input_name : & str , secs_took : u64 ) -> InputStats {
2024-10-25 00:39:22 +02:00
InputStats {
name : input_name . to_string (),
2024-10-30 11:14:55 +01:00
input_type ,
error_count ,
2024-10-25 00:39:22 +02:00
raw_stats : PlaylistStats {
group_count ,
channel_count ,
},
processed_stats : PlaylistStats {
group_count : 0 ,
channel_count : 0 ,
},
2024-11-04 18:46:56 +01:00
secs_took ,
2024-10-25 00:39:22 +02:00
}
}
2025-08-04 14:18:46 +02:00
async fn process_sources ( client : Arc < reqwest ::Client > , config : & Arc < AppConfig > , user_targets : Arc < ProcessTargets > , event_manager : Option < Arc < EventManager >> ) -> ( Vec < SourceStats > , Vec < TuliproxError > ) {
2023-02-25 16:23:26 +01:00
let mut handle_list = vec! [];
2025-06-24 01:06:25 +02:00
let thread_num = config . config . load (). threads ;
let sources = config . sources . load ();
let process_parallel = thread_num > 1 && sources . sources . len () > 1 ;
2024-03-26 16:07:31 +01:00
if process_parallel && log_enabled! ( Level ::Debug ) {
2025-04-10 00:06:48 +02:00
debug! ( "Using {thread_num} threads" );
2023-11-03 18:05:13 +01:00
}
2025-05-13 21:50:08 +02:00
let errors = Arc ::new ( Mutex ::< Vec < TuliproxError >> ::new ( vec! []));
2025-01-05 13:29:28 +01:00
let stats = Arc ::new ( Mutex ::< Vec < SourceStats >> ::new ( vec! []));
2025-06-24 01:06:25 +02:00
for ( index , _ ) in sources . sources . iter (). enumerate () {
2025-01-06 09:17:05 +01:00
// We're using the file lock this way on purpose
let source_lock_path = PathBuf ::from ( format! ( "source_ {index} " ));
2025-03-11 14:56:10 +01:00
let Ok ( update_lock ) = config . file_locks . try_write_lock ( & source_lock_path ). await else {
2025-01-06 09:17:05 +01:00
warn! ( "The update operation for the source at index {index} was skipped because an update is already in progress." );
continue ;
};
2023-10-13 13:59:22 +02:00
let shared_errors = errors . clone ();
2023-10-19 18:56:08 +02:00
let shared_stats = stats . clone ();
let cfg = config . clone ();
let usr_trgts = user_targets . clone ();
2025-08-04 14:18:46 +02:00
let event_manager = event_manager . clone ();
2023-02-13 12:19:58 +01:00
if process_parallel {
2025-01-09 15:58:38 +01:00
let http_client = Arc ::clone ( & client );
2023-02-13 19:08:28 +01:00
let handles = & mut handle_list ;
2023-11-03 18:05:13 +01:00
let process = move || {
2025-03-11 14:56:10 +01:00
// TODO better way ?
2025-07-22 11:32:46 +02:00
match tokio ::runtime ::Runtime ::new () {
Ok ( rt ) => {
rt . block_on ( async {
2025-08-04 14:18:46 +02:00
let ( input_stats , target_stats , mut res_errors ) =
process_source ( Arc ::clone ( & http_client ), cfg , index , usr_trgts , event_manager ). await ;
2025-07-22 11:32:46 +02:00
shared_errors . lock (). await . append ( & mut res_errors );
let process_stats = SourceStats ::new ( input_stats , target_stats );
shared_stats . lock (). await . push ( process_stats );
});
},
Err ( err ) => error! ( "Could not create runtime !!! {err}" ),
}
2023-11-03 18:05:13 +01:00
};
2023-05-03 17:46:39 +02:00
handles . push ( thread ::spawn ( process ));
2024-05-10 12:01:48 +02:00
if handles . len () >= thread_num as usize {
2023-10-13 13:59:22 +02:00
handles . drain ( .. ). for_each ( | handle | { let _ = handle . join (); });
2023-02-13 19:08:28 +01:00
}
2023-02-13 12:19:58 +01:00
} else {
2025-08-04 14:18:46 +02:00
let ( input_stats , target_stats , mut res_errors ) = process_source ( Arc ::clone ( & client ), cfg , index , usr_trgts , event_manager ). await ;
2025-03-11 14:56:10 +01:00
shared_errors . lock (). await . append ( & mut res_errors );
2025-01-05 13:29:28 +01:00
let process_stats = SourceStats ::new ( input_stats , target_stats );
2025-03-11 14:56:10 +01:00
shared_stats . lock (). await . push ( process_stats );
2023-02-13 12:19:58 +01:00
}
2025-01-06 09:17:05 +01:00
drop ( update_lock );
2023-02-13 12:19:58 +01:00
}
2023-02-13 19:08:28 +01:00
for handle in handle_list {
2023-02-25 16:23:26 +01:00
let _ = handle . join ();
2022-03-24 14:08:25 +01:00
}
2025-07-22 11:32:46 +02:00
if let ( Ok ( s ), Ok ( e )) = ( Arc ::try_unwrap ( stats ), Arc ::try_unwrap ( errors )) {
( s . into_inner (), e . into_inner ())
} else {
( vec! [], vec! [])
}
2022-03-24 14:08:25 +01:00
}
2023-09-29 16:03:20 +02:00
2024-05-04 20:30:20 +02:00
pub type ProcessingPipe = Vec < fn ( playlist : & mut [ PlaylistGroup ], target : & ConfigTarget ) -> Option < Vec < PlaylistGroup >>> ;
2023-09-29 16:03:20 +02:00
2023-10-13 13:59:22 +02:00
fn get_processing_pipe ( target : & ConfigTarget ) -> ProcessingPipe {
match & target . processing_order {
ProcessingOrder ::Frm => vec! [ filter_playlist , rename_playlist , map_playlist ],
ProcessingOrder ::Fmr => vec! [ filter_playlist , map_playlist , rename_playlist ],
ProcessingOrder ::Rfm => vec! [ rename_playlist , filter_playlist , map_playlist ],
ProcessingOrder ::Rmf => vec! [ rename_playlist , map_playlist , filter_playlist ],
ProcessingOrder ::Mfr => vec! [ map_playlist , filter_playlist , rename_playlist ],
ProcessingOrder ::Mrf => vec! [ map_playlist , rename_playlist , filter_playlist ]
}
}
2025-01-06 10:11:14 +01:00
fn duplicate_hash ( item : & PlaylistItem ) -> UUIDType {
2025-01-29 17:33:06 +01:00
item . get_uuid ()
2025-01-06 10:11:14 +01:00
}
2024-09-10 20:05:29 +02:00
2025-01-06 10:11:14 +01:00
fn execute_pipe < 'a > ( target : & ConfigTarget , pipe : & ProcessingPipe , fpl : & FetchedPlaylist < 'a > , duplicates : & mut HashSet < UUIDType > ) -> FetchedPlaylist < 'a > {
2024-09-10 20:05:29 +02:00
let mut new_fpl = FetchedPlaylist {
input : fpl . input ,
playlistgroups : fpl . playlistgroups . clone (), // we need to clone, because of multiple target definitions, we cant change the initial playlist.
epg : fpl . epg . clone (),
};
2025-01-06 10:11:14 +01:00
if target . options . as_ref (). is_some_and ( | opt | opt . remove_duplicates ) {
for group in & mut new_fpl . playlistgroups {
// `HashSet::insert` returns true for first insert, otherweise false
group . channels . retain ( | item | duplicates . insert ( duplicate_hash ( item )));
}
}
2024-09-10 20:05:29 +02:00
for f in pipe {
if let Some ( groups ) = f ( & mut new_fpl . playlistgroups , target ) {
new_fpl . playlistgroups = groups ;
}
}
new_fpl
}
// This method is needed, because of duplicate group names in different inputs.
// We merge the same group names considering cluster together.
2024-11-04 18:46:56 +01:00
fn flatten_groups ( playlistgroups : Vec < PlaylistGroup > ) -> Vec < PlaylistGroup > {
2024-09-11 12:13:05 +02:00
let mut sort_order : Vec < PlaylistGroup > = vec! [];
let mut idx : usize = 0 ;
2025-03-11 14:56:10 +01:00
let mut group_map : HashMap < ( String , XtreamCluster ), usize > = HashMap ::new ();
2024-11-04 18:46:56 +01:00
for group in playlistgroups {
2025-08-26 13:53:11 +02:00
let key = ( group . title . clone (), group . xtream_cluster );
2024-09-10 20:05:29 +02:00
match group_map . entry ( key ) {
2024-09-11 12:13:05 +02:00
std ::collections ::hash_map ::Entry ::Vacant ( v ) => {
v . insert ( idx );
idx += 1 ;
sort_order . push ( group );
}
2024-09-10 20:05:29 +02:00
std ::collections ::hash_map ::Entry ::Occupied ( o ) => {
2025-07-22 11:32:46 +02:00
if let Some ( pl_group ) = sort_order . get_mut ( * o . get ()) {
pl_group . channels . extend ( group . channels );
}
2024-09-10 20:05:29 +02:00
}
2025-04-10 00:06:48 +02:00
}
2024-11-04 18:46:56 +01:00
}
2024-09-11 12:13:05 +02:00
sort_order
2024-09-10 20:05:29 +02:00
}
2025-06-24 01:06:25 +02:00
async fn process_playlist_for_target ( app_config : & AppConfig ,
client : Arc < reqwest ::Client > ,
2025-01-09 15:58:38 +01:00
playlists : & mut [ FetchedPlaylist < '_ > ],
2025-01-06 10:11:14 +01:00
target : & ConfigTarget ,
2025-01-27 15:58:57 +01:00
stats : & mut HashMap < String , InputStats > ,
2025-08-04 14:18:46 +02:00
errors : & mut Vec < TuliproxError > ,
event_manager : Option < Arc < EventManager >> ) -> Result < (), Vec < TuliproxError >> {
2023-10-13 13:59:22 +02:00
let pipe = get_processing_pipe ( target );
2024-12-09 19:26:39 +01:00
debug_if_enabled! ( "Processing order is {}" , & target . processing_order );
2023-09-29 16:03:20 +02:00
2025-01-06 10:11:14 +01:00
let mut duplicates : HashSet < UUIDType > = HashSet ::new ();
2024-12-10 23:16:04 +01:00
let mut processed_fetched_playlists : Vec < FetchedPlaylist > = vec! [];
2025-04-23 20:21:36 +02:00
debug! ( "Executing processing pipes" );
2025-08-04 14:18:46 +02:00
let broadcast_step = {
let event_manager = event_manager . clone ();
move | context : & str , msg : & str | {
if let Some ( events ) = & event_manager {
events . send_event ( EventMessage ::PlaylistUpdateProgress ( context . to_owned (), msg . to_owned ()));
}
}
};
2025-04-23 20:21:36 +02:00
2025-08-04 14:18:46 +02:00
let mut step = StepMeasure ::new ( & target . name , broadcast_step );
2024-12-10 23:16:04 +01:00
for provider_fpl in playlists . iter_mut () {
2025-01-06 10:11:14 +01:00
let mut processed_fpl = execute_pipe ( target , & pipe , provider_fpl , & mut duplicates );
2025-06-24 01:06:25 +02:00
playlist_resolve_series ( app_config , Arc ::clone ( & client ), target , errors , & pipe , provider_fpl , & mut processed_fpl ). await ;
playlist_resolve_vod ( app_config , Arc ::clone ( & client ), target , errors , & mut processed_fpl ). await ;
2023-10-19 18:56:08 +02:00
// stats
2025-01-27 15:58:57 +01:00
let input_stats = stats . get_mut ( & processed_fpl . input . name );
2023-10-19 18:56:08 +02:00
if let Some ( stat ) = input_stats {
2024-12-10 23:16:04 +01:00
stat . processed_stats . group_count = processed_fpl . playlistgroups . len ();
stat . processed_stats . channel_count = processed_fpl . playlistgroups . iter ()
2023-10-19 18:56:08 +02:00
. map ( | group | group . channels . len ())
. sum ();
}
2024-12-10 23:16:04 +01:00
processed_fetched_playlists . push ( processed_fpl );
2023-12-08 19:08:04 +01:00
}
2025-08-04 14:18:46 +02:00
step . tick ( "filter rename map" );
2025-06-17 13:06:52 +02:00
let ( new_epg , mut new_playlist ) = process_epg ( & mut processed_fetched_playlists );
2025-08-04 14:18:46 +02:00
step . tick ( "epg" );
2025-04-15 18:38:18 +02:00
if new_playlist . is_empty () {
2025-08-04 14:18:46 +02:00
step . stop ( "" );
2025-04-15 18:38:18 +02:00
info! ( "Playlist is empty: {}" , & target . name );
Ok (())
} else {
2025-06-17 13:06:52 +02:00
// Process Trakt categories
trakt_playlist ( & client , target , errors , & mut new_playlist ). await ;
2025-08-04 14:18:46 +02:00
step . tick ( "trakt categories" );
2025-06-17 13:06:52 +02:00
2025-04-15 18:38:18 +02:00
let mut flat_new_playlist = flatten_groups ( new_playlist );
2025-08-04 14:18:46 +02:00
step . tick ( "playlist merge" );
2025-06-17 13:06:52 +02:00
2025-04-15 18:38:18 +02:00
sort_playlist ( target , & mut flat_new_playlist );
2025-08-04 14:18:46 +02:00
step . tick ( "playlist sort" );
2025-05-04 10:59:08 +02:00
assign_channel_no_playlist ( & mut flat_new_playlist );
2025-08-04 14:18:46 +02:00
step . tick ( "assigning channel numbers" );
2025-04-15 18:38:18 +02:00
map_playlist_counter ( target , & mut flat_new_playlist );
2025-08-04 14:18:46 +02:00
step . tick ( "assigning channel counter" );
2025-06-16 11:49:17 +02:00
2025-06-24 01:06:25 +02:00
let config = app_config . config . load ();
process_watch ( & config , & client , target , & flat_new_playlist );
2025-08-04 14:18:46 +02:00
step . tick ( "group watches" );
2025-06-24 01:06:25 +02:00
let result = persist_playlist ( app_config , & mut flat_new_playlist , flatten_tvguide ( & new_epg ). as_ref (), target ). await ;
2025-08-04 14:18:46 +02:00
step . stop ( "Persisting playlists" );
2025-04-23 20:21:36 +02:00
result
2025-04-15 18:38:18 +02:00
}
}
2025-06-17 13:06:52 +02:00
async fn trakt_playlist ( client : & Arc < Client > , target : & ConfigTarget , errors : & mut Vec < TuliproxError > , playlist : & mut Vec < PlaylistGroup > ) {
match process_trakt_categories_for_target ( Arc ::clone ( client ), playlist , target ). await {
2025-06-16 17:53:03 +02:00
Ok ( trakt_categories ) => {
if ! trakt_categories . is_empty () {
info! ( "Adding {} Trakt categories to playlist" , trakt_categories . len ());
2025-06-17 13:06:52 +02:00
playlist . extend ( trakt_categories );
2025-06-16 17:53:03 +02:00
}
}
Err ( trakt_errors ) => {
warn! ( "Trakt processing failed with {} errors" , trakt_errors . len ());
errors . extend ( trakt_errors );
}
}
}
2025-04-19 13:50:41 +02:00
fn process_epg ( processed_fetched_playlists : & mut Vec < FetchedPlaylist > ) -> ( Vec < Epg > , Vec < PlaylistGroup > ) {
2023-10-13 15:56:00 +02:00
let mut new_playlist = vec! [];
2023-10-27 19:03:29 +02:00
let mut new_epg = vec! [];
2024-05-05 19:15:54 +02:00
2024-11-02 16:44:38 +01:00
// each fetched playlist can have its own epgl url.
// we need to process each input epg.
2025-04-15 18:38:18 +02:00
for fp in processed_fetched_playlists {
2025-04-23 20:21:36 +02:00
process_playlist_epg ( fp , & mut new_epg );
2025-04-06 16:50:43 +02:00
new_playlist . append ( & mut fp . playlistgroups );
2024-11-04 18:46:56 +01:00
}
2025-04-15 18:38:18 +02:00
( new_epg , new_playlist )
2024-09-10 20:05:29 +02:00
}
2023-10-13 18:13:16 +02:00
2025-06-24 01:06:25 +02:00
fn process_watch ( cfg : & Config , client : & Arc < reqwest ::Client > , target : & ConfigTarget , new_playlist : & Vec < PlaylistGroup > ) {
if let Some ( watches ) = & target . watch {
2024-09-10 20:05:29 +02:00
if default_as_default (). eq_ignore_ascii_case ( & target . name ) {
error! ( "cant watch a target with no unique name" );
} else {
for pl in new_playlist {
2025-06-24 01:06:25 +02:00
if watches . iter (). any ( | r | r . is_match ( & pl . title )) {
2025-04-30 13:32:55 +02:00
process_group_watch ( client , cfg , & target . name , pl );
2024-05-10 12:01:48 +02:00
}
2023-10-13 18:13:16 +02:00
}
}
2023-10-06 19:24:06 +02:00
}
2023-09-29 16:03:20 +02:00
}
2023-10-19 18:56:08 +02:00
2025-08-03 10:24:23 +02:00
pub async fn exec_processing ( client : Arc < reqwest ::Client > , app_config : Arc < AppConfig > , targets : Arc < ProcessTargets > , event_manager : Option < Arc < EventManager >> ) {
2024-12-25 14:12:30 +01:00
let start_time = Instant ::now ();
2025-08-04 14:18:46 +02:00
let event_manager_clone = event_manager . clone ();
let ( stats , errors ) = process_sources ( Arc ::clone ( & client ), & app_config , targets . clone (), event_manager_clone ). await ;
2024-12-13 17:51:47 +01:00
// log errors
2024-12-28 20:01:28 +01:00
for err in & errors {
error! ( "{}" , err . message );
}
2025-08-03 10:24:23 +02:00
let config = app_config . config . load ();
2025-06-24 01:06:25 +02:00
let messaging = config . messaging . as_ref ();
2024-12-25 14:12:30 +01:00
if let Ok ( stats_msg ) = serde_json ::to_string ( & serde_json ::Value ::Object ( serde_json ::map ::Map ::from_iter ([( "stats" . to_string (), serde_json ::to_value ( stats ). unwrap ())]))) {
// print stats
2025-04-10 00:06:48 +02:00
info! ( "{stats_msg}" );
2024-12-25 14:12:30 +01:00
// send stats
2025-06-24 01:06:25 +02:00
send_message ( & client , & MsgKind ::Stats , messaging , stats_msg . as_str ());
2024-12-13 17:51:47 +01:00
}
2023-10-20 08:23:25 +02:00
// send errors
2023-10-19 18:56:08 +02:00
if let Some ( message ) = get_errors_notify_message! ( errors , 255 ) {
2025-08-03 10:24:23 +02:00
if let Some ( events ) = event_manager {
events . send_event ( EventMessage ::PlaylistUpdate ( PlaylistUpdateState ::Failure ));
}
2024-12-25 14:12:30 +01:00
if let Ok ( error_msg ) = serde_json ::to_string ( & serde_json ::Value ::Object ( serde_json ::map ::Map ::from_iter ([( "errors" . to_string (), serde_json ::Value ::String ( message ))]))) {
2025-06-24 01:06:25 +02:00
send_message ( & client , & MsgKind ::Error , messaging , error_msg . as_str ());
2024-12-25 14:12:30 +01:00
}
2025-08-03 10:24:23 +02:00
} else if let Some ( events ) = event_manager {
events . send_event ( EventMessage ::PlaylistUpdate ( PlaylistUpdateState ::Success ));
2023-10-19 18:56:08 +02:00
}
2024-12-25 14:12:30 +01:00
let elapsed = start_time . elapsed (). as_secs ();
2025-05-12 21:38:02 +02:00
info! ( "🌷 Update process finished! Took {elapsed} secs." );
2025-02-01 21:50:13 -06:00
}
2025-04-11 20:25:09 +02:00
2025-08-13 13:40:41 +02:00
// #[cfg(test)]
// mod tests {
2025-04-30 17:16:53 +02:00
// #[test]
// fn test_jaro_winkeler() {
// let data = [("yessport5", "heyessport5gold"), ("yessport5", "heyesport5gold")];
//
// data.iter().for_each(|(first, second)|
// println!("jaro_winkler {} = {} => {}", first, second, strsim::jaro_winkler(first, second)));
// // println!("jaro {}", strsim::jaro(data.0, data.1));
// // println!("levenhstein {}", strsim::levenshtein(data.0, data.1));
// // println!("damerau_levenshtein {:?}", strsim::damerau_levenshtein(data.0, data.1));
// // println!("osa distance {:?}", strsim::osa_distance(data.0, data.1));
// // println!("sorensen dice {:?}", strsim::sorensen_dice(data.0, data.1));
// }
2025-08-13 13:40:41 +02:00
// }