Merge branch 'feature/source_editor' into feature/local_media

This commit is contained in:
euzu
2025-12-12 17:05:12 +01:00
59 changed files with 1168 additions and 671 deletions
+24 -25
View File
@@ -24,12 +24,11 @@ use crate::processing::processor::trakt::process_trakt_categories_for_target;
use crate::processing::processor::xtream_series::playlist_resolve_series;
use crate::processing::processor::xtream_vod::playlist_resolve_vod;
use crate::repository::playlist_repository::persist_playlist;
use crate::utils::debug_if_enabled;
use crate::utils::{debug_if_enabled, trace_if_enabled};
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};
@@ -135,7 +134,7 @@ fn exec_rename(pli: &mut PlaylistItem, rename: Option<&Vec<ConfigRename>>) {
let value = get_field_value(result, r.field);
let cap = r.pattern.replace_all(value.as_str(), &r.new_name);
if log_enabled!(log::Level::Debug) && *value != cap {
debug_if_enabled!("Renamed {}={} to {}", &r.field, value, cap);
trace_if_enabled!("Renamed {}={value} to {cap}", &r.field);
}
let value = cap.into_owned();
set_field_value(result, r.field, value);
@@ -154,7 +153,7 @@ fn rename_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Opt
for r in renames {
if matches!(r.field, ItemField::Group) {
let cap = r.pattern.replace_all(&grp.title, &r.new_name);
debug_if_enabled!("Renamed group {} to {} for {}", &grp.title, cap, target.name);
trace_if_enabled!("Renamed group {} to {cap} for {}", &grp.title, target.name);
grp.title = cap.into_owned();
}
}
@@ -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,
@@ -315,7 +314,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>) {
@@ -333,9 +332,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![])
};
@@ -374,7 +373,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));
}
@@ -408,7 +407,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();
@@ -433,13 +432,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);
@@ -448,7 +447,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);
@@ -532,7 +531,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>,
@@ -559,8 +558,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 {
@@ -581,7 +580,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");
}
@@ -597,7 +596,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;
@@ -606,8 +605,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());
@@ -638,7 +637,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");
@@ -658,10 +657,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);
@@ -678,7 +677,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}"),
}
@@ -693,7 +692,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,13 +14,12 @@ 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};
use crate::utils;
use crate::processing::processor::xtream::normalize_json_content;
use crate::utils::bincode_serialize;
use crate::utils::{bincode_serialize, IO_BUFFER_SIZE};
create_resolve_options_function_for_xtream_target!(series);
@@ -32,27 +31,28 @@ fn write_series_episode_record_to_wal_file(
writer: &mut BufWriter<&File>,
provider_id: u32,
episode: &XtreamSeriesInfoEpisode,
) -> std::io::Result<()> {
) -> std::io::Result<usize> {
let series_episode = XtreamSeriesEpisode::from(episode);
if let Ok(content_bytes) = bincode_serialize(&series_episode) {
writer.write_all(&provider_id.to_le_bytes())?;
if let Ok(len) = u32::try_from(content_bytes.len()) {
let content_len = content_bytes.len();
if let Ok(len) = u32::try_from(content_len) {
writer.write_all(&len.to_le_bytes())?;
writer.write_all(&content_bytes)?;
} else {
error!("Cant write to WAL file, content length exceeds u32");
return Ok(content_len + 4usize)
}
error!("Cant write to WAL file, content length exceeds u32");
}
Ok(())
Ok(0)
}
fn should_update_series_info(pli: &mut PlaylistItem, processed_provider_ids: &HashMap<u32, u64>) -> (bool, u32, u64) {
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 processed_info_ids: HashMap<u32, u64> = read_processed_series_info_ids(cfg, errors, fpl).await;
let mut fetched_in_run: HashSet<u32> = HashSet::new();
// we cant write to the indexed-document directly because of the write lock and time-consuming operation.
// All readers would be waiting for the lock and the app would be unresponsive.
@@ -64,57 +64,70 @@ async fn playlist_resolve_series_info(cfg: &AppConfig, client: Arc<reqwest::Clie
let mut record_writer = utils::file_writer(&wal_record_file);
let mut content_updated = false;
// TODO merge both filters to one
let series_info_count = fpl.playlistgroups.iter()
.filter(|&plg| plg.xtream_cluster == XtreamCluster::Series)
.flat_map(|plg| &plg.channels)
.filter(|&pli| pli.header.item_type == PlaylistItemType::SeriesInfo).count();
let series_info_iter = fpl.playlistgroups.iter_mut()
.filter(|plg| plg.xtream_cluster == XtreamCluster::Series)
.flat_map(|plg| &mut plg.channels)
.filter(|pli| pli.header.item_type == PlaylistItemType::SeriesInfo);
info!("Found {series_info_count} series info to resolve");
let start_time = Instant::now();
let mut last_log_time = Instant::now();
let mut processed_series_info_count = 0;
let mut last_processed_series_info_count = 0;
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 {
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}"))));
processed_info_ids.insert(provider_id, ts);
content_updated = true;
}
let mut write_counter = 0usize;
for plg in &mut fpl.playlistgroups {
if plg.xtream_cluster != XtreamCluster::Series {
continue;
}
if log_enabled!(Level::Info) {
processed_series_info_count += 1;
let elapsed = start_time.elapsed().as_secs();
if elapsed > 0 && ((processed_series_info_count - last_processed_series_info_count) > 50) && elapsed.is_multiple_of(30) {
info!("resolved {processed_series_info_count}/{series_info_count} series info");
last_processed_series_info_count = processed_series_info_count;
for pli in &mut plg.channels {
if pli.header.item_type != PlaylistItemType::SeriesInfo {
continue;
}
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(client, pli, fpl.input, errors, resolve_delay, XtreamCluster::Series).await {
let normalized_content = normalize_json_content(content);
let normalized_str = normalized_content.as_str();
handle_error_and_return!(write_series_info_to_wal_file(provider_id, ts, normalized_str, &mut content_writer, &mut record_writer),
|err| errors.push(notify_err!(format!("Failed to resolve series, could not write to wal file {err}"))));
processed_info_ids.insert(provider_id, ts);
content_updated = true;
write_counter += normalized_str.len();
// periodic flush to bound BufWriter memory
if write_counter >= IO_BUFFER_SIZE {
write_counter = 0;
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}")));
}
}
}
}
if log_enabled!(Level::Info) {
processed_series_info_count += 1;
if last_log_time.elapsed().as_secs() >= 30 {
info!("resolved {processed_series_info_count}/{series_info_count} series info");
last_log_time = Instant::now();
}
}
}
}
if last_processed_series_info_count != processed_series_info_count {
info!("resolved {processed_series_info_count}/{series_info_count} series info");
}
info!("resolved {processed_series_info_count}/{series_info_count} series info");
// content_wal contains the provider_id and series_info with episode listing
// record_wal contains provider_id and timestamp
if content_updated {
handle_error!(content_writer.flush(),
|err| errors.push(notify_err!(format!("Failed to resolve vod, could not write to wal file {err}"))));
|err| errors.push(notify_err!(format!("Failed to resolve series, 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}"))));
|err| errors.push(notify_err!(format!("Failed to resolve series 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 series info to wal file {err}"))));
handle_error!(record_writer.get_ref().sync_all(), |err| errors.push(notify_err!(format!("Failed to sync series info record to wal file {err}"))));
drop(content_writer);
drop(wal_content_file);
drop(record_writer);
drop(wal_content_file);
drop(wal_record_file);
handle_error!(xtream_update_input_info_file(cfg, fpl.input, &wal_content_path, XtreamCluster::Series).await,
|err| errors.push(err));
@@ -143,6 +156,8 @@ async fn process_series_info(
return result;
};
let mut write_counter = 0usize;
let _file_lock = app_config.file_locks.read_lock(&info_path).await;
// Contains the Series Info with episode listing
@@ -182,8 +197,19 @@ async fn process_series_info(
Ok(Some(mut series)) => {
for (episode, pli_episode) in &mut series {
let Some(provider_id) = &pli_episode.header.get_provider_id() else { continue; };
handle_error!(write_series_episode_record_to_wal_file(&mut wal_writer, *provider_id, episode),
|err| errors.push(info_err!(format!("Failed to write to series episode wal file: {err}"))));
match write_series_episode_record_to_wal_file(&mut wal_writer, *provider_id, episode) {
Ok(written_bytes) => {
write_counter += written_bytes;
// periodic flush to bound BufWriter memory
if write_counter >= IO_BUFFER_SIZE {
write_counter = 0;
if let Err(err) = wal_writer.flush() {
errors.push(notify_err!(format!("Failed periodic flush of wal content writer {err}")));
}
}
}
Err(err) => {errors.push(info_err!(format!("Failed to write to series episode wal file: {err}"))) }
}
}
group_series.extend(series.into_iter().map(|(_, pli)| pli));
}
@@ -206,8 +232,9 @@ async fn process_series_info(
}
}
handle_error!(wal_writer.flush(),
|err| errors.push(notify_err!(format!("Failed to resolve series episodes, could not write to wal file {err}"))));
handle_error!(wal_writer.flush(), |err| errors.push(notify_err!(format!("Failed to resolve series episodes, could not write to wal file {err}"))));
handle_error!(wal_writer.get_ref().sync_all(), |err| errors.push(notify_err!(format!("Failed to sync series info to wal file {err}"))));
drop(wal_writer);
drop(wal_file);
handle_error!(xtream_update_input_series_episodes_record_from_wal_file(app_config, input, &wal_path).await,
@@ -217,7 +244,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,
+56 -54
View File
@@ -8,15 +8,14 @@ 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;
use crate::processing::processor::xtream::normalize_json_content;
use crate::utils::IO_BUFFER_SIZE;
create_resolve_options_function_for_xtream_target!(vod);
@@ -65,10 +64,12 @@ 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<'_>) {
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; }
// TODO read existing WAL File and import it to avoid duplicate requests
// we cant write to the indexed-document directly because of the write lock and time-consuming operation.
// All readers would be waiting for the lock and the app would be unresponsive.
// We collect the content into a wal file and write it once we collected everything.
@@ -76,7 +77,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 +88,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 += normalized_str.len();
// periodic flush to bound BufWriter memory
if write_counter >= IO_BUFFER_SIZE {
write_counter = 0;
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)
})
});
}
}
}
}