mirror of
https://github.com/euzu/tuliprox.git
synced 2026-09-30 13:02:10 +02:00
WebUI Source Explorer custom source download
This commit is contained in:
@@ -30,7 +30,8 @@ use log::{debug, error, info, log_enabled, trace, warn, Level};
|
||||
use reqwest::Client;
|
||||
use shared::error::{get_errors_notify_message, notify_err, TuliproxError, TuliproxErrorKind};
|
||||
use shared::foundation::filter::{get_field_value, set_field_value, ValueAccessor, ValueProvider};
|
||||
use shared::model::{CounterModifier, FieldGetAccessor, FieldSetAccessor, InputType, ItemField, MsgKind, PlaylistEntry, PlaylistGroup, PlaylistItem, PlaylistUpdateState, ProcessingOrder, UUIDType, XtreamCluster};
|
||||
use shared::model::{CounterModifier, FieldGetAccessor, FieldSetAccessor, InputType, ItemField, MsgKind, PlaylistEntry,
|
||||
PlaylistGroup, PlaylistItem, PlaylistUpdateState, ProcessingOrder, UUIDType, XtreamCluster};
|
||||
use shared::utils::default_as_default;
|
||||
use std::time::Instant;
|
||||
use crate::api::model::{EventManager, EventMessage};
|
||||
@@ -233,7 +234,9 @@ fn is_target_enabled(target: &ConfigTarget, user_targets: &ProcessTargets) -> bo
|
||||
(!user_targets.enabled && target.enabled) || (user_targets.enabled && user_targets.has_target(target.id))
|
||||
}
|
||||
|
||||
async fn process_source(client: Arc<reqwest::Client>, cfg: Arc<AppConfig>, source_idx: usize, user_targets: Arc<ProcessTargets>) -> (Vec<InputStats>, Vec<TargetStats>, Vec<TuliproxError>) {
|
||||
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>) {
|
||||
let sources = cfg.sources.load();
|
||||
let mut errors = vec![];
|
||||
let mut input_stats = HashMap::<String, InputStats>::new();
|
||||
@@ -290,9 +293,11 @@ async fn process_source(client: Arc<reqwest::Client>, cfg: Arc<AppConfig>, sourc
|
||||
errors.push(notify_err!(format!("Source at index {source_idx} is empty: {}", source.inputs.iter().map(|i| i.name.as_str()).collect::<Vec<_>>().join(", "))));
|
||||
} else {
|
||||
debug_if_enabled!("Source has {} groups", source_playlists.iter().map(|fpl| fpl.playlistgroups.len()).sum::<usize>());
|
||||
let event_manager_clone = event_manager.clone();
|
||||
for target in &source.targets {
|
||||
let event_manager_clone = event_manager_clone.clone();
|
||||
if is_target_enabled(target, &user_targets) {
|
||||
match process_playlist_for_target(&cfg, Arc::clone(&client), &mut source_playlists, target, &mut input_stats, &mut errors).await {
|
||||
match process_playlist_for_target(&cfg, Arc::clone(&client), &mut source_playlists, target, &mut input_stats, &mut errors, event_manager_clone).await {
|
||||
Ok(()) => {
|
||||
target_stats.push(TargetStats::success(&target.name));
|
||||
}
|
||||
@@ -326,7 +331,7 @@ fn create_input_stat(group_count: usize, channel_count: usize, error_count: usiz
|
||||
}
|
||||
}
|
||||
|
||||
async fn process_sources(client: Arc<reqwest::Client>, config: &Arc<AppConfig>, user_targets: Arc<ProcessTargets>) -> (Vec<SourceStats>, Vec<TuliproxError>) {
|
||||
async fn process_sources(client: Arc<reqwest::Client>, config: &Arc<AppConfig>, user_targets: Arc<ProcessTargets>, event_manager: Option<Arc<EventManager>>) -> (Vec<SourceStats>, Vec<TuliproxError>) {
|
||||
let mut handle_list = vec![];
|
||||
let thread_num = config.config.load().threads;
|
||||
let sources = config.sources.load();
|
||||
@@ -348,6 +353,7 @@ async fn process_sources(client: Arc<reqwest::Client>, config: &Arc<AppConfig>,
|
||||
let shared_stats = stats.clone();
|
||||
let cfg = config.clone();
|
||||
let usr_trgts = user_targets.clone();
|
||||
let event_manager = event_manager.clone();
|
||||
if process_parallel {
|
||||
let http_client = Arc::clone(&client);
|
||||
let handles = &mut handle_list;
|
||||
@@ -356,7 +362,8 @@ async fn process_sources(client: Arc<reqwest::Client>, config: &Arc<AppConfig>,
|
||||
match tokio::runtime::Runtime::new() {
|
||||
Ok(rt) => {
|
||||
rt.block_on(async {
|
||||
let (input_stats, target_stats, mut res_errors) = process_source(Arc::clone(&http_client), cfg, index, usr_trgts).await;
|
||||
let (input_stats, target_stats, mut res_errors) =
|
||||
process_source(Arc::clone(&http_client), cfg, index, usr_trgts, event_manager).await;
|
||||
shared_errors.lock().await.append(&mut res_errors);
|
||||
let process_stats = SourceStats::new(input_stats, target_stats);
|
||||
shared_stats.lock().await.push(process_stats);
|
||||
@@ -370,7 +377,7 @@ async fn process_sources(client: Arc<reqwest::Client>, config: &Arc<AppConfig>,
|
||||
handles.drain(..).for_each(|handle| { let _ = handle.join(); });
|
||||
}
|
||||
} else {
|
||||
let (input_stats, target_stats, mut res_errors) = process_source(Arc::clone(&client), cfg, index, usr_trgts).await;
|
||||
let (input_stats, target_stats, mut res_errors) = process_source(Arc::clone(&client), cfg, index, usr_trgts, event_manager).await;
|
||||
shared_errors.lock().await.append(&mut res_errors);
|
||||
let process_stats = SourceStats::new(input_stats, target_stats);
|
||||
shared_stats.lock().await.push(process_stats);
|
||||
@@ -454,7 +461,8 @@ async fn process_playlist_for_target(app_config: &AppConfig,
|
||||
playlists: &mut [FetchedPlaylist<'_>],
|
||||
target: &ConfigTarget,
|
||||
stats: &mut HashMap<String, InputStats>,
|
||||
errors: &mut Vec<TuliproxError>) -> Result<(), Vec<TuliproxError>> {
|
||||
errors: &mut Vec<TuliproxError>,
|
||||
event_manager: Option<Arc<EventManager>>) -> Result<(), Vec<TuliproxError>> {
|
||||
let pipe = get_processing_pipe(target);
|
||||
debug_if_enabled!("Processing order is {}", &target.processing_order);
|
||||
|
||||
@@ -462,8 +470,16 @@ async fn process_playlist_for_target(app_config: &AppConfig,
|
||||
let mut processed_fetched_playlists: Vec<FetchedPlaylist> = vec![];
|
||||
|
||||
debug!("Executing processing pipes");
|
||||
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()));
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
let mut step = StepMeasure::new("Pipes processed");
|
||||
let mut step = StepMeasure::new(&target.name, broadcast_step);
|
||||
for provider_fpl in playlists.iter_mut() {
|
||||
let mut processed_fpl = execute_pipe(target, &pipe, provider_fpl, &mut duplicates);
|
||||
playlist_resolve_series(app_config, Arc::clone(&client), target, errors, &pipe, provider_fpl, &mut processed_fpl).await;
|
||||
@@ -478,35 +494,35 @@ async fn process_playlist_for_target(app_config: &AppConfig,
|
||||
}
|
||||
processed_fetched_playlists.push(processed_fpl);
|
||||
}
|
||||
|
||||
step.tick("Processed epg");
|
||||
step.tick("filter rename map");
|
||||
let (new_epg, mut new_playlist) = process_epg(&mut processed_fetched_playlists);
|
||||
step.tick("epg");
|
||||
|
||||
if new_playlist.is_empty() {
|
||||
step.stop("");
|
||||
info!("Playlist is empty: {}", &target.name);
|
||||
Ok(())
|
||||
} else {
|
||||
|
||||
// Process Trakt categories
|
||||
step.tick("Processing Trakt categories");
|
||||
trakt_playlist(&client, target, errors, &mut new_playlist).await;
|
||||
step.tick("trakt categories");
|
||||
|
||||
step.tick("Merged playlists");
|
||||
let mut flat_new_playlist = flatten_groups(new_playlist);
|
||||
step.tick("playlist merge");
|
||||
|
||||
step.tick("Sorted playlists");
|
||||
sort_playlist(target, &mut flat_new_playlist);
|
||||
step.tick("Assigned channel number");
|
||||
step.tick("playlist sort");
|
||||
assign_channel_no_playlist(&mut flat_new_playlist);
|
||||
step.tick("Assigned channel counter");
|
||||
step.tick("assigning channel numbers");
|
||||
map_playlist_counter(target, &mut flat_new_playlist);
|
||||
step.tick("assigning channel counter");
|
||||
|
||||
step.tick("Processed group watches");
|
||||
let config = app_config.config.load();
|
||||
process_watch(&config, &client, target, &flat_new_playlist);
|
||||
step.tick("Persisting playlists");
|
||||
step.tick("group watches");
|
||||
let result = persist_playlist(app_config, &mut flat_new_playlist, flatten_tvguide(&new_epg).as_ref(), target).await;
|
||||
step.stop();
|
||||
step.stop("Persisting playlists");
|
||||
result
|
||||
}
|
||||
}
|
||||
@@ -555,7 +571,8 @@ fn process_watch(cfg: &Config, client: &Arc<reqwest::Client>, target: &ConfigTar
|
||||
|
||||
pub async fn exec_processing(client: Arc<reqwest::Client>, app_config: Arc<AppConfig>, targets: Arc<ProcessTargets>, event_manager: Option<Arc<EventManager>>) {
|
||||
let start_time = Instant::now();
|
||||
let (stats, errors) = process_sources(Arc::clone(&client), &app_config, targets.clone()).await;
|
||||
let event_manager_clone = event_manager.clone();
|
||||
let (stats, errors) = process_sources(Arc::clone(&client), &app_config, targets.clone(), event_manager_clone).await;
|
||||
// log errors
|
||||
for err in &errors {
|
||||
error!("{}", err.message);
|
||||
|
||||
Reference in New Issue
Block a user