vod resolve refactored

Arc<Client> is now Client
This commit is contained in:
euzu
2025-12-11 09:54:28 +01:00
parent f0ec224f52
commit 55191d6c84
25 changed files with 163 additions and 168 deletions
+21 -22
View File
@@ -29,7 +29,6 @@ use crate::utils::StepMeasure;
use deunicode::deunicode;
use futures::StreamExt;
use log::{debug, error, info, log_enabled, trace, warn, Level};
use reqwest::Client;
use shared::error::{get_errors_notify_message, notify_err, TuliproxError};
use shared::foundation::filter::{get_field_value, set_field_value, Filter, ValueAccessor, ValueProvider};
use shared::model::{CounterModifier, FieldGetAccessor, FieldSetAccessor, InputType, ItemField, MsgKind, PlaylistEntry, PlaylistGroup, PlaylistItem, PlaylistUpdateState, ProcessingOrder, UUIDType, XtreamCluster};
@@ -305,7 +304,7 @@ 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 playlist_download_from_input(client: &Arc<reqwest::Client>, config: &Arc<Config>, input: &Arc<ConfigInput>) -> (Vec<PlaylistGroup>, Vec<TuliproxError>) {
async fn playlist_download_from_input(client: &reqwest::Client, config: &Arc<Config>, input: &Arc<ConfigInput>) -> (Vec<PlaylistGroup>, Vec<TuliproxError>) {
let working_dir = &config.working_dir;
match input.input_type {
InputType::M3u => m3u::get_m3u_playlist(client, config, input, working_dir).await,
@@ -314,7 +313,7 @@ async fn playlist_download_from_input(client: &Arc<reqwest::Client>, config: &Ar
}
}
async fn process_source(client: Arc<reqwest::Client>, cfg: Arc<AppConfig>, source_idx: usize,
async fn process_source(client: &reqwest::Client, cfg: Arc<AppConfig>, source_idx: usize,
user_targets: Arc<ProcessTargets>, event_manager: Option<Arc<EventManager>>,
playlist_state: Option<&Arc<PlaylistStorageState>>,
) -> (Vec<InputStats>, Vec<TargetStats>, Vec<TuliproxError>) {
@@ -332,9 +331,9 @@ async fn process_source(client: Arc<reqwest::Client>, cfg: Arc<AppConfig>, sourc
let working_dir = &config.working_dir;
source_downloaded = true;
let start_time = Instant::now();
let (mut playlistgroups, mut error_list) = playlist_download_from_input(&client, &config, input).await;
let (mut playlistgroups, mut error_list) = playlist_download_from_input(client, &config, input).await;
let (tvguide, mut tvguide_errors) = if error_list.is_empty() {
epg::get_xmltv(Arc::clone(&client), input, working_dir).await
epg::get_xmltv(client, input, working_dir).await
} else {
(None, vec![])
};
@@ -373,7 +372,7 @@ async fn process_source(client: Arc<reqwest::Client>, cfg: Arc<AppConfig>, sourc
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, event_manager_clone, playlist_state).await {
match process_playlist_for_target(&cfg, client, &mut source_playlists, target, &mut input_stats, &mut errors, event_manager_clone, playlist_state).await {
Ok(()) => {
target_stats.push(TargetStats::success(&target.name));
}
@@ -407,7 +406,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>,
async fn process_sources(client: &reqwest::Client, config: &Arc<AppConfig>, user_targets: Arc<ProcessTargets>,
event_manager: Option<Arc<EventManager>>, playlist_state: Option<&Arc<PlaylistStorageState>>,
) -> (Vec<SourceStats>, Vec<TuliproxError>) {
let mut async_tasks = JoinSet::new();
@@ -432,13 +431,13 @@ async fn process_sources(client: Arc<reqwest::Client>, config: &Arc<AppConfig>,
let usr_trgts = user_targets.clone();
let event_manager = event_manager.clone();
if process_parallel {
let http_client = Arc::clone(&client);
let http_client = client.clone();
let playlist_state = playlist_state.cloned();
async_tasks.spawn(async move {
// Hold the per-source lock for the full duration of this update.
let current_update_lock = update_lock;
let (input_stats, target_stats, mut res_errors) =
process_source(Arc::clone(&http_client), cfg, index, usr_trgts, event_manager, playlist_state.as_ref()).await;
process_source(&http_client, cfg, index, usr_trgts, event_manager, playlist_state.as_ref()).await;
shared_errors.lock().await.append(&mut res_errors);
if let Some(process_stats) = SourceStats::try_new(input_stats, target_stats) {
shared_stats.lock().await.push(process_stats);
@@ -447,7 +446,7 @@ async fn process_sources(client: Arc<reqwest::Client>, config: &Arc<AppConfig>,
});
} else {
let (input_stats, target_stats, mut res_errors) =
process_source(Arc::clone(&client), cfg, index, usr_trgts, event_manager, playlist_state).await;
process_source(client, cfg, index, usr_trgts, event_manager, playlist_state).await;
shared_errors.lock().await.append(&mut res_errors);
if let Some(process_stats) = SourceStats::try_new(input_stats, target_stats) {
shared_stats.lock().await.push(process_stats);
@@ -531,7 +530,7 @@ fn flatten_groups(playlistgroups: Vec<PlaylistGroup>) -> Vec<PlaylistGroup> {
#[allow(clippy::too_many_arguments)]
async fn process_playlist_for_target(app_config: &AppConfig,
client: Arc<reqwest::Client>,
client: &reqwest::Client,
playlists: &mut [FetchedPlaylist<'_>],
target: &ConfigTarget,
stats: &mut HashMap<String, InputStats>,
@@ -558,8 +557,8 @@ async fn process_playlist_for_target(app_config: &AppConfig,
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;
playlist_resolve_vod(app_config, Arc::clone(&client), target, errors, &mut processed_fpl).await;
playlist_resolve_series(app_config, client, target, errors, &pipe, provider_fpl, &mut processed_fpl).await;
playlist_resolve_vod(app_config, client, target, errors, &mut processed_fpl).await;
// stats
let input_stats = stats.get_mut(&processed_fpl.input.name);
if let Some(stat) = input_stats {
@@ -580,7 +579,7 @@ async fn process_playlist_for_target(app_config: &AppConfig,
Ok(())
} else {
// Process Trakt categories
if trakt_playlist(&client, target, errors, &mut new_playlist).await {
if trakt_playlist(client, target, errors, &mut new_playlist).await {
step.tick("trakt categories");
}
@@ -596,7 +595,7 @@ async fn process_playlist_for_target(app_config: &AppConfig,
step.tick("assigning channel counter");
let config = app_config.config.load();
if process_watch(&config, &client, target, &flat_new_playlist).await {
if process_watch(&config, client, target, &flat_new_playlist).await {
step.tick("group watches");
}
let result = persist_playlist(app_config, &mut flat_new_playlist, flatten_tvguide(&new_epg).as_ref(), target, playlist_state).await;
@@ -605,8 +604,8 @@ async fn process_playlist_for_target(app_config: &AppConfig,
}
}
async fn trakt_playlist(client: &Arc<Client>, target: &ConfigTarget, errors: &mut Vec<TuliproxError>, playlist: &mut Vec<PlaylistGroup>) -> bool {
match process_trakt_categories_for_target(Arc::clone(client), playlist, target).await {
async fn trakt_playlist(client: &reqwest::Client, target: &ConfigTarget, errors: &mut Vec<TuliproxError>, playlist: &mut Vec<PlaylistGroup>) -> bool {
match process_trakt_categories_for_target(client, playlist, target).await {
Ok(Some(trakt_categories)) => {
if !trakt_categories.is_empty() {
info!("Adding {} Trakt categories to playlist", trakt_categories.len());
@@ -637,7 +636,7 @@ async fn process_epg(processed_fetched_playlists: &mut Vec<FetchedPlaylist<'_>>)
(new_epg, new_playlist)
}
async fn process_watch(cfg: &Config, client: &Arc<reqwest::Client>, target: &ConfigTarget, new_playlist: &[PlaylistGroup]) -> bool {
async fn process_watch(cfg: &Config, client: &reqwest::Client, target: &ConfigTarget, new_playlist: &[PlaylistGroup]) -> bool {
if let Some(watches) = &target.watch {
if default_as_default().eq_ignore_ascii_case(&target.name) {
error!("can't watch a target with no unique name");
@@ -657,10 +656,10 @@ async fn process_watch(cfg: &Config, client: &Arc<reqwest::Client>, target: &Con
}
}
pub async fn exec_processing(client: Arc<reqwest::Client>, app_config: Arc<AppConfig>, targets: Arc<ProcessTargets>, event_manager: Option<Arc<EventManager>>, playlist_state: Option<Arc<PlaylistStorageState>>) {
pub async fn exec_processing(client: &reqwest::Client, app_config: Arc<AppConfig>, targets: Arc<ProcessTargets>, event_manager: Option<Arc<EventManager>>, playlist_state: Option<Arc<PlaylistStorageState>>) {
let start_time = Instant::now();
let event_manager_clone = event_manager.clone();
let (stats, errors) = process_sources(Arc::clone(&client), &app_config, targets.clone(), event_manager_clone, playlist_state.as_ref()).await;
let (stats, errors) = process_sources(client, &app_config, targets.clone(), event_manager_clone, playlist_state.as_ref()).await;
// log errors
for err in &errors {
error!("{}", err.message);
@@ -677,7 +676,7 @@ pub async fn exec_processing(client: Arc<reqwest::Client>, app_config: Arc<AppCo
// print stats
info!("{stats_msg}");
// send stats
send_message_json(&client, MsgKind::Stats, messaging, stats_msg.as_str()).await;
send_message_json(client, MsgKind::Stats, messaging, stats_msg.as_str()).await;
}
Err(err) => error!("Failed to serialize playlist stats {err}"),
}
@@ -692,7 +691,7 @@ pub async fn exec_processing(client: Arc<reqwest::Client>, app_config: Arc<AppCo
events.send_event(EventMessage::PlaylistUpdate(PlaylistUpdateState::Failure));
}
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))]))) {
send_message_json(&client, MsgKind::Error, messaging, error_msg.as_str()).await;
send_message_json(client, MsgKind::Error, messaging, error_msg.as_str()).await;
}
} else if let Some(events) = event_manager {
events.send_event(EventMessage::PlaylistUpdate(PlaylistUpdateState::Success));
+3 -4
View File
@@ -8,7 +8,6 @@ use shared::utils::{get_u32_from_serde_value};
use shared::utils::{CONSTANTS};
use crate::utils::{trace_if_enabled, with};
use log::{debug, info, trace, warn};
use std::sync::Arc;
use strsim::normalized_levenshtein;
fn extract_quality(value: &str) -> Option<&str> {
@@ -243,8 +242,8 @@ pub struct TraktCategoriesProcessor {
}
impl TraktCategoriesProcessor {
pub fn new(http_client: Arc<reqwest::Client>, trakt_config: &TraktConfig) -> Self {
let client = TraktClient::new(http_client, trakt_config.api.clone());
pub fn new(http_client: &reqwest::Client, trakt_config: &TraktConfig) -> Self {
let client = TraktClient::new(http_client.clone(), trakt_config.api.clone());
Self { client }
}
@@ -292,7 +291,7 @@ impl TraktCategoriesProcessor {
}
}
pub async fn process_trakt_categories_for_target(
http_client: Arc<reqwest::Client>,
http_client: &reqwest::Client,
playlist: &[PlaylistGroup],
target: &ConfigTarget,
) -> Result<Option<Vec<PlaylistGroup>>, Vec<TuliproxError>> {
+1 -2
View File
@@ -9,7 +9,6 @@ use serde::{Deserialize, Serialize};
use std::collections::HashMap;
use std::fs::File;
use std::path::PathBuf;
use std::sync::Arc;
use crate::repository::bplustree::BPlusTree;
use crate::repository::storage_const;
use crate::repository::xtream_repository::xtream_get_record_file_path;
@@ -17,7 +16,7 @@ use crate::utils;
use crate::utils::xtream;
use serde_json::{from_str, to_string, Value};
pub(in crate::processing) async fn playlist_resolve_download_playlist_item(client: Arc<reqwest::Client>, pli: &PlaylistItem, input: &ConfigInput, errors: &mut Vec<TuliproxError>, resolve_delay: u16, cluster: XtreamCluster) -> Option<String> {
pub(in crate::processing) async fn playlist_resolve_download_playlist_item(client: &reqwest::Client, pli: &PlaylistItem, input: &ConfigInput, errors: &mut Vec<TuliproxError>, resolve_delay: u16, cluster: XtreamCluster) -> Option<String> {
let mut result = None;
let provider_id = pli.get_provider_id()?;
if let Some(info_url) = xtream::get_xtream_player_api_info_url(input, cluster, provider_id) {
@@ -14,7 +14,6 @@ use crate::processing::processor::{handle_error, handle_error_and_return, create
use std::collections::{HashMap, HashSet};
use std::fs::File;
use std::io::{BufWriter, Write};
use std::sync::Arc;
use std::time::Instant;
use log::{error, info, log_enabled, warn, Level};
use crate::model::{XtreamSeriesEpisode, XtreamSeriesInfoEpisode};
@@ -50,7 +49,7 @@ fn should_update_series_info(pli: &mut PlaylistItem, processed_provider_ids: &Ha
should_update_info(pli, processed_provider_ids, crate::model::XC_TAG_SERIES_INFO_LAST_MODIFIED)
}
async fn playlist_resolve_series_info(cfg: &AppConfig, client: Arc<reqwest::Client>, errors: &mut Vec<TuliproxError>,
async fn playlist_resolve_series_info(cfg: &AppConfig, client: &reqwest::Client, errors: &mut Vec<TuliproxError>,
fpl: &mut FetchedPlaylist<'_>, resolve_delay: u16) -> bool {
let mut processed_info_ids = read_processed_series_info_ids(cfg, errors, fpl).await;
let mut fetched_in_run: HashSet<u32> = HashSet::new();
@@ -83,7 +82,7 @@ async fn playlist_resolve_series_info(cfg: &AppConfig, client: Arc<reqwest::Clie
for pli in series_info_iter {
let (should_update, provider_id, ts) = should_update_series_info(pli, &processed_info_ids);
if should_update && provider_id != 0 && fetched_in_run.insert(provider_id) {
if let Some(content) = playlist_resolve_download_playlist_item(Arc::clone(&client), pli, fpl.input, errors, resolve_delay, XtreamCluster::Series).await {
if let Some(content) = playlist_resolve_download_playlist_item(client, pli, fpl.input, errors, resolve_delay, XtreamCluster::Series).await {
let normalized_content = normalize_json_content(content);
handle_error_and_return!(write_series_info_to_wal_file(provider_id, ts, &normalized_content, &mut content_writer, &mut record_writer),
|err| errors.push(notify_err!(format!("Failed to resolve series, could not write to wal file {err}"))));
@@ -217,7 +216,7 @@ async fn process_series_info(
pub async fn playlist_resolve_series(cfg: &AppConfig,
client: Arc<reqwest::Client>,
client: &reqwest::Client,
target: &ConfigTarget,
errors: &mut Vec<TuliproxError>,
pipe: &ProcessingPipe,
+55 -54
View File
@@ -8,11 +8,9 @@ use crate::repository::xtream_repository::{write_vod_info_to_wal_file, xtream_up
use shared::error::{notify_err};
use crate::processing::processor::{handle_error, handle_error_and_return, create_resolve_options_function_for_xtream_target};
use shared::utils::{get_u32_from_serde_value, get_u64_from_serde_value, get_string_from_serde_value};
use crate::repository::xtream_repository::xtream_get_input_info;
use serde_json::{from_str, Map, Value};
use std::collections::{HashMap, HashSet};
use std::io::{Write};
use std::sync::Arc;
use std::time::Instant;
use log::{info, log_enabled, Level};
use crate::utils;
@@ -65,7 +63,9 @@ fn should_update_vod_info(pli: &mut PlaylistItem, processed_provider_ids: &HashM
should_update_info(pli, processed_provider_ids, crate::model::XC_TAG_VOD_INFO_ADDED)
}
pub async fn playlist_resolve_vod(app_config: &AppConfig, client: Arc<reqwest::Client>, target: &ConfigTarget, errors: &mut Vec<TuliproxError>, fpl: &mut FetchedPlaylist<'_>) {
const FLUSH_INTERVAL: usize = 50;
pub async fn playlist_resolve_vod(app_config: &AppConfig, client: &reqwest::Client, target: &ConfigTarget, errors: &mut Vec<TuliproxError>, fpl: &mut FetchedPlaylist<'_>) {
let (resolve_movies, resolve_delay) = get_resolve_vod_options(target, fpl);
if !resolve_movies { return; }
@@ -76,7 +76,7 @@ pub async fn playlist_resolve_vod(app_config: &AppConfig, client: Arc<reqwest::C
let Some((wal_content_file, wal_record_file, wal_content_path, wal_record_path)) = create_resolve_info_wal_files(&config, fpl.input, XtreamCluster::Video)
else { return; };
let mut processed_info_ids = read_processed_vod_info_ids(app_config, errors, fpl).await;
let mut processed_info_ids: HashMap<u32, u64> = read_processed_vod_info_ids(app_config, errors, fpl).await;
let mut fetched_in_run: HashSet<u32> = HashSet::new();
let mut content_writer = utils::file_writer(&wal_content_file);
let mut record_writer = utils::file_writer(&wal_record_file);
@@ -87,78 +87,79 @@ pub async fn playlist_resolve_vod(app_config: &AppConfig, client: Arc<reqwest::C
.flat_map(|plg| &plg.channels)
.filter(|&pli| pli.header.xtream_cluster == XtreamCluster::Video).count();
let vod_info_iter = fpl.playlistgroups.iter_mut()
.flat_map(|plg| plg.channels.iter_mut())
.filter(|pli| pli.header.xtream_cluster == XtreamCluster::Video);
info!("Found {vod_info_count} vod info to resolve");
let start_time = Instant::now();
let mut last_log_time = Instant::now();
let mut processed_vod_info_count = 0;
let mut last_processed_vod_info_count = 0;
let mut write_counter = 0usize;
for pli in vod_info_iter {
let (should_update, provider_id, _ts) = should_update_vod_info(pli, &processed_info_ids);
if should_update && provider_id != 0 && fetched_in_run.insert(provider_id) {
if let Some(content) = playlist_resolve_download_playlist_item(Arc::clone(&client), pli, fpl.input, errors, resolve_delay, XtreamCluster::Video).await {
let normalized_content = normalize_json_content(content);
if let Some((provider_id, info_record)) = extract_info_record_from_vod_info(&normalized_content) {
let ts = info_record.ts;
handle_error_and_return!(write_vod_info_to_wal_file(provider_id, &normalized_content, &info_record, &mut content_writer, &mut record_writer),
|err| errors.push(notify_err!(format!("Failed to resolve vod, could not write to wal file {err}"))));
processed_info_ids.insert(provider_id, ts);
content_updated = true;
for plg in &mut fpl.playlistgroups {
for pli in &mut plg.channels {
if pli.header.xtream_cluster != XtreamCluster::Video {
continue;
}
let (should_update, provider_id, _ts) = should_update_vod_info(pli, &processed_info_ids);
if should_update && provider_id != 0 && fetched_in_run.insert(provider_id) {
if let Some(content) = playlist_resolve_download_playlist_item(client, pli, fpl.input, errors, resolve_delay, XtreamCluster::Video).await {
let normalized_content = normalize_json_content(content);
let normalized_str: &str = &normalized_content;
if let Some((provider_id, info_record)) = extract_info_record_from_vod_info(normalized_str) {
let ts = info_record.ts;
handle_error_and_return!(write_vod_info_to_wal_file(provider_id, normalized_str, &info_record, &mut content_writer, &mut record_writer),
|err| errors.push(notify_err!(format!("Failed to resolve vod, could not write to wal file {err}"))));
processed_info_ids.insert(provider_id, ts);
content_updated = true;
write_counter += 1;
// periodic flush to bound BufWriter memory
if write_counter.is_multiple_of(FLUSH_INTERVAL) {
if let Err(err) = content_writer.flush() {
errors.push(notify_err!(format!("Failed periodic flush of wal content writer {err}")));
}
if let Err(err) = record_writer.flush() {
errors.push(notify_err!(format!("Failed periodic flush of wal record writer {err}")));
}
}
// Update in-memory playlist items with the newly fetched vod info.
// This makes the data available for subsequent processing steps like STRM export.
pli.header.additional_properties = from_str::<Map<String, Value>>(normalized_str).ok().and_then(|info_doc| {
info_doc.get("info").cloned().map(|info_content| {
let mut wrapped_info = Map::new();
wrapped_info.insert("info".to_string(), info_content);
Value::Object(wrapped_info)
})
});
}
}
}
if log_enabled!(Level::Info) {
processed_vod_info_count += 1;
if last_log_time.elapsed().as_secs() >= 30 {
info!("resolved {processed_vod_info_count}/{vod_info_count} vod info");
last_log_time = Instant::now();
}
}
}
if log_enabled!(Level::Info) {
processed_vod_info_count += 1;
let elapsed = start_time.elapsed().as_secs();
if elapsed > 0 && ((processed_vod_info_count - last_processed_vod_info_count) > 50) && elapsed.is_multiple_of(30) {
info!("resolved {processed_vod_info_count}/{vod_info_count} vod info");
last_processed_vod_info_count = processed_vod_info_count;
}
}
}
if last_processed_vod_info_count != processed_vod_info_count {
info!("resolved {processed_vod_info_count}/{vod_info_count} vod info");
}
info!("resolved {processed_vod_info_count}/{vod_info_count} vod info");
if content_updated {
// TODO better approach for transactional updates is multiplexed WAL file.
// final flush & sync with proper error handling
handle_error!(content_writer.flush(),
|err| errors.push(notify_err!(format!("Failed to resolve vod, could not write to wal file {err}"))));
handle_error!(record_writer.flush(),
|err| errors.push(notify_err!(format!("Failed to resolve vod tmdb, could not write to wal file {err}"))));
handle_error!(content_writer.get_ref().sync_all(), |err| errors.push(notify_err!(format!("Failed to sync vod info to wal file {err}"))));
handle_error!(record_writer.get_ref().sync_all(), |err| errors.push(notify_err!(format!("Failed to sync vod info record to wal file {err}"))));
// drop writers and files to release handles
drop(content_writer);
drop(record_writer);
drop(wal_content_file);
drop(wal_record_file);
handle_error!(xtream_update_input_info_file(app_config, fpl.input, &wal_content_path, XtreamCluster::Video).await,
|err| errors.push(err));
handle_error!(xtream_update_input_vod_record_from_wal_file(app_config, fpl.input, &wal_record_path).await,
|err| errors.push(err));
}
// Update in-memory playlist items with the newly fetched vod info.
// This makes the data available for subsequent processing steps like STRM export.
let vod_info_iter = fpl.playlistgroups.iter_mut()
.flat_map(|plg| &mut plg.channels)
.filter(|pli| pli.header.xtream_cluster == XtreamCluster::Video);
for pli in vod_info_iter {
if let Some(provider_id) = pli.header.get_provider_id() {
if let Some(content) = xtream_get_input_info(app_config, fpl.input, provider_id, XtreamCluster::Video).await {
// Add the "info" section to the playlist item additional properties.
pli.header.additional_properties = from_str::<Map<String, Value>>(&content).ok().and_then(|info_doc| {
info_doc.get("info").cloned().map(|info_content| {
let mut wrapped_info = Map::new();
wrapped_info.insert("info".to_string(), info_content);
Value::Object(wrapped_info)
})
});
}
}
}
}