Code restructuring

This commit is contained in:
euzu
2025-02-05 20:01:49 +01:00
parent dc607d5520
commit bbf6d6e014
82 changed files with 616 additions and 507 deletions
+59
View File
@@ -0,0 +1,59 @@
use crate::model::config::{ConfigInput, InputAffix, AFFIX_FIELDS, valid_property};
use crate::model::playlist::{FetchedPlaylist, FieldGetAccessor, FieldSetAccessor, PlaylistItem};
use crate::utils::{debug_if_enabled};
type AffixProcessor<'a> = Box<dyn Fn(&mut PlaylistItem) + 'a>;
fn create_affix_processor(affix: &InputAffix, is_prefix: bool) -> AffixProcessor {
Box::new(move |channel: &mut PlaylistItem| {
let header = &mut channel.header.borrow_mut();
let value = header.get_field(affix.field.as_str()).map_or_else(|| String::from(&affix.value), |field_value| if is_prefix {
format!("{}{}", &affix.value, field_value.as_str())
} else {
format!("{}{}", field_value.as_str(), &affix.value)
});
debug_if_enabled!("Applying input {}: {}={}", if is_prefix {"prefix"} else {"suffix"}, &affix.field, &value);
header.set_field(&affix.field, value.as_str());
})
}
fn validate_and_create_affix_processor(affix: Option<&InputAffix>, is_prefix: bool) -> Option<AffixProcessor> {
if let Some(affix_def) = affix {
if (valid_property!(&affix_def.field.as_str(), AFFIX_FIELDS) && !affix_def.value.is_empty()) {
return Some(create_affix_processor(affix_def, is_prefix));
}
};
None
}
fn get_affix_processor(input: &ConfigInput) -> Option<AffixProcessor> {
if input.suffix.is_some() || input.prefix.is_some() {
let processors: Vec<AffixProcessor> = vec![
validate_and_create_affix_processor(input.prefix.as_ref(), true),
validate_and_create_affix_processor(input.suffix.as_ref(), false)
].into_iter().flatten().collect();
if !processors.is_empty() {
let apply_affix: AffixProcessor = Box::new(move |channel: &mut PlaylistItem| {
for x in &processors {
x(channel);
}
});
return Some(apply_affix);
}
}
None
}
pub fn apply_affixes(fetched_playlists: &mut [FetchedPlaylist]) {
for fetched_playlist in fetched_playlists.iter_mut() {
let FetchedPlaylist { input, playlistgroups: playlist, epg: _ } = fetched_playlist;
if let Some(affix_processor) = get_affix_processor(input) {
for group in playlist.iter_mut() {
group.channels.iter_mut().for_each(|channel| {
affix_processor(channel);
});
}
}
}
}
+44
View File
@@ -0,0 +1,44 @@
pub mod playlist;
mod xtream;
mod affix;
mod xtream_vod;
mod xtream_series;
#[macro_export]
macro_rules! handle_error {
($stmt:expr, $map_err:expr) => {
if let Err(err) = $stmt {
$map_err(err);
}
};
}
use handle_error;
#[macro_export]
macro_rules! handle_error_and_return {
($stmt:expr, $map_err:expr) => {
if let Err(err) = $stmt {
$map_err(err);
return Default::default();
}
};
}
use handle_error_and_return;
#[macro_export]
macro_rules! create_resolve_options_function_for_xtream_target {
($cluster:ident) => {
paste::paste! {
fn [<get_resolve_ $cluster _options>](target: &ConfigTarget, fpl: &FetchedPlaylist) -> (bool, u16) {
let (resolve, resolve_delay) =
target.options.as_ref().map_or((false, 0), |opt| {
(opt.[<xtream_resolve_ $cluster>] && fpl.input.input_type == InputType::Xtream,
opt.[<xtream_resolve_ $cluster _delay>])
});
(resolve, resolve_delay)
}
}
};
}
use create_resolve_options_function_for_xtream_target;
+599
View File
@@ -0,0 +1,599 @@
extern crate unidecode;
use crate::Config;
use crate::model::config::ConfigRename;
use crate::utils::network::epg;
use crate::utils::network::m3u;
use crate::utils::network::xtream;
use async_std::sync::Mutex;
use core::cmp::Ordering;
use std::cell::RefCell;
use std::collections::{HashMap, HashSet};
use std::path::PathBuf;
use std::rc::Rc;
use std::sync::Arc;
use std::thread;
use actix_rt::System;
use log::{debug, error, info, log_enabled, trace, warn, Level};
use std::time::Instant;
use unidecode::unidecode;
use crate::foundation::filter::{get_field_value, set_field_value, MockValueProcessor, ValueProvider};
use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind, get_errors_notify_message, notify_err};
use crate::messaging::{send_message, MsgKind};
use crate::model::config::{ConfigSortChannel, ConfigSortGroup, ConfigTarget, InputType,
ItemField, ProcessTargets, ProcessingOrder, SortOrder::{Asc, Desc}};
use crate::model::mapping::{CounterModifier, Mapping, MappingValueProcessor};
use crate::model::playlist::{FetchedPlaylist, FieldGetAccessor, FieldSetAccessor, PlaylistEntry, PlaylistGroup, PlaylistItem, UUIDType, XtreamCluster};
use crate::model::stats::{InputStats, PlaylistStats, SourceStats, TargetStats};
use crate::processing::processor::affix::apply_affixes;
use crate::processing::playlist_watch::process_group_watch;
use crate::processing::parser::xmltv::flatten_tvguide;
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::default_utils::default_as_default;
use crate::utils::{debug_if_enabled};
fn is_valid(pli: &PlaylistItem, target: &ConfigTarget) -> bool {
let provider = ValueProvider { pli: RefCell::new(pli) };
target.filter(&provider)
}
#[allow(clippy::unnecessary_wraps)]
fn filter_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option<Vec<PlaylistGroup>> {
debug!("Filtering {} groups", playlist.len());
let mut new_playlist = Vec::with_capacity(128);
for pg in playlist.iter_mut() {
let channels = pg.channels.iter()
.filter(|&pli| is_valid(pli, target)).cloned().collect::<Vec<PlaylistItem>>();
trace!("Filtered group {} has now {}/{} items", pg.title, channels.len(), pg.channels.len());
if !channels.is_empty() {
new_playlist.push(PlaylistGroup {
id: pg.id,
title: pg.title.clone(),
channels,
xtream_cluster: pg.xtream_cluster,
});
}
}
Some(new_playlist)
}
fn playlistgroup_comparator(a: &PlaylistGroup, b: &PlaylistGroup, group_sort: &ConfigSortGroup, match_as_ascii: bool) -> Ordering {
let value_a = if match_as_ascii { Rc::new(unidecode(&a.title)) } else { Rc::clone(&a.title) };
let value_b = if match_as_ascii { Rc::new(unidecode(&b.title)) } else { Rc::clone(&b.title) };
let ordering = value_a.partial_cmp(&value_b).unwrap();
match group_sort.order {
Asc => ordering,
Desc => ordering.reverse()
}
}
fn playlistitem_comparator(a: &PlaylistItem, b: &PlaylistItem, channel_sort: &ConfigSortChannel, match_as_ascii: bool) -> Ordering {
let raw_value_a = get_field_value(a, &channel_sort.field);
let raw_value_b = get_field_value(b, &channel_sort.field);
let value_a = if match_as_ascii { Rc::new(unidecode(&raw_value_a)) } else { raw_value_a };
let value_b = if match_as_ascii { Rc::new(unidecode(&raw_value_b)) } else { raw_value_b };
channel_sort.sequence.as_ref().map_or_else(|| {
let ordering = value_a.partial_cmp(&value_b).unwrap();
match channel_sort.order {
Asc => ordering,
Desc => ordering.reverse()
}
}, |custom_order| {
// Check indices in the custom order vector
let index_a = custom_order.iter().position(|s| s == value_a.as_ref());
let index_b = custom_order.iter().position(|s| s == value_b.as_ref());
match (index_a, index_b) {
(Some(idx_a), Some(idx_b)) => {
// Both items found in custom order, compare indices
idx_a.cmp(&idx_b)
}
(Some(_), None) => {
// Only 'a' found in custom order, it comes first
Ordering::Less
}
(None, Some(_)) => {
// Only 'b' found in custom order, it comes first
Ordering::Greater
}
(None, None) => {
// Neither found, fall back to default ordering
let ordering = value_a.partial_cmp(&value_b).unwrap();
match channel_sort.order {
Asc => ordering,
Desc => ordering.reverse(),
}
}
}
})
}
fn sort_playlist(target: &ConfigTarget, new_playlist: &mut [PlaylistGroup]) {
if let Some(sort) = &target.sort {
let match_as_ascii = sort.match_as_ascii;
if let Some(group_sort) = &sort.groups {
new_playlist.sort_by(|a, b| playlistgroup_comparator(a, b, group_sort, match_as_ascii));
}
if let Some(channel_sorts) = &sort.channels {
for channel_sort in channel_sorts {
let regexp = channel_sort.re.as_ref().unwrap();
for group in new_playlist.iter_mut() {
let group_title = if match_as_ascii { Rc::new(unidecode(&group.title)) } else { Rc::clone(&group.title) };
if regexp.is_match(group_title.as_str()) {
group.channels.sort_by(|chan1, chan2| playlistitem_comparator(chan1, chan2, channel_sort, match_as_ascii));
}
}
}
}
}
}
fn channel_no_playlist(new_playlist: &[PlaylistGroup]) {
let mut chno = 1;
for group in new_playlist {
for chan in &group.channels {
chan.header.borrow_mut().chno = Rc::new(chno.to_string());
chno += 1;
}
}
}
fn exec_rename(pli: &PlaylistItem, rename: Option<&Vec<ConfigRename>>) {
if let Some(renames) = rename {
if !renames.is_empty() {
let result = pli;
for r in renames {
let value = get_field_value(result, &r.field);
let cap = r.re.as_ref().unwrap().replace_all(value.as_str(), &r.new_name);
if log::log_enabled!(log::Level::Debug) && *value != cap {
debug_if_enabled!("Renamed {}={} to {}", &r.field, value, cap);
}
let value = cap.into_owned();
set_field_value(result, &r.field, Rc::new(value));
}
}
}
}
fn rename_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option<Vec<PlaylistGroup>> {
match &target.rename {
Some(renames) => {
if !renames.is_empty() {
let mut new_playlist: Vec<PlaylistGroup> = Vec::with_capacity(playlist.len());
for g in playlist {
let mut grp = g.clone();
for r in renames {
if matches!(r.field, ItemField::Group) {
let cap = r.re.as_ref().unwrap().replace_all(&grp.title, &r.new_name);
debug_if_enabled!("Renamed group {} to {} for {}", &grp.title, cap, target.name);
grp.title = Rc::new(cap.into_owned());
}
}
grp.channels.iter_mut().for_each(|pli| exec_rename(pli, target.rename.as_ref()));
new_playlist.push(grp);
}
return Some(new_playlist);
}
None
}
_ => None
}
}
macro_rules! apply_pattern {
($pattern:expr, $provider:expr, $processor:expr) => {{
if let Some(ptrn) = $pattern {
ptrn.filter($provider, $processor);
};
}};
}
fn map_channel(channel: PlaylistItem, mapping: &Mapping) -> PlaylistItem {
if !mapping.mapper.is_empty() {
let header = channel.header.borrow();
let channel_name = if mapping.match_as_ascii { Rc::new(unidecode(&header.name)) } else { header.name.clone() };
if mapping.match_as_ascii && log_enabled!(Level::Trace) { trace!("Decoded {} for matching to {}", &header.name, &channel_name); };
drop(header);
let ref_chan = RefCell::new(&channel);
let provider = ValueProvider { pli: ref_chan.clone() };
let mut mock_processor = MockValueProcessor {};
for m in &mapping.mapper {
let mut processor = MappingValueProcessor { pli: ref_chan.clone(), mapper: m };
match &m.t_filter {
Some(filter) => {
if filter.filter(&provider, &mut mock_processor) {
apply_pattern!(&m.t_pattern, &provider, &mut processor);
}
}
_ => {
apply_pattern!(&m.t_pattern, &provider, &mut processor);
}
};
}
}
channel
}
fn map_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option<Vec<PlaylistGroup>> {
if target.t_mapping.is_some() {
let new_playlist: Vec<PlaylistGroup> = playlist.iter().map(|playlist_group| {
let mut grp = playlist_group.clone();
let mappings = target.t_mapping.as_ref().unwrap();
mappings.iter().filter(|&mapping| !mapping.mapper.is_empty()).for_each(|mapping|
grp.channels = grp.channels.drain(..).map(|chan| map_channel(chan, mapping)).collect());
grp
}).collect();
// if the group names are changed, restructure channels to the right groups
// we use
let mut new_groups: Vec<PlaylistGroup> = Vec::with_capacity(128);
let mut grp_id: u32 = 0;
for playlist_group in new_playlist {
for channel in &playlist_group.channels {
let cluster = &channel.header.borrow().xtream_cluster;
let title = &channel.header.borrow().group;
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,
title: Rc::clone(title),
channels: vec![channel.clone()],
xtream_cluster: *cluster,
});
}
}
}
Some(new_groups)
} else {
None
}
}
fn map_playlist_counter(target: &ConfigTarget, playlist: &[PlaylistGroup]) {
if target.t_mapping.is_some() {
let mut mock_processor = MockValueProcessor {};
let mappings = target.t_mapping.as_ref().unwrap();
for mapping in mappings {
if let Some(counter_list) = &mapping.t_counter {
for counter in counter_list {
for plg in playlist {
for channel in &plg.channels {
let provider = ValueProvider { pli: RefCell::new(channel) };
if counter.filter.filter(&provider, &mut mock_processor) {
let cntval = counter.value.load(core::sync::atomic::Ordering::Relaxed);
let new_value = if counter.modifier == CounterModifier::Assign {
cntval.to_string()
} else {
let value = channel.header.borrow_mut().get_field(&counter.field).map_or_else(String::new, |field_value| field_value.to_string());
if counter.modifier == CounterModifier::Suffix {
format!("{value}{}{cntval}", counter.concat)
} else {
format!("{cntval}{}{value}", counter.concat)
}
};
channel.header.borrow_mut().set_field(&counter.field, new_value.as_str());
counter.value.fetch_add(1, core::sync::atomic::Ordering::Relaxed);
}
}
}
}
}
}
}
}
// 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.
fn is_input_enabled(enabled_inputs: usize, input_enabled: bool, input_id: u16, user_targets: &ProcessTargets) -> bool {
if enabled_inputs == 0 {
return user_targets.enabled && user_targets.has_input(input_id);
}
input_enabled
}
fn is_target_enabled(target: &ConfigTarget, user_targets: &ProcessTargets) -> bool {
(!user_targets.enabled && target.enabled) || (user_targets.enabled && user_targets.has_target(target.id))
}
async fn process_source(client: Arc<reqwest::Client>, cfg: Arc<Config>, source_idx: usize, user_targets: Arc<ProcessTargets>) -> (Vec<InputStats>, Vec<TargetStats>, Vec<M3uFilterError>) {
let source = cfg.sources.get(source_idx).unwrap();
let mut errors = vec![];
let mut input_stats = HashMap::<String, InputStats>::new();
let mut target_stats = Vec::<TargetStats>::new();
let mut source_playlists = Vec::with_capacity(128);
let enabled_inputs = source.inputs.iter().filter(|item| item.enabled).count();
// Downlod the sources
for input in &source.inputs {
if is_input_enabled(enabled_inputs, input.enabled, input.id, &user_targets) {
let start_time = Instant::now();
let (mut playlistgroups, mut error_list) = match input.input_type {
InputType::M3u => m3u::get_m3u_playlist(Arc::clone(&client), &cfg, input, &cfg.working_dir).await,
InputType::Xtream => xtream::get_xtream_playlist(Arc::clone(&client), input, &cfg.working_dir).await,
};
let (tvguide, mut tvguide_errors) = if error_list.is_empty() {
epg::get_xmltv(Arc::clone(&client), &cfg, input, &cfg.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();
input_stats.insert(input_name.to_string(), create_input_stat(group_count, channel_count, error_list.len(),
input.input_type.clone(), input_name, elapsed));
}
}
if source_playlists.is_empty() {
debug!("Source at index {source_idx} is empty");
errors.push(notify_err!(format!("Source at {source_idx} is empty")));
} else {
debug_if_enabled!("Source has {} groups", source_playlists.iter().map(|fpl| fpl.playlistgroups.len()).sum::<usize>());
for target in &source.targets {
if is_target_enabled(target, &user_targets) {
match process_playlist_for_target(Arc::clone(&client), &mut source_playlists, target, &cfg, &mut input_stats, &mut errors).await {
Ok(()) => {
target_stats.push(TargetStats::success(&target.name));
}
Err(mut err) => {
target_stats.push(TargetStats::failure(&target.name));
errors.append(&mut err);
}
}
}
}
}
(input_stats.into_values().collect(), target_stats, errors)
}
fn create_input_stat(group_count: usize, channel_count: usize, error_count: usize, input_type: InputType, input_name: &str, secs_took: u64) -> InputStats {
InputStats {
name: input_name.to_string(),
input_type,
error_count,
raw_stats: PlaylistStats {
group_count,
channel_count,
},
processed_stats: PlaylistStats {
group_count: 0,
channel_count: 0,
},
secs_took,
}
}
async fn process_sources(client: Arc<reqwest::Client>, config: Arc<Config>, user_targets: Arc<ProcessTargets>) -> (Vec<SourceStats>, Vec<M3uFilterError>) {
let mut handle_list = vec![];
let thread_num = config.threads;
let process_parallel = thread_num > 1 && config.sources.len() > 1;
if process_parallel && log_enabled!(Level::Debug) {
debug!("Using {} threads", thread_num);
}
let errors = Arc::new(Mutex::<Vec<M3uFilterError>>::new(vec![]));
let stats = Arc::new(Mutex::<Vec<SourceStats>>::new(vec![]));
for (index, _) in config.sources.iter().enumerate() {
// We're using the file lock this way on purpose
let source_lock_path = PathBuf::from(format!("source_{index}"));
let Ok(update_lock) = config.file_locks.try_write_lock(&source_lock_path).await else {
warn!("The update operation for the source at index {index} was skipped because an update is already in progress.");
continue;
};
let shared_errors = errors.clone();
let shared_stats = stats.clone();
let cfg = config.clone();
let usr_trgts = user_targets.clone();
if process_parallel {
let http_client = Arc::clone(&client);
let handles = &mut handle_list;
let process = move || {
System::new().block_on(async {
let (input_stats, target_stats, mut res_errors) = process_source(Arc::clone(&http_client), cfg, index, usr_trgts).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);
});
};
handles.push(thread::spawn(process));
if handles.len() >= thread_num as usize {
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;
shared_errors.lock().await.append(&mut res_errors);
let process_stats = SourceStats::new(input_stats, target_stats);
shared_stats.lock().await.push(process_stats);
}
drop(update_lock);
}
for handle in handle_list {
let _ = handle.join();
}
(Arc::try_unwrap(stats).unwrap().into_inner(), Arc::try_unwrap(errors).unwrap().into_inner())
}
pub type ProcessingPipe = Vec<fn(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option<Vec<PlaylistGroup>>>;
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]
}
}
fn duplicate_hash(item: &PlaylistItem) -> UUIDType {
item.get_uuid()
}
fn execute_pipe<'a>(target: &ConfigTarget, pipe: &ProcessingPipe, fpl: &FetchedPlaylist<'a>, duplicates: &mut HashSet<UUIDType>) -> FetchedPlaylist<'a> {
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(),
};
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)));
}
}
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.
fn flatten_groups(playlistgroups: Vec<PlaylistGroup>) -> Vec<PlaylistGroup> {
let mut sort_order: Vec<PlaylistGroup> = vec![];
let mut idx: usize = 0;
let mut group_map: HashMap<(Rc<String>, XtreamCluster), usize> = HashMap::new();
for group in playlistgroups {
let key = (Rc::clone(&group.title), group.xtream_cluster);
match group_map.entry(key) {
std::collections::hash_map::Entry::Vacant(v) => {
v.insert(idx);
idx += 1;
sort_order.push(group);
}
std::collections::hash_map::Entry::Occupied(o) => {
sort_order.get_mut(*o.get()).unwrap().channels.extend(group.channels);
}
};
}
sort_order
}
async fn process_playlist_for_target(client: Arc<reqwest::Client>,
playlists: &mut [FetchedPlaylist<'_>],
target: &ConfigTarget,
cfg: &Config,
stats: &mut HashMap<String, InputStats>,
errors: &mut Vec<M3uFilterError>) -> Result<(), Vec<M3uFilterError>> {
let pipe = get_processing_pipe(target);
debug_if_enabled!("Processing order is {}", &target.processing_order);
let mut duplicates: HashSet<UUIDType> = HashSet::new();
let mut processed_fetched_playlists: Vec<FetchedPlaylist> = vec![];
for provider_fpl in playlists.iter_mut() {
let mut processed_fpl = execute_pipe(target, &pipe, provider_fpl, &mut duplicates);
playlist_resolve_series(Arc::clone(&client), cfg, target, errors, &pipe, provider_fpl, &mut processed_fpl).await;
playlist_resolve_vod(Arc::clone(&client), cfg, target, errors, &processed_fpl).await;
// stats
let input_stats = stats.get_mut(&processed_fpl.input.name);
if let Some(stat) = input_stats {
stat.processed_stats.group_count = processed_fpl.playlistgroups.len();
stat.processed_stats.channel_count = processed_fpl.playlistgroups.iter()
.map(|group| group.channels.len())
.sum();
}
processed_fetched_playlists.push(processed_fpl);
}
apply_affixes(&mut processed_fetched_playlists);
let mut new_playlist = vec![];
let mut new_epg = vec![];
// each fetched playlist can have its own epgl url.
// we need to process each input epg.
for mut fp in processed_fetched_playlists {
// collect all epg_channel ids
let epg_channel_ids: HashSet<_> = fp.playlistgroups.iter().flat_map(|g| &g.channels)
.filter_map(|c| c.header.borrow().epg_channel_id.clone()).collect();
new_playlist.append(&mut fp.playlistgroups);
if epg_channel_ids.is_empty() {
debug_if_enabled!("channel ids are empty");
} else if let Some(tv_guide) = fp.epg {
debug!("found epg information for {}", &target.name);
if let Some(epg) = tv_guide.filter(&epg_channel_ids) {
new_epg.push(epg);
}
}
}
if new_playlist.is_empty() {
info!("Playlist is empty: {}", &target.name);
Ok(())
} else {
let mut flat_new_playlist = flatten_groups(new_playlist);
sort_playlist(target, &mut flat_new_playlist);
channel_no_playlist(&flat_new_playlist);
map_playlist_counter(target, &flat_new_playlist);
process_watch(target, cfg, &flat_new_playlist);
persist_playlist(&mut flat_new_playlist, flatten_tvguide(&new_epg).as_ref(), target, cfg).await
}
}
fn process_watch(target: &ConfigTarget, cfg: &Config, new_playlist: &Vec<PlaylistGroup>) {
if target.t_watch_re.is_some() {
if default_as_default().eq_ignore_ascii_case(&target.name) {
error!("cant watch a target with no unique name");
} else {
let watch_re = target.t_watch_re.as_ref().unwrap();
for pl in new_playlist {
if watch_re.iter().any(|r| r.is_match(&pl.title)) {
process_group_watch(cfg, &target.name, pl);
}
}
}
}
}
pub async fn exec_processing(client: Arc<reqwest::Client>, cfg: Arc<Config>, targets: Arc<ProcessTargets>) {
let start_time = Instant::now();
let (stats, errors) = process_sources(client, cfg.clone(), targets.clone()).await;
// log errors
for err in &errors {
error!("{}", err.message);
}
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
info!("{}", stats_msg);
// send stats
send_message(&MsgKind::Stats, cfg.messaging.as_ref(), stats_msg.as_str());
}
// send errors
if let Some(message) = get_errors_notify_message!(errors, 255) {
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(&MsgKind::Error, cfg.messaging.as_ref(), error_msg.as_str());
}
}
let elapsed = start_time.elapsed().as_secs();
info!("Update process finished! Took {elapsed} secs.");
}
+122
View File
@@ -0,0 +1,122 @@
use crate::m3u_filter_error::{str_to_io_error, to_io_error, M3uFilterError, M3uFilterErrorKind};
use crate::model::config::{Config, ConfigInput};
use crate::model::playlist::{FetchedPlaylist, PlaylistEntry, PlaylistItem, PlaylistItemType, XtreamCluster};
use crate::repository::storage::get_input_storage_path;
use crate::m3u_filter_error::{info_err, notify_err};
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
use std::fs::File;
use std::io::{BufWriter, Write};
use std::path::PathBuf;
use std::sync::Arc;
const FILE_SERIES_INFO: &str = "xtream_series_info";
const FILE_VOD_INFO: &str = "xtream_vod_info";
const FILE_SUFFIX_WAL: &str = "wal";
const FILE_SERIES_EPISODE_RECORD: &str = "series_episode_record";
use crate::repository::bplustree::BPlusTree;
use crate::repository::xtream_repository::xtream_get_record_file_path;
use crate::utils::file::file_utils::append_or_crate_file;
use crate::utils::network::xtream;
pub(in crate::processing) async fn playlist_resolve_download_playlist_item(client: Arc<reqwest::Client>, pli: &PlaylistItem, input: &ConfigInput, errors: &mut Vec<M3uFilterError>, 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) {
result = match xtream::get_xtream_stream_info_content(client, &info_url, input).await {
Ok(content) => Some(content),
Err(err) => {
errors.push(info_err!(format!("{err}")));
None
}
};
}
if resolve_delay > 0 {
actix_web::rt::time::sleep(std::time::Duration::new(u64::from(resolve_delay), 0)).await;
}
result
}
pub(in crate::processing) fn write_info_content_to_wal_file(writer: &mut BufWriter<&File>, provider_id: u32, content: &str) -> std::io::Result<()> {
let length = u32::try_from(content.len()).map_err(to_io_error)?;
if length > 0 {
writer.write_all(&provider_id.to_le_bytes())?;
writer.write_all(&length.to_le_bytes())?;
writer.write_all(content.as_bytes())?;
}
Ok(())
}
pub(in crate::processing) fn create_resolve_episode_wal_files(cfg: &Config, input: &ConfigInput) -> Option<(File, PathBuf)> {
match get_input_storage_path(input, &cfg.working_dir) {
Ok(storage_path) => {
let info_path = storage_path.join(format!("{FILE_SERIES_EPISODE_RECORD}.{FILE_SUFFIX_WAL}"));
let info_file = append_or_crate_file(&info_path).ok()?;
Some((info_file, info_path))
}
Err(_) => None
}
}
pub(in crate::processing) fn create_resolve_info_wal_files(cfg: &Config, input: &ConfigInput, cluster: XtreamCluster) -> Option<(File, File, PathBuf, PathBuf)> {
match get_input_storage_path(input, &cfg.working_dir) {
Ok(storage_path) => {
if let Some(file_prefix) = match cluster {
XtreamCluster::Live => None,
XtreamCluster::Video => Some(FILE_VOD_INFO),
XtreamCluster::Series => Some(FILE_SERIES_INFO)
} {
let content_path = storage_path.join(format!("{file_prefix}_content.{FILE_SUFFIX_WAL}"));
let info_path = storage_path.join(format!("{file_prefix}_record.{FILE_SUFFIX_WAL}"));
let content_file = append_or_crate_file(&content_path).ok()?;
let info_file = append_or_crate_file(&info_path).ok()?;
return Some((content_file, info_file, content_path, info_path));
}
None
}
Err(_) => None
}
}
pub(in crate::processing) fn should_update_info(pli: &PlaylistItem, processed_provider_ids: &HashMap<u32, u64>, field: &str) -> (bool, u32, u64) {
let Some(provider_id) = pli.header.borrow_mut().get_provider_id() else { return (false, 0, 0) };
let last_modified = pli.header.borrow().get_additional_property_as_u64(field);
let old_timestamp = processed_provider_ids.get(&provider_id);
(old_timestamp.is_none()
|| last_modified.is_none()
|| *old_timestamp.unwrap() != last_modified.unwrap(), provider_id, last_modified.unwrap_or(0))
}
pub(in crate::processing) async fn read_processed_info_ids<V, F>(cfg: &Config, errors: &mut Vec<M3uFilterError>, fpl: &FetchedPlaylist<'_>,
item_type: PlaylistItemType, extract_ts: F) -> HashMap<u32, u64>
where
F: Fn(&V) -> u64,
V: Serialize + for<'de> Deserialize<'de> + Clone,
{
let mut processed_info_ids = HashMap::new();
let file_path = match get_input_storage_path(fpl.input, &cfg.working_dir)
.map(|storage_path| xtream_get_record_file_path(&storage_path, item_type)).and_then(|opt| opt.ok_or_else(|| str_to_io_error("Not supported")))
{
Ok(file_path) => file_path,
Err(err) => {
let fpl_name = &fpl.input.name;
errors.push(notify_err!(format!("Could not create storage path for input {fpl_name}: {err}")));
return processed_info_ids;
}
};
match cfg.file_locks.read_lock(&file_path).await {
Ok(file_lock) => {
if let Ok(info_records) = BPlusTree::<u32, V>::load(&file_path) {
info_records.iter().for_each(|(provider_id, record)| {
processed_info_ids.insert(*provider_id, extract_ts(record));
});
}
drop(file_lock);
}
Err(err) => errors.push(info_err!(format!("{err}"))),
}
processed_info_ids
}
+241
View File
@@ -0,0 +1,241 @@
use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind};
use crate::model::config::{Config, ConfigTarget, InputType};
use crate::model::playlist::{FetchedPlaylist, PlaylistGroup, PlaylistItem, PlaylistItemType, XtreamCluster};
use crate::processing::processor::playlist::ProcessingPipe;
use crate::processing::parser::xtream::parse_xtream_series_info;
use crate::processing::processor::xtream::{create_resolve_episode_wal_files, create_resolve_info_wal_files, playlist_resolve_download_playlist_item, read_processed_info_ids, should_update_info, write_info_content_to_wal_file};
use crate::repository::storage::get_input_storage_path;
use crate::repository::xtream_repository::{xtream_get_info_file_paths, xtream_update_input_info_file, xtream_update_input_series_episodes_record_from_wal_file, xtream_update_input_series_record_from_wal_file};
use crate::repository::IndexedDocumentReader;
use crate::m3u_filter_error::{notify_err, info_err};
use crate::processing::processor::{handle_error, handle_error_and_return, create_resolve_options_function_for_xtream_target};
use std::collections::HashMap;
use std::fs::File;
use std::io::{BufWriter, Write};
use std::sync::Arc;
use std::time::Instant;
use log::{info, log_enabled, Level};
use crate::model::xtream::{XtreamSeriesEpisode, XtreamSeriesInfoEpisode};
use crate::utils::file::file_utils::file_writer;
const TAG_SERIES_INFO_LAST_MODIFIED: &str = "last_modified";
create_resolve_options_function_for_xtream_target!(series);
async fn read_processed_series_info_ids(cfg: &Config, errors: &mut Vec<M3uFilterError>, fpl: &FetchedPlaylist<'_>) -> HashMap<u32, u64> {
read_processed_info_ids(cfg, errors, fpl, PlaylistItemType::SeriesInfo, |ts: &u64| *ts).await
}
fn write_series_info_record_to_wal_file(
writer: &mut BufWriter<&File>,
provider_id: u32,
ts: u64,
) -> std::io::Result<()> {
writer.write_all(&provider_id.to_le_bytes())?;
writer.write_all(&ts.to_le_bytes())?;
Ok(())
}
fn write_series_episode_record_to_wal_file(
writer: &mut BufWriter<&File>,
provider_id: u32,
episode: &XtreamSeriesInfoEpisode,
) -> std::io::Result<()> {
let series_episode = XtreamSeriesEpisode::from(episode);
if let Ok(content_bytes) = bincode::serialize(&series_episode) {
writer.write_all(&provider_id.to_le_bytes())?;
let len = u32::try_from(content_bytes.len()).unwrap();
writer.write_all(&len.to_le_bytes())?;
writer.write_all(&content_bytes)?;
}
Ok(())
}
fn should_update_series_info(pli: &PlaylistItem, processed_provider_ids: &HashMap<u32, u64>) -> (bool, u32, u64) {
should_update_info(pli, processed_provider_ids, TAG_SERIES_INFO_LAST_MODIFIED)
}
async fn playlist_resolve_series_info(client: Arc<reqwest::Client>, cfg: &Config, errors: &mut Vec<M3uFilterError>,
fpl: &mut FetchedPlaylist<'_>, resolve_delay: u16) -> bool {
let mut processed_info_ids = read_processed_series_info_ids(cfg, errors, fpl).await;
// 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.
let Some((wal_content_file, wal_record_file, wal_content_path, wal_record_path)) = create_resolve_info_wal_files(cfg, fpl.input, XtreamCluster::Series)
else { return !processed_info_ids.is_empty(); };
let mut content_writer = file_writer(&wal_content_file);
let mut record_writer = file_writer(&wal_record_file);
let mut content_updated = false;
let series_info_iter = fpl.playlistgroups.iter()
.filter(|&plg| plg.xtream_cluster == XtreamCluster::Series)
.flat_map(|plg| &plg.channels)
.filter(|&pli| pli.header.borrow().item_type == PlaylistItemType::SeriesInfo);
let series_info_count = series_info_iter.clone().count();
info!("Found {series_info_count} series info to resolve");
let start_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 {
if let Some(content) = playlist_resolve_download_playlist_item(Arc::clone(&client), pli, fpl.input, errors, resolve_delay, XtreamCluster::Series).await {
handle_error_and_return!(write_info_content_to_wal_file(&mut content_writer, provider_id, &content),
|err| errors.push(notify_err!(format!("Failed to resolve series, could not write to content wal file {err}"))));
processed_info_ids.insert(provider_id, ts);
handle_error_and_return!(write_series_info_record_to_wal_file(&mut record_writer, provider_id, ts),
|err| errors.push(notify_err!(format!("Failed to resolve series wal, could not write to record wal file {err}"))));
content_updated = true;
}
}
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 % 30 == 0) {
info!("resolved {processed_series_info_count}/{series_info_count} series info");
last_processed_series_info_count = processed_series_info_count;
}
}
}
if last_processed_series_info_count != processed_series_info_count {
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}"))));
drop(content_writer);
drop(wal_content_file);
handle_error!(record_writer.flush(),
|err| errors.push(notify_err!(format!("Failed to resolve vod tmdb, could not write to wal file {err}"))));
drop(record_writer);
drop(wal_record_file);
handle_error!(xtream_update_input_info_file(cfg, fpl.input, &wal_content_path, XtreamCluster::Series).await,
|err| errors.push(err));
handle_error!(xtream_update_input_series_record_from_wal_file(cfg, fpl.input, &wal_record_path).await,
|err| errors.push(err));
}
// we updated now
// - series_info.db which contains the original series_info json
// - series_record.db which contains the series_info provider_id and timestamp
!processed_info_ids.is_empty()
}
async fn process_series_info(
cfg: &Config,
fpl: &mut FetchedPlaylist<'_>,
errors: &mut Vec<M3uFilterError>,
) -> Vec<PlaylistGroup> {
let mut result: Vec<PlaylistGroup> = vec![];
let input = fpl.input;
let Ok(Some((info_path, idx_path))) = get_input_storage_path(input, &cfg.working_dir)
.map(|storage_path| xtream_get_info_file_paths(&storage_path, XtreamCluster::Series))
else {
errors.push(notify_err!("Failed to open input info file for series".to_string()));
return result;
};
let Ok(_file_lock) = cfg.file_locks.read_lock(&info_path).await else {
errors.push(notify_err!("Could not lock input info file for series".to_string()));
return result;
};
// Contains the Series Info with episode listing
let Ok(mut info_reader) = IndexedDocumentReader::<u32, String>::new(&info_path, &idx_path) else { return result; };
let Some((wal_file, wal_path)) = create_resolve_episode_wal_files(cfg, input) else {
errors.push(notify_err!("Could not create wal file for series episodes record".to_string()));
return result;
};
let mut wal_writer = file_writer(&wal_file);
for plg in fpl
.playlistgroups
.iter_mut()
.filter(|plg| plg.xtream_cluster == XtreamCluster::Series)
{
let mut group_series = vec![];
for pli in plg
.channels
.iter()
.filter(|pli| pli.header.borrow().item_type == PlaylistItemType::SeriesInfo)
{
let Some(provider_id) = pli.header.borrow_mut().get_provider_id() else { continue; };
let Ok(content) = info_reader.get(&provider_id) else { continue; };
match serde_json::from_str::<serde_json::Value>(&content) {
Ok(series_content) => {
let (group, series_name) = {
let header = pli.header.borrow();
(header.group.clone(), if header.name.is_empty() {header.title.clone()} else { header.name.clone()})
};
match parse_xtream_series_info(&series_content, &group, &series_name, input) {
Ok(Some(series)) => {
for (episode, pli_episode) in &series {
let Some(provider_id) = &pli_episode.header.borrow_mut().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}"))));
}
group_series.extend(series.into_iter().map(|(_, pli)| pli));
}
Ok(None) => {}
Err(err) => {
errors.push(err);
}
}
}
Err(err) => errors.push(info_err!(format!("Failed to parse JSON: {err}"))),
}
}
if !group_series.is_empty() {
result.push(PlaylistGroup {
id: plg.id,
title: plg.title.clone(),
channels: group_series,
xtream_cluster: XtreamCluster::Series,
});
}
}
handle_error!(wal_writer.flush(),
|err| errors.push(notify_err!(format!("Failed to resolve series episodes, could not write to wal file {err}"))));
drop(wal_writer);
drop(wal_file);
handle_error!(xtream_update_input_series_episodes_record_from_wal_file(cfg, input, &wal_path).await,
|err| errors.push(err));
result
}
pub async fn playlist_resolve_series(client: Arc<reqwest::Client>, cfg: &Config, target: &ConfigTarget,
errors: &mut Vec<M3uFilterError>,
pipe: &ProcessingPipe,
provider_fpl: &mut FetchedPlaylist<'_>,
processed_fpl: &mut FetchedPlaylist<'_>,
) {
let (resolve_series, resolve_delay) = get_resolve_series_options(target, processed_fpl);
if !resolve_series { return; }
if !playlist_resolve_series_info(client, cfg, errors, processed_fpl, resolve_delay).await { return; }
let series_playlist = process_series_info(cfg, provider_fpl, errors).await;
if series_playlist.is_empty() { return; }
// original content saved into original list
for plg in &series_playlist {
provider_fpl.update_playlist(plg);
}
// run processing pipe over new items
let mut new_playlist = series_playlist;
for f in pipe {
if let Some(v) = f(&mut new_playlist, target) {
new_playlist = v;
}
}
// assign new items to the new playlist
for plg in &new_playlist {
processed_fpl.update_playlist(plg);
}
}
+136
View File
@@ -0,0 +1,136 @@
use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind};
use crate::model::config::{Config, ConfigTarget, InputType};
use crate::model::playlist::{FetchedPlaylist, PlaylistItem, PlaylistItemType, XtreamCluster};
use crate::processing::processor::xtream::{create_resolve_info_wal_files, playlist_resolve_download_playlist_item, read_processed_info_ids, should_update_info, write_info_content_to_wal_file};
use crate::repository::xtream_repository::{xtream_update_input_info_file, xtream_update_input_vod_record_from_wal_file, InputVodInfoRecord};
use crate::m3u_filter_error::{notify_err};
use crate::processing::processor::{handle_error, handle_error_and_return, create_resolve_options_function_for_xtream_target};
use crate::utils::json_utils::{get_u32_from_serde_value, get_u64_from_serde_value};
use serde_json::{Map, Value};
use std::collections::HashMap;
use std::fs::File;
use std::io::{BufWriter, Write};
use std::sync::Arc;
use std::time::Instant;
use log::{info, log_enabled, Level};
use crate::utils::file::file_utils::file_writer;
const TAG_VOD_INFO_INFO: &str = "info";
const TAG_VOD_INFO_MOVIE_DATA: &str = "movie_data";
const TAG_VOD_INFO_TMDB_ID: &str = "tmdb_id";
const TAG_VOD_INFO_STREAM_ID: &str = "stream_id";
const TAG_VOD_INFO_ADDED: &str = "added";
create_resolve_options_function_for_xtream_target!(vod);
async fn read_processed_vod_info_ids(cfg: &Config, errors: &mut Vec<M3uFilterError>, fpl: &FetchedPlaylist<'_>) -> HashMap<u32, u64> {
read_processed_info_ids(cfg, errors, fpl, PlaylistItemType::Video, |record: &InputVodInfoRecord| record.ts).await
}
fn extract_info_record_from_vod_info(content: &str) -> Option<(u32, InputVodInfoRecord)> {
let doc = serde_json::from_str::<Map<String, Value>>(content).ok()?;
let movie_data = doc.get(TAG_VOD_INFO_MOVIE_DATA)?.as_object()?;
let provider_id = get_u32_from_serde_value(
movie_data.get(TAG_VOD_INFO_STREAM_ID)?,
)?;
let added = movie_data
.get(TAG_VOD_INFO_ADDED)
.and_then(get_u64_from_serde_value)
.unwrap_or(0);
let tmdb_id = doc.get(TAG_VOD_INFO_INFO)?.as_object()
.and_then(|info| info.get(TAG_VOD_INFO_TMDB_ID))
.and_then(get_u32_from_serde_value)
.unwrap_or(0);
Some((provider_id, InputVodInfoRecord {
tmdb_id,
ts: added,
}))
}
fn write_vod_info_record_to_wal_file(
writer: &mut BufWriter<&File>,
provider_id: u32,
record: &InputVodInfoRecord,
) -> std::io::Result<()> {
writer.write_all(&provider_id.to_le_bytes())?;
writer.write_all(&record.tmdb_id.to_le_bytes())?;
writer.write_all(&record.ts.to_le_bytes())?;
Ok(())
}
fn should_update_vod_info(pli: &PlaylistItem, processed_provider_ids: &HashMap<u32, u64>) -> (bool, u32, u64) {
should_update_info(pli, processed_provider_ids, TAG_VOD_INFO_ADDED)
}
pub async fn playlist_resolve_vod(client: Arc<reqwest::Client>, cfg: &Config, target: &ConfigTarget, errors: &mut Vec<M3uFilterError>, fpl: &FetchedPlaylist<'_>) {
let (resolve_movies, resolve_delay) = get_resolve_vod_options(target, fpl);
if !resolve_movies { return; }
// 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.
let Some((wal_content_file, wal_record_file, wal_content_path, wal_record_path)) = create_resolve_info_wal_files(cfg, fpl.input, XtreamCluster::Video)
else { return; };
let mut processed_info_ids = read_processed_vod_info_ids(cfg, errors, fpl).await;
let mut content_writer = file_writer(&wal_content_file);
let mut record_writer = file_writer(&wal_record_file);
let mut content_updated = false;
let vod_info_iter = fpl.playlistgroups.iter()
.flat_map(|plg| &plg.channels)
.filter(|&pli| pli.header.borrow().xtream_cluster == XtreamCluster::Video);
let vod_info_count = vod_info_iter.clone().count();
info!("Found {vod_info_count} vod info to resolve");
let start_time = Instant::now();
let mut processed_vod_info_count = 0;
let mut last_processed_vod_info_count = 0;
for pli in vod_info_iter {
let (should_update, _provider_id, _ts) = should_update_vod_info(pli, &processed_info_ids);
if should_update {
if let Some(content) = playlist_resolve_download_playlist_item(Arc::clone(&client), pli, fpl.input, errors, resolve_delay, XtreamCluster::Video).await {
if let Some((provider_id, info_record)) = extract_info_record_from_vod_info(&content) {
let ts = info_record.ts;
handle_error_and_return!(write_info_content_to_wal_file(&mut content_writer, provider_id, &content),
|err| errors.push(notify_err!(format!("Failed to resolve vod, could not write to content wal file {err}"))));
processed_info_ids.insert(provider_id, ts);
handle_error_and_return!(write_vod_info_record_to_wal_file(&mut record_writer, provider_id, &info_record),
|err| errors.push(notify_err!(format!("Failed to resolve vod wal, could not write to record wal file {err}"))));
content_updated = true;
}
}
}
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 % 30 == 0) {
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");
}
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}"))));
handle_error!(record_writer.flush(),
|err| errors.push(notify_err!(format!("Failed to resolve vod tmdb, could not write to wal file {err}"))));
drop(content_writer);
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::Video).await,
|err| errors.push(err));
handle_error!(xtream_update_input_vod_record_from_wal_file(cfg, fpl.input, &wal_record_path).await,
|err| errors.push(err));
}
}