mirror of
https://github.com/euzu/tuliprox.git
synced 2026-09-30 13:02:10 +02:00
project shared added
This commit is contained in:
@@ -0,0 +1,59 @@
|
||||
use crate::model::{ConfigInput, InputAffix, valid_property, MAPPER_FIELDS};
|
||||
use crate::model::{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;
|
||||
let value = header.get_field(affix.field.as_str()).map_or_else(|| String::from(&affix.value), |field_value| if is_prefix {
|
||||
format!("{}{field_value}", &affix.value, )
|
||||
} else {
|
||||
format!("{field_value}{}", , &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(), MAPPER_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);
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,255 @@
|
||||
use crate::model::{Epg, TVGuide, XmlTag, XmlTagIcon, EPG_ATTRIB_ID};
|
||||
use crate::model::{EpgConfig, EpgSmartMatchConfig};
|
||||
use crate::model::{FetchedPlaylist, PlaylistItem};
|
||||
use crate::processing::parser::xmltv::normalize_channel_name;
|
||||
use log::{debug, trace};
|
||||
use rphonetic::{DoubleMetaphone, Encoder};
|
||||
use std::borrow::Cow;
|
||||
use std::collections::{HashMap, HashSet};
|
||||
use shared::model::XtreamCluster;
|
||||
|
||||
pub struct EpgIdCache<'a> {
|
||||
pub channel_epg_id: HashSet<Cow<'a, str>>,
|
||||
pub normalized: HashMap<String, Option<String>>,
|
||||
pub phonetics: HashMap<String, HashSet<String>>,
|
||||
pub processed: HashSet<String>,
|
||||
pub smart_match_config: EpgSmartMatchConfig,
|
||||
pub metaphone: DoubleMetaphone,
|
||||
pub smart_match_enabled: bool, // smart match is enabled, normalizing names
|
||||
pub fuzzy_match_enabled: bool, // fuzzy matching enabled
|
||||
}
|
||||
|
||||
impl EpgIdCache<'_> {
|
||||
/// Creates a new `EpgIdCache` with configuration for smart and fuzzy matching.
|
||||
///
|
||||
/// Initializes all internal caches and sets matching options based on the provided EPG configuration. If no configuration is given, defaults are used.
|
||||
///
|
||||
/// # Examples
|
||||
///
|
||||
/// ```
|
||||
/// let cache = EpgIdCache::new(None);
|
||||
/// assert!(cache.is_empty());
|
||||
/// ```
|
||||
pub fn new(epg_config: Option<&EpgConfig>) -> Self {
|
||||
let normalize_config = epg_config.map_or_else(EpgSmartMatchConfig::default, |epg_config| epg_config.t_smart_match.clone());
|
||||
EpgIdCache {
|
||||
channel_epg_id: HashSet::new(), // contains the epg_ids collected from playlist channels
|
||||
normalized: HashMap::new(),
|
||||
phonetics: HashMap::new(),
|
||||
processed: HashSet::new(),
|
||||
metaphone: DoubleMetaphone::default(),
|
||||
smart_match_enabled: normalize_config.enabled,
|
||||
fuzzy_match_enabled: normalize_config.enabled && normalize_config.fuzzy_matching,
|
||||
smart_match_config: normalize_config,
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
fn is_empty(&self) -> bool {
|
||||
self.channel_epg_id.is_empty() && self.normalized.is_empty()
|
||||
}
|
||||
|
||||
/// Normalizes a channel name, computes its phonetic encoding, and stores both in the cache for later EPG matching.
|
||||
///
|
||||
/// The normalized name is mapped to the provided EPG ID (if any), and the phonetic encoding is added to the phonetics map.
|
||||
/// This facilitates efficient lookup and fuzzy matching of channel names during EPG assignment.
|
||||
///
|
||||
/// # Examples
|
||||
///
|
||||
/// ```
|
||||
/// let mut cache = EpgIdCache::new(None);
|
||||
/// cache.normalize_and_store("Discovery Channel", Some(&"discovery.epg".to_string()));
|
||||
/// assert!(cache.normalized.contains_key(&cache.normalize("Discovery Channel")));
|
||||
/// ```
|
||||
fn normalize_and_store(&mut self, name: &str, epg_id: Option<&String>) {
|
||||
let normalized_name = self.normalize(name);
|
||||
let phonetic = self.phonetic(&normalized_name);
|
||||
self.normalized.insert(normalized_name.to_string(), epg_id.map(std::string::ToString::to_string));
|
||||
self.phonetics.entry(phonetic.to_string()).or_default().insert(normalized_name);
|
||||
}
|
||||
|
||||
/// Returns the normalized form of a channel name using the configured smart match settings.
|
||||
///
|
||||
/// # Examples
|
||||
///
|
||||
/// ```
|
||||
/// let cache = EpgIdCache::new(None);
|
||||
/// let normalized = cache.normalize("HBO HD");
|
||||
/// assert!(!normalized.is_empty());
|
||||
/// ```
|
||||
fn normalize(&self, name: &str) -> String {
|
||||
normalize_channel_name(name, &self.smart_match_config)
|
||||
}
|
||||
|
||||
pub(crate) fn phonetic(&self, name: &str) -> String {
|
||||
self.metaphone.encode(name)
|
||||
}
|
||||
|
||||
pub fn collect_epg_id(&mut self, fp: &mut FetchedPlaylist) {
|
||||
let smart_match_enabled = self.smart_match_enabled;
|
||||
let fuzzy_matching = self.fuzzy_match_enabled;
|
||||
|
||||
for channel in fp.playlistgroups.iter().flat_map(|g| &g.channels) {
|
||||
let mut missing_epg_id = true;
|
||||
// insert epg_id to known channel epg_ids
|
||||
if let Some(id) = channel.header.epg_channel_id.as_deref() {
|
||||
if !id.is_empty() {
|
||||
missing_epg_id = false;
|
||||
self.channel_epg_id.insert(Cow::Owned(id.to_string()));
|
||||
}
|
||||
}
|
||||
|
||||
// for fuzzy_matching we need to put the normalized name even if there is an epg_id, because the epg_id
|
||||
// could not match to the epg file. And then we try to guess it based on normalized name
|
||||
let needs_normalization = smart_match_enabled && (fuzzy_matching || missing_epg_id);
|
||||
|
||||
if needs_normalization {
|
||||
let name = &channel.header.name;
|
||||
self.normalize_and_store(name, channel.header.epg_channel_id.as_ref());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub fn match_with_normalized(&mut self, epg_id: &str, normalized_epg_ids: &[String]) -> bool {
|
||||
for key in normalized_epg_ids {
|
||||
if let Some(entry) = self.normalized.get_mut(key) {
|
||||
entry.replace(epg_id.to_string());
|
||||
self.channel_epg_id.insert(epg_id.to_string().into());
|
||||
return true;
|
||||
}
|
||||
}
|
||||
false
|
||||
}
|
||||
}
|
||||
|
||||
/// Assigns EPG IDs and logos to live playlist channels by matching them with EPG data.
|
||||
///
|
||||
/// For each live channel in the playlist missing an EPG ID, attempts to assign one using normalized name matching if smart matching is enabled. If a channel has an EPG ID but lacks logos, assigns logos from the corresponding EPG icon tags. Adds the matched EPG data to the provided vector.
|
||||
///
|
||||
/// # Examples
|
||||
///
|
||||
/// ```
|
||||
/// let mut new_epg = Vec::new();
|
||||
/// let mut playlist = FetchedPlaylist::default();
|
||||
/// let mut id_cache = EpgIdCache::new(None);
|
||||
/// assign_channel_epg(&mut new_epg, &mut playlist, &mut id_cache);
|
||||
/// ```
|
||||
fn assign_channel_epg(new_epg: &mut Vec<Epg>, fp: &mut FetchedPlaylist, id_cache: &mut EpgIdCache) {
|
||||
id_cache.normalized.retain(|_, v| v.is_some());
|
||||
if let Some(tv_guide) = &fp.epg {
|
||||
let mut processed_epgs = vec![];
|
||||
if let Some(epg_sources) = tv_guide.filter(id_cache) {
|
||||
let mut icon_assigned = HashSet::new();
|
||||
for epg_source in epg_sources {
|
||||
// icon tags
|
||||
let icon_tags: HashMap<&String, &XmlTag> = epg_source.children.iter()
|
||||
.filter(|tag| tag.icon != XmlTagIcon::Undefined && tag.get_attribute_value(EPG_ATTRIB_ID).is_some())
|
||||
.map(|t| (t.get_attribute_value(EPG_ATTRIB_ID).unwrap(), t)).collect();
|
||||
|
||||
let assign_values = |chan: &mut PlaylistItem| {
|
||||
if id_cache.smart_match_enabled && chan.header.epg_channel_id.is_none() {
|
||||
// if the channel has no epg_id or the epg_id is not present in xmltv/tvguide then we need to match one from existing tvguide
|
||||
let not_processed = match &chan.header.epg_channel_id {
|
||||
None => true,
|
||||
Some(epg_id) => !id_cache.processed.contains(epg_id),
|
||||
};
|
||||
if not_processed {
|
||||
let normalized = id_cache.normalize(&chan.header.name);
|
||||
if let Some(epg_id) = id_cache.normalized.get(&normalized) {
|
||||
if epg_id.is_some() {
|
||||
trace!("Matched channel {} to epg {epg_id:?}", chan.header.name);
|
||||
chan.header.epg_channel_id.clone_from(epg_id);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if let Some(epg_channel_id) = chan.header.epg_channel_id.as_ref() {
|
||||
if !icon_assigned.contains(epg_channel_id) &&
|
||||
(epg_source.logo_override || chan.header.logo.is_empty() || chan.header.logo_small.is_empty()) {
|
||||
if let Some(icon_tag) = icon_tags.get(chan.header.epg_channel_id.as_ref().unwrap()) {
|
||||
if let XmlTagIcon::Src(icon) = &icon_tag.icon {
|
||||
icon_assigned.insert(epg_channel_id.to_string());
|
||||
if epg_source.logo_override || chan.header.logo.is_empty() {
|
||||
trace!("Matched channel {} to epg icon {icon}", chan.header.name);
|
||||
chan.header.logo = (*icon).to_string();
|
||||
}
|
||||
if epg_source.logo_override || chan.header.logo_small.is_empty() {
|
||||
chan.header.logo_small = (*icon).to_string();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
let filter_live = |c: &&mut PlaylistItem| c.header.xtream_cluster == XtreamCluster::Live;
|
||||
fp.playlistgroups.iter_mut()
|
||||
.flat_map(|g| &mut g.channels)
|
||||
.filter(filter_live)
|
||||
.for_each(assign_values);
|
||||
processed_epgs.push(epg_source);
|
||||
}
|
||||
}
|
||||
|
||||
if let Some(epg) = TVGuide::merge(processed_epgs) {
|
||||
new_epg.push(epg);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Processes a fetched playlist and assigns EPG data to its channels.
|
||||
///
|
||||
/// Collects EPG channel IDs from the playlist, initializes an EPG ID cache, and assigns EPG data to channels using normalization and smart matching if enabled. Logs a debug message if no EPG IDs are found and smart matching is disabled.
|
||||
///
|
||||
/// # Examples
|
||||
///
|
||||
/// ```
|
||||
/// let mut playlist = FetchedPlaylist::default();
|
||||
/// let mut epg_data = Vec::new();
|
||||
/// process_playlist_epg(&mut playlist, &mut epg_data);
|
||||
/// ```
|
||||
pub fn process_playlist_epg(fp: &mut FetchedPlaylist, epg: &mut Vec<Epg>) {
|
||||
// collect all epg_channel ids
|
||||
let mut id_cache = EpgIdCache::new(fp.input.epg.as_ref());
|
||||
id_cache.collect_epg_id(fp);
|
||||
|
||||
if id_cache.is_empty() && !id_cache.smart_match_enabled {
|
||||
debug!("No epg ids found");
|
||||
} else {
|
||||
assign_channel_epg(epg, fp, &mut id_cache);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use rand::distr::Alphanumeric;
|
||||
use rand::Rng;
|
||||
use rphonetic::{DoubleMetaphone, Encoder};
|
||||
use tokio::time::Instant;
|
||||
|
||||
fn random_string() -> String {
|
||||
rand::rng()
|
||||
.sample_iter(&Alphanumeric)
|
||||
.take(30)
|
||||
.map(char::from)
|
||||
.collect()
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_phonetic() {
|
||||
let strings: Vec<String> = (0..5_000)
|
||||
.map(|_| random_string())
|
||||
.collect();
|
||||
|
||||
let phonetic = DoubleMetaphone::new(Some(6));
|
||||
|
||||
let now = Instant::now();
|
||||
for value in &strings {
|
||||
let _ = phonetic.encode(value);
|
||||
}
|
||||
|
||||
let elapsed = now.elapsed();
|
||||
println!("Elapsed time: {}.{:03} secs", elapsed.as_secs(), elapsed.subsec_millis());
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,47 @@
|
||||
pub mod playlist;
|
||||
mod xtream;
|
||||
// mod affix;
|
||||
mod xtream_vod;
|
||||
mod xtream_series;
|
||||
pub mod epg;
|
||||
mod sort;
|
||||
pub mod trakt;
|
||||
|
||||
#[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) {
|
||||
match target.get_xtream_output() {
|
||||
Some(xtream_output) => (xtream_output.[<resolve_ $cluster>] && fpl.input.input_type == InputType::Xtream,
|
||||
xtream_output.[<resolve_ $cluster _delay>]),
|
||||
None => (false, 0)
|
||||
}
|
||||
}
|
||||
}
|
||||
};
|
||||
}
|
||||
use create_resolve_options_function_for_xtream_target;
|
||||
|
||||
@@ -0,0 +1,578 @@
|
||||
use crate::model::{ConfigInput, ConfigRename};
|
||||
use crate::utils::epg;
|
||||
use crate::utils::m3u;
|
||||
use crate::utils::xtream;
|
||||
use crate::Config;
|
||||
use std::collections::{HashMap, HashSet};
|
||||
use std::path::PathBuf;
|
||||
use std::sync::Arc;
|
||||
use std::thread;
|
||||
use tokio::sync::Mutex;
|
||||
|
||||
use crate::foundation::filter::{get_field_value, set_field_value, ValueProvider, ValueAccessor};
|
||||
use crate::messaging::{send_message, MsgKind};
|
||||
use crate::model::{ConfigTarget, InputType, ItemField, ProcessTargets, ProcessingOrder};
|
||||
use crate::model::{CounterModifier, Mapping};
|
||||
use crate::model::{FetchedPlaylist, PlaylistGroup, PlaylistItem};
|
||||
use shared::model::{FieldGetAccessor, FieldSetAccessor, PlaylistEntry, UUIDType, XtreamCluster};
|
||||
use crate::model::{InputStats, PlaylistStats, SourceStats, TargetStats};
|
||||
use crate::processing::playlist_watch::process_group_watch;
|
||||
use crate::processing::processor::xtream_series::playlist_resolve_series;
|
||||
use crate::processing::processor::trakt::process_trakt_categories_for_target;
|
||||
use crate::repository::playlist_repository::persist_playlist;
|
||||
use crate::tuliprox_error::{get_errors_notify_message, notify_err, TuliproxError, TuliproxErrorKind};
|
||||
use crate::utils::debug_if_enabled;
|
||||
use crate::utils::default_as_default;
|
||||
use deunicode::deunicode;
|
||||
use log::{debug, error, info, log_enabled, trace, warn, Level};
|
||||
use std::time::Instant;
|
||||
use reqwest::Client;
|
||||
use crate::model::Epg;
|
||||
use crate::processing::parser::xmltv::flatten_tvguide;
|
||||
use crate::processing::processor::epg::process_playlist_epg;
|
||||
use crate::processing::processor::xtream_vod::playlist_resolve_vod;
|
||||
use crate::processing::processor::sort::sort_playlist;
|
||||
use crate::utils::StepMeasure;
|
||||
|
||||
fn is_valid(pli: &PlaylistItem, target: &ConfigTarget) -> bool {
|
||||
let provider = ValueProvider { 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 assign_channel_no_playlist(new_playlist: &mut [PlaylistGroup]) {
|
||||
let assigned_chnos: HashSet<u32> = new_playlist.iter().flat_map(|g| &g.channels)
|
||||
.filter(|c| !c.header.chno.is_empty())
|
||||
.map(|c| c.header.chno.as_str())
|
||||
.flat_map(str::parse::<u32>).collect();
|
||||
let mut chno = 1;
|
||||
for group in new_playlist {
|
||||
for chan in &mut group.channels {
|
||||
if chan.header.chno.is_empty() {
|
||||
while assigned_chnos.contains(&chno) {
|
||||
chno += 1;
|
||||
}
|
||||
chan.header.chno = chno.to_string();
|
||||
chno += 1;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn exec_rename(pli: &mut 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_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, 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 = 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
|
||||
}
|
||||
}
|
||||
|
||||
fn map_channel(mut channel: PlaylistItem, mapping: &Mapping) -> PlaylistItem {
|
||||
if let Some(mapper) = &mapping.mapper {
|
||||
if !mapper.is_empty() {
|
||||
let header = &channel.header;
|
||||
let channel_name = if mapping.match_as_ascii { deunicode(&header.name) } else { header.name.to_string() };
|
||||
if mapping.match_as_ascii && log_enabled!(Level::Trace) { trace!("Decoded {} for matching to {}", &header.name, &channel_name); }
|
||||
let ref_chan = &mut channel;
|
||||
let templates = mapping.templates.as_ref();
|
||||
for m in mapper {
|
||||
if let Some(script) = m.t_script.as_ref() {
|
||||
if let Some(filter) = &m.t_filter {
|
||||
let provider = ValueProvider { pli: ref_chan };
|
||||
if filter.filter(&provider) {
|
||||
let mut accessor = ValueAccessor { pli: ref_chan };
|
||||
script.eval(&mut accessor, templates);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
channel
|
||||
}
|
||||
|
||||
fn map_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option<Vec<PlaylistGroup>> {
|
||||
if let Some(mappings) = target.t_mapping.load().as_ref() {
|
||||
let new_playlist: Vec<PlaylistGroup> = playlist.iter().map(|playlist_group| {
|
||||
let mut grp = playlist_group.clone();
|
||||
mappings.iter().filter(|&mapping| mapping.mapper.as_ref().is_some_and(|v| !v.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.xtream_cluster;
|
||||
let title = &channel.header.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: title.to_string(),
|
||||
channels: vec![channel.clone()],
|
||||
xtream_cluster: *cluster,
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
Some(new_groups)
|
||||
} else {
|
||||
None
|
||||
}
|
||||
}
|
||||
|
||||
fn map_playlist_counter(target: &ConfigTarget, playlist: &mut [PlaylistGroup]) {
|
||||
if target.t_mapping.load().is_some() {
|
||||
let guard = target.t_mapping.load();
|
||||
let mappings = guard.as_ref().unwrap();
|
||||
for mapping in mappings.iter() {
|
||||
if let Some(counter_list) = &mapping.t_counter {
|
||||
for counter in counter_list {
|
||||
for plg in &mut *playlist {
|
||||
for channel in &mut plg.channels {
|
||||
let provider = ValueProvider { pli: channel };
|
||||
if counter.filter.filter(&provider) {
|
||||
let cntval = counter.value.fetch_add(1, core::sync::atomic::Ordering::SeqCst);
|
||||
let padded_cntval = if counter.padding > 0 {
|
||||
format!("{:0width$}", cntval, width = counter.padding as usize)
|
||||
} else {
|
||||
cntval.to_string()
|
||||
};
|
||||
let new_value = if counter.modifier == CounterModifier::Assign {
|
||||
padded_cntval
|
||||
} else {
|
||||
let value = channel.header.get_field(&counter.field).map_or_else(String::new, |field_value| field_value.to_string());
|
||||
if counter.modifier == CounterModifier::Suffix {
|
||||
format!("{value}{}{padded_cntval}", counter.concat)
|
||||
} else {
|
||||
format!("{padded_cntval}{}{value}", counter.concat)
|
||||
}
|
||||
};
|
||||
channel.header.set_field(&counter.field, new_value.as_str());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// 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(input: &ConfigInput, user_targets: &ProcessTargets) -> bool {
|
||||
let input_enabled = input.enabled;
|
||||
let input_id = input.id;
|
||||
(!user_targets.enabled && input_enabled) || user_targets.has_input(input_id)
|
||||
}
|
||||
|
||||
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<TuliproxError>) {
|
||||
let source = cfg.sources.get_source_at(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);
|
||||
// Download the sources
|
||||
let mut source_downloaded = false;
|
||||
for input in &source.inputs {
|
||||
if is_input_enabled(input, &user_targets) {
|
||||
source_downloaded = true;
|
||||
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(&cfg, Arc::clone(&client), input, &cfg.working_dir).await,
|
||||
InputType::M3uBatch | InputType::XtreamBatch => (vec![], vec![])
|
||||
};
|
||||
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, input_name, elapsed));
|
||||
}
|
||||
}
|
||||
if source_downloaded {
|
||||
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<TuliproxError>) {
|
||||
let mut handle_list = vec![];
|
||||
let thread_num = config.threads;
|
||||
let process_parallel = thread_num > 1 && config.sources.sources.len() > 1;
|
||||
if process_parallel && log_enabled!(Level::Debug) {
|
||||
debug!("Using {thread_num} threads");
|
||||
}
|
||||
let errors = Arc::new(Mutex::<Vec<TuliproxError>>::new(vec![]));
|
||||
let stats = Arc::new(Mutex::<Vec<SourceStats>>::new(vec![]));
|
||||
for (index, _) in config.sources.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 || {
|
||||
// TODO better way ?
|
||||
let rt = tokio::runtime::Runtime::new().unwrap();
|
||||
rt.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<(String, XtreamCluster), usize> = HashMap::new();
|
||||
for group in playlistgroups {
|
||||
let key = (group.title.to_string(), 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<TuliproxError>) -> Result<(), Vec<TuliproxError>> {
|
||||
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![];
|
||||
|
||||
debug!("Executing processing pipes");
|
||||
|
||||
let mut step = StepMeasure::new("Pipes processed");
|
||||
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, &mut 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);
|
||||
}
|
||||
|
||||
step.tick("Processed epg");
|
||||
let (new_epg, mut new_playlist) = process_epg(&mut processed_fetched_playlists);
|
||||
|
||||
if new_playlist.is_empty() {
|
||||
info!("Playlist is empty: {}", &target.name);
|
||||
Ok(())
|
||||
} else {
|
||||
|
||||
// Process Trakt categories
|
||||
step.tick("Processing Trakt categories");
|
||||
trakt_playlist(&client, target, errors, &mut new_playlist).await;
|
||||
|
||||
step.tick("Merged playlists");
|
||||
let mut flat_new_playlist = flatten_groups(new_playlist);
|
||||
|
||||
step.tick("Sorted playlists");
|
||||
sort_playlist(target, &mut flat_new_playlist);
|
||||
step.tick("Assigned channel number");
|
||||
assign_channel_no_playlist(&mut flat_new_playlist);
|
||||
step.tick("Assigned channel counter");
|
||||
map_playlist_counter(target, &mut flat_new_playlist);
|
||||
|
||||
step.tick("Processed group watches");
|
||||
process_watch(&client, target, cfg, &flat_new_playlist);
|
||||
step.tick("Persisting playlists");
|
||||
let result = persist_playlist(&mut flat_new_playlist, flatten_tvguide(&new_epg).as_ref(), target, cfg).await;
|
||||
step.stop();
|
||||
result
|
||||
}
|
||||
}
|
||||
|
||||
async fn trakt_playlist(client: &Arc<Client>, target: &ConfigTarget, errors: &mut Vec<TuliproxError>, playlist: &mut Vec<PlaylistGroup>) {
|
||||
match process_trakt_categories_for_target(Arc::clone(client), playlist, target).await {
|
||||
Ok(trakt_categories) => {
|
||||
if !trakt_categories.is_empty() {
|
||||
info!("Adding {} Trakt categories to playlist", trakt_categories.len());
|
||||
playlist.extend(trakt_categories);
|
||||
}
|
||||
}
|
||||
Err(trakt_errors) => {
|
||||
warn!("Trakt processing failed with {} errors", trakt_errors.len());
|
||||
errors.extend(trakt_errors);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn process_epg(processed_fetched_playlists: &mut Vec<FetchedPlaylist>) -> (Vec<Epg>, Vec<PlaylistGroup>) {
|
||||
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 fp in processed_fetched_playlists {
|
||||
process_playlist_epg(fp, &mut new_epg);
|
||||
new_playlist.append(&mut fp.playlistgroups);
|
||||
}
|
||||
(new_epg, new_playlist)
|
||||
}
|
||||
|
||||
fn process_watch(client: &Arc<reqwest::Client>, 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(client, 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(Arc::clone(&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(&client, &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(&client, &MsgKind::Error, cfg.messaging.as_ref(), error_msg.as_str());
|
||||
}
|
||||
}
|
||||
let elapsed = start_time.elapsed().as_secs();
|
||||
info!("🌷 Update process finished! Took {elapsed} secs.");
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
// #[test]
|
||||
// fn test_jaro_winkeler() {
|
||||
// let data = [("yessport5", "heyessport5gold"), ("yessport5", "heyesport5gold")];
|
||||
//
|
||||
// data.iter().for_each(|(first, second)|
|
||||
// println!("jaro_winkler {} = {} => {}", first, second, strsim::jaro_winkler(first, second)));
|
||||
// // println!("jaro {}", strsim::jaro(data.0, data.1));
|
||||
// // println!("levenhstein {}", strsim::levenshtein(data.0, data.1));
|
||||
// // println!("damerau_levenshtein {:?}", strsim::damerau_levenshtein(data.0, data.1));
|
||||
// // println!("osa distance {:?}", strsim::osa_distance(data.0, data.1));
|
||||
// // println!("sorensen dice {:?}", strsim::sorensen_dice(data.0, data.1));
|
||||
// }
|
||||
|
||||
}
|
||||
@@ -0,0 +1,223 @@
|
||||
use crate::foundation::filter::get_field_value;
|
||||
use crate::model::{ConfigSortChannel, ConfigSortGroup, ConfigTarget, SortOrder};
|
||||
use crate::model::{PlaylistGroup, PlaylistItem};
|
||||
use deunicode::deunicode;
|
||||
use std::cmp::Ordering;
|
||||
|
||||
fn playlist_comparator(
|
||||
sequence: Option<&Vec<regex::Regex>>,
|
||||
order: SortOrder,
|
||||
value_a: &str,
|
||||
value_b: &str,
|
||||
) -> Ordering {
|
||||
if let Some(regex_list) = sequence {
|
||||
let mut match_a = None;
|
||||
let mut match_b = None;
|
||||
|
||||
for (i, regex) in regex_list.iter().enumerate() {
|
||||
if match_a.is_none() {
|
||||
if let Some(caps) = regex.captures(value_a) {
|
||||
match_a = Some((i, caps));
|
||||
}
|
||||
}
|
||||
if match_b.is_none() {
|
||||
if let Some(caps) = regex.captures(value_b) {
|
||||
match_b = Some((i, caps));
|
||||
}
|
||||
}
|
||||
|
||||
// If both matches found → break
|
||||
if match_a.is_some() && match_b.is_some() {
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
match (match_a, match_b) {
|
||||
(Some((idx_a, caps_a)), Some((idx_b, caps_b))) => {
|
||||
// Different regex indices → sort by their sequence order.
|
||||
if idx_a != idx_b {
|
||||
return match order {
|
||||
SortOrder::Asc => idx_a.cmp(&idx_b),
|
||||
SortOrder::Desc => idx_b.cmp(&idx_a),
|
||||
};
|
||||
}
|
||||
|
||||
// Same regex → sort by captures (c1, c2, …)
|
||||
let mut named: Vec<_> = regex_list[idx_a]
|
||||
.capture_names()
|
||||
.flatten()
|
||||
.filter(|name| name.starts_with('c'))
|
||||
.collect();
|
||||
|
||||
named.sort_by_key(|name| name[1..].parse::<u32>().unwrap_or(0));
|
||||
|
||||
for name in named {
|
||||
let va = caps_a.name(name).map(|m| m.as_str());
|
||||
let vb = caps_b.name(name).map(|m| m.as_str());
|
||||
if let (Some(va), Some(vb)) = (va, vb) {
|
||||
let o = va.cmp(vb);
|
||||
if o != Ordering::Equal {
|
||||
return match order {
|
||||
SortOrder::Asc => o,
|
||||
SortOrder::Desc => o.reverse(),
|
||||
};
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Ordering::Equal
|
||||
}
|
||||
(Some(_), None) => match order {
|
||||
SortOrder::Asc => Ordering::Less,
|
||||
SortOrder::Desc => Ordering::Greater,
|
||||
},
|
||||
(None, Some(_)) => match order {
|
||||
SortOrder::Asc => Ordering::Greater,
|
||||
SortOrder::Desc => Ordering::Less,
|
||||
},
|
||||
(None, None) => {
|
||||
// NP match → fallback
|
||||
let o = value_a.cmp(value_b);
|
||||
match order {
|
||||
SortOrder::Asc => o,
|
||||
SortOrder::Desc => o.reverse(),
|
||||
}
|
||||
}
|
||||
}
|
||||
} else {
|
||||
// No Regex-Sequence defined → fallback
|
||||
let o = value_a.cmp(value_b);
|
||||
match order {
|
||||
SortOrder::Asc => o,
|
||||
SortOrder::Desc => o.reverse(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn playlistgroup_comparator(a: &PlaylistGroup, b: &PlaylistGroup, group_sort: &ConfigSortGroup, match_as_ascii: bool) -> Ordering {
|
||||
let value_a = if match_as_ascii { deunicode(&a.title) } else { a.title.to_string() };
|
||||
let value_b = if match_as_ascii { deunicode(&b.title) } else { b.title.to_string() };
|
||||
|
||||
playlist_comparator(group_sort.t_re_sequence.as_ref(), group_sort.order, &value_a, &value_b)
|
||||
}
|
||||
|
||||
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 { deunicode(&raw_value_a) } else { raw_value_a };
|
||||
let value_b = if match_as_ascii { deunicode(&raw_value_b) } else { raw_value_b };
|
||||
|
||||
playlist_comparator(channel_sort.t_re_sequence.as_ref(), channel_sort.order, &value_a, &value_b)
|
||||
}
|
||||
|
||||
pub(in crate::processing::processor) 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.t_re_group_pattern.as_ref().unwrap();
|
||||
for group in new_playlist.iter_mut() {
|
||||
let group_title = if match_as_ascii { deunicode(&group.title) } else { group.title.to_string() };
|
||||
if regexp.is_match(group_title.as_str()) {
|
||||
group.channels.sort_by(|chan1, chan2| playlistitem_comparator(chan1, chan2, channel_sort, match_as_ascii));
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use crate::model::{ConfigSortChannel, ItemField, SortOrder};
|
||||
use crate::model::{PlaylistItem, PlaylistItemHeader};
|
||||
use crate::processing::processor::sort::playlistitem_comparator;
|
||||
use regex::Regex;
|
||||
|
||||
#[test]
|
||||
fn test_sort() {
|
||||
let mut channels: Vec<PlaylistItem> = vec![
|
||||
("D", "HD"), ("A", "FHD"), ("Z", "HD"), ("K", "HD"), ("B", "HD"), ("A", "HD"),
|
||||
("K", "UHD"), ("C", "HD"), ("L", "FHD"), ("R", "UHD"), ("T", "SD"), ("A", "FHD"),
|
||||
].into_iter().map(|(name, quality)| PlaylistItem { header: PlaylistItemHeader { title: format!("Chanel {name} [{quality}]"), ..Default::default() } }).collect::<Vec<PlaylistItem>>();
|
||||
|
||||
let channel_sort = ConfigSortChannel {
|
||||
field: ItemField::Caption,
|
||||
group_pattern: ".*".to_string(),
|
||||
order: SortOrder::Asc,
|
||||
sequence: None,
|
||||
t_re_sequence: Some(vec![
|
||||
Regex::new(r"(?P<c1>.*?)\bUHD\b").unwrap(),
|
||||
Regex::new(r"(?P<c1>.*?)\bFHD\b").unwrap(),
|
||||
Regex::new(r"(?P<c1>.*?)\bHD\b").unwrap(),
|
||||
]),
|
||||
t_re_group_pattern: Some(Regex::new(".*").unwrap()),
|
||||
};
|
||||
|
||||
channels.sort_by(|chan1, chan2| playlistitem_comparator(chan1, chan2, &channel_sort, true));
|
||||
let expected = vec!["Chanel K [UHD]", "Chanel R [UHD]", "Chanel A [FHD]", "Chanel A [FHD]", "Chanel L [FHD]", "Chanel A [HD]", "Chanel B [HD]", "Chanel C [HD]", "Chanel D [HD]", "Chanel K [HD]", "Chanel Z [HD]", "Chanel T [SD]"];
|
||||
let sorted = channels.into_iter().map(|pli| pli.header.title.clone()).collect::<Vec<String>>();
|
||||
assert_eq!(expected, sorted);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_sort2() {
|
||||
let mut channels: Vec<PlaylistItem> = vec![
|
||||
"US| EAST [FHD] abc",
|
||||
"US| EAST [FHD] def",
|
||||
"US| EAST [FHD] ghi",
|
||||
"US| EAST [HD] jkl",
|
||||
"US| EAST [HD] mno",
|
||||
"US| EAST [HD] pqrs",
|
||||
"US| EAST [HD] tuv",
|
||||
"US| EAST [HD] wxy",
|
||||
"US| EAST [HD] z",
|
||||
"US| EAST [SD] a",
|
||||
"US| EAST [FHD] bc",
|
||||
"US| EAST [FHD] de",
|
||||
"US| EAST [HD] f",
|
||||
"US| EAST [HD] h",
|
||||
"US| EAST [SD] ijk",
|
||||
"US| EAST [SD] l",
|
||||
"US| EAST [UHD] m",
|
||||
"US| WEST [FHD] no",
|
||||
"US| WEST [HD] qrst",
|
||||
"US| WEST [HD] uvw",
|
||||
"US| (West) xv",
|
||||
"US| East d",
|
||||
"US| West e",
|
||||
"US| West f",
|
||||
].into_iter().map(|name| PlaylistItem { header: PlaylistItemHeader { title: name.to_string(), ..Default::default() } }).collect::<Vec<PlaylistItem>>();
|
||||
|
||||
let channel_sort = ConfigSortChannel {
|
||||
field: ItemField::Caption,
|
||||
group_pattern: ".*US.*".to_string(),
|
||||
order: SortOrder::Asc,
|
||||
sequence: None,
|
||||
t_re_sequence: Some(vec").unwrap(),
|
||||
Regex::new(r"^US\| EAST.*?\[\bFHD\b\](?P<c1>.*)").unwrap(),
|
||||
Regex::new(r"^US\| EAST.*?\[\bHD\b\](?P<c1>.*)").unwrap(),
|
||||
Regex::new(r"^US\| EAST.*?\[\bSD\b\](?P<c1>.*)").unwrap(),
|
||||
Regex::new(r"^US\| WEST.*?\[\bUHD\b\](?P<c1>.*)").unwrap(),
|
||||
Regex::new(r"^US\| WEST.*?\[\bFHD\b\](?P<c1>.*)").unwrap(),
|
||||
Regex::new(r"^US\| WEST.*?\[\bHD\b\](?P<c1>.*)").unwrap(),
|
||||
Regex::new(r"^US\| WEST.*?\[\bSD\b\](?P<c1>.*)").unwrap(),
|
||||
]),
|
||||
t_re_group_pattern: Some(Regex::new(".*").unwrap()),
|
||||
};
|
||||
|
||||
channels.sort_by(|chan1, chan2| playlistitem_comparator(chan1, chan2, &channel_sort, true));
|
||||
let sorted = channels.into_iter().map(|pli| pli.header.title.clone()).collect::<Vec<String>>();
|
||||
let expected = vec!["US| EAST [UHD] m", "US| EAST [FHD] abc", "US| EAST [FHD] bc", "US| EAST [FHD] de", "US| EAST [FHD] def", "US| EAST [FHD] ghi", "US| EAST [HD] f", "US| EAST [HD] h", "US| EAST [HD] jkl", "US| EAST [HD] mno", "US| EAST [HD] pqrs", "US| EAST [HD] tuv", "US| EAST [HD] wxy", "US| EAST [HD] z", "US| EAST [SD] a", "US| EAST [SD] ijk", "US| EAST [SD] l", "US| WEST [FHD] no", "US| WEST [HD] qrst", "US| WEST [HD] uvw", "US| (West) xv", "US| East d", "US| West e", "US| West f"];
|
||||
assert_eq!(expected, sorted);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,321 @@
|
||||
use std::borrow::Cow;
|
||||
use crate::model::{ConfigTarget, PlaylistGroup, PlaylistItem};
|
||||
use shared::model::{FieldGetAccessor, FieldSetAccessor, XtreamCluster};
|
||||
use crate::model::{TraktConfig, TraktContentType, TraktListConfig, TraktListItem, TraktMatchItem, TraktMatchResult};
|
||||
use crate::tuliprox_error::TuliproxError;
|
||||
use crate::utils::trakt::client::TraktClient;
|
||||
use crate::utils::trakt::extract_year_from_title;
|
||||
use crate::utils::trakt::normalize_title_for_matching;
|
||||
use crate::utils::{get_u32_from_serde_value, 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> {
|
||||
if let Some(caps) = CONSTANTS.re_quality.captures(value) {
|
||||
if let Some(val) = caps.get(0) {
|
||||
return Some(val.as_str());
|
||||
}
|
||||
}
|
||||
None
|
||||
}
|
||||
|
||||
|
||||
/// Utility functions for content type compatibility
|
||||
fn should_include_item(item: &TraktListItem, content_type: &TraktContentType) -> bool {
|
||||
match content_type {
|
||||
TraktContentType::Vod => item.content_type == TraktContentType::Vod,
|
||||
TraktContentType::Series => item.content_type == TraktContentType::Series,
|
||||
TraktContentType::Both => true,
|
||||
}
|
||||
}
|
||||
|
||||
fn is_compatible_content_type(cluster: XtreamCluster, content_type: &TraktContentType) -> bool {
|
||||
match content_type {
|
||||
TraktContentType::Vod => cluster == XtreamCluster::Video,
|
||||
TraktContentType::Series => cluster == XtreamCluster::Series,
|
||||
TraktContentType::Both => matches!(cluster, XtreamCluster::Video | XtreamCluster::Series),
|
||||
}
|
||||
}
|
||||
|
||||
/// Extract TMDB ID from playlist item
|
||||
fn extract_tmdb_id_from_playlist_item(item: &PlaylistItem) -> Option<u32> {
|
||||
if let Some(additional_props) = &item.header.additional_properties {
|
||||
if let Some(props_str) = additional_props.as_str() {
|
||||
if let Ok(props) = serde_json::from_str::<serde_json::Map<String, serde_json::Value>>(props_str) {
|
||||
if let Some(tmdb_value) = props.get("tmdb") {
|
||||
return get_u32_from_serde_value(tmdb_value);
|
||||
}
|
||||
if let Some(tmdb_id_value) = props.get("tmdb_id") {
|
||||
return get_u32_from_serde_value(tmdb_id_value);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
None
|
||||
}
|
||||
|
||||
fn calculate_year_bonus(playlist_year: Option<u32>, trakt_year: Option<u32>) -> f64 {
|
||||
if let (Some(p_year), Some(t_year)) = (playlist_year, trakt_year) {
|
||||
if p_year == t_year {
|
||||
// Perfect year match gets substantial bonus
|
||||
return 0.5;
|
||||
}
|
||||
return -0.5;
|
||||
}
|
||||
0.0
|
||||
}
|
||||
|
||||
fn find_best_fuzzy_match_for_item<'a>(channel: (&'a PlaylistItem, String, Option<u32>, Option<u32>), trakt_items: &'a [TraktMatchItem], list_config: &'a TraktListConfig) -> Option<TraktMatchResult<'a>> {
|
||||
// Try fuzzy matching if no exact match found
|
||||
let normalized_playlist_title = channel.1;
|
||||
let playlist_year = channel.2;
|
||||
let threshold = f64::from(list_config.fuzzy_match_threshold) / 100.0;
|
||||
let mut best_match: Option<(&TraktMatchItem, f64)> = None;
|
||||
|
||||
for trakt_item in trakt_items {
|
||||
let title_score = normalized_levenshtein(&normalized_playlist_title, &trakt_item.normalized_title);
|
||||
|
||||
if title_score >= threshold {
|
||||
// Calculate year bonus
|
||||
let year_bonus = calculate_year_bonus(playlist_year, trakt_item.year);
|
||||
let mut combined_score = title_score + year_bonus;
|
||||
|
||||
// Clamp score to [0.0, 1.0]
|
||||
combined_score = combined_score.clamp(0.0, 1.0);
|
||||
|
||||
// Check if this is the best match so far and meets threshold
|
||||
if combined_score >= threshold {
|
||||
if let Some((_, current_best_score)) = &best_match {
|
||||
if combined_score > *current_best_score {
|
||||
best_match = Some((trakt_item, combined_score));
|
||||
}
|
||||
} else {
|
||||
best_match = Some((trakt_item, combined_score));
|
||||
}
|
||||
// early exit strategy
|
||||
if combined_score >= 0.99 {
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if let Some((trakt_item, combined_score)) = best_match {
|
||||
// let match_type = if playlist_year.is_some() && trakt_item.year.is_some() {
|
||||
// MatchType::FuzzyTitleYear
|
||||
// } else {
|
||||
// MatchType::FuzzyTitle
|
||||
// };
|
||||
|
||||
trace_if_enabled!("Fuzzy match: '{}' -> '{}' (final: {combined_score:.3}" /*, type: {match_type:?})"*/, channel.0.header.title, trakt_item.title);
|
||||
|
||||
return Some(TraktMatchResult {
|
||||
playlist_item: channel.0,
|
||||
trakt_item,
|
||||
match_score: combined_score,
|
||||
// match_type: match_type.clone(),
|
||||
});
|
||||
}
|
||||
|
||||
None
|
||||
}
|
||||
|
||||
fn find_best_match_for_item<'a>(
|
||||
channel: (&'a PlaylistItem, String, Option<u32>, Option<u32>),
|
||||
trakt_items: &'a [TraktMatchItem<'a>],
|
||||
list_config: &'a TraktListConfig,
|
||||
) -> Option<TraktMatchResult<'a>> {
|
||||
// Try TMDB exact matching first
|
||||
if let Some(playlist_tmdb_id) = channel.3 {
|
||||
for trakt_item in trakt_items {
|
||||
if Some(playlist_tmdb_id) == trakt_item.tmdb_id {
|
||||
trace!("TMDB exact match: '{}' (TMDB: {})", channel.0.header.title, playlist_tmdb_id);
|
||||
return Some(TraktMatchResult {
|
||||
playlist_item: channel.0,
|
||||
trakt_item,
|
||||
match_score: 1.0,
|
||||
// match_type: MatchType::TmdbExact,
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
find_best_fuzzy_match_for_item(channel, trakt_items, list_config)
|
||||
}
|
||||
|
||||
fn create_category_from_matches<'a>(
|
||||
matches: Vec<TraktMatchResult<'a>>,
|
||||
list_config: &'a TraktListConfig,
|
||||
) -> Option<PlaylistGroup> {
|
||||
if matches.is_empty() { return None; }
|
||||
|
||||
let mut matched_items = Vec::new();
|
||||
|
||||
let mut sorted_matches = matches;
|
||||
sorted_matches.sort_by(|a, b| {
|
||||
(
|
||||
a.trakt_item.rank.unwrap_or(9999),
|
||||
a.trakt_item.title.to_lowercase(),
|
||||
).cmp(&(
|
||||
b.trakt_item.rank.unwrap_or(9999),
|
||||
b.trakt_item.title.to_lowercase(),
|
||||
))
|
||||
});
|
||||
|
||||
let group_title = &list_config.category_name;
|
||||
|
||||
for match_result in sorted_matches {
|
||||
let mut modified_item = match_result.playlist_item.clone();
|
||||
// Use the (possibly numbered) title from the match result (which now contains the original playlist title)
|
||||
with!(mut modified_item.header => header {
|
||||
// Synchronize name with title so both fields show the same value
|
||||
// header.title.clone_from(&match_result.trakt_item.title.to_string());
|
||||
// header.name.clone_from(&match_result.trakt_item.title.to_string());
|
||||
let title = header.get_field("caption").unwrap_or_else(|| Cow::Borrowed(&header.title));
|
||||
if extract_quality(&title).is_none() {
|
||||
if let Some(quality) = extract_quality(&header.group) {
|
||||
let mut caption = String::with_capacity(title.len() + 6);
|
||||
caption.push('[');
|
||||
caption.push_str(quality);
|
||||
caption.push_str("] ");
|
||||
caption.push_str(&title);
|
||||
header.set_field("caption", &caption);
|
||||
}
|
||||
}
|
||||
header.group = String::from(group_title);
|
||||
header.gen_uuid();
|
||||
});
|
||||
matched_items.push(modified_item);
|
||||
}
|
||||
|
||||
if matched_items.is_empty() { return None; }
|
||||
|
||||
|
||||
let cluster = match list_config.content_type {
|
||||
TraktContentType::Vod => XtreamCluster::Video,
|
||||
TraktContentType::Series => XtreamCluster::Series,
|
||||
TraktContentType::Both => {
|
||||
matched_items.first()
|
||||
.map_or(XtreamCluster::Video, |item| item.header.xtream_cluster)
|
||||
}
|
||||
};
|
||||
|
||||
Some(PlaylistGroup {
|
||||
id: 0,
|
||||
title: String::from(group_title),
|
||||
channels: matched_items,
|
||||
xtream_cluster: cluster,
|
||||
})
|
||||
}
|
||||
|
||||
fn match_trakt_items_with_playlist<'a>(
|
||||
trakt_items: &'a [TraktListItem],
|
||||
playlist: &'a [PlaylistGroup],
|
||||
list_config: &'a TraktListConfig,
|
||||
) -> Option<PlaylistGroup> {
|
||||
let trakt_match_items: Vec<TraktMatchItem<'a>> = trakt_items
|
||||
.iter()
|
||||
.filter(|item| should_include_item(item, &list_config.content_type))
|
||||
.filter_map(TraktMatchItem::from_trakt_list_item)
|
||||
.collect();
|
||||
|
||||
debug!("Matching {} Trakt items against playlist for content type {:?}", trakt_match_items.len(), list_config.content_type);
|
||||
|
||||
let mut matches = Vec::new();
|
||||
for playlist_group in playlist {
|
||||
for channel in &playlist_group.channels {
|
||||
if is_compatible_content_type(channel.header.xtream_cluster, &list_config.content_type) {
|
||||
let normalized_title = normalize_title_for_matching(&channel.header.title);
|
||||
let channel_year = extract_year_from_title(&channel.header.title);
|
||||
let channel_tmdb_id = extract_tmdb_id_from_playlist_item(channel);
|
||||
if let Some(matched) = find_best_match_for_item((channel, normalized_title, channel_year, channel_tmdb_id), &trakt_match_items, list_config) {
|
||||
matches.push(matched);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
create_category_from_matches(matches, list_config)
|
||||
}
|
||||
|
||||
pub struct TraktCategoriesProcessor {
|
||||
client: TraktClient,
|
||||
}
|
||||
|
||||
impl TraktCategoriesProcessor {
|
||||
pub fn new(http_client: Arc<reqwest::Client>, trakt_config: &TraktConfig) -> Self {
|
||||
let client = TraktClient::new(http_client, trakt_config.api.clone());
|
||||
Self { client }
|
||||
}
|
||||
|
||||
pub async fn process_trakt_categories(
|
||||
&self,
|
||||
playlist: &[PlaylistGroup],
|
||||
target: &ConfigTarget,
|
||||
trakt_config: &TraktConfig,
|
||||
) -> Result<Vec<PlaylistGroup>, Vec<TuliproxError>> {
|
||||
if trakt_config.lists.is_empty() {
|
||||
debug!("No Trakt lists configured for target {}", target.name);
|
||||
return Ok(vec![]);
|
||||
}
|
||||
|
||||
info!("Processing {} Trakt lists for target {}", trakt_config.lists.len(), target.name);
|
||||
let mut new_categories = Vec::new();
|
||||
let mut total_matches = 0;
|
||||
for list_config in &trakt_config.lists {
|
||||
let cache_key = format!("{}:{}", list_config.user, list_config.list_slug);
|
||||
|
||||
match self.client.get_list_items(list_config).await {
|
||||
Ok(trakt_items) => {
|
||||
debug!("Processing Trakt list {cache_key} with {} items", trakt_items.len());
|
||||
|
||||
if let Some(category) = match_trakt_items_with_playlist(&trakt_items, playlist, list_config) {
|
||||
if !category.channels.is_empty() {
|
||||
total_matches += category.channels.len();
|
||||
let category_len = category.channels.len();
|
||||
new_categories.push(category);
|
||||
debug!("Created Trakt category '{}' with {category_len} items", list_config.category_name);
|
||||
}
|
||||
}
|
||||
}
|
||||
Err(err) => {
|
||||
warn!("Failed to fetch Trakt list {cache_key}: {}", err.message);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
info!("Trakt processing complete: created {} categories with {total_matches} total matches",
|
||||
new_categories.len());
|
||||
|
||||
Ok(new_categories)
|
||||
}
|
||||
}
|
||||
pub async fn process_trakt_categories_for_target(
|
||||
http_client: Arc<reqwest::Client>,
|
||||
playlist: &[PlaylistGroup],
|
||||
target: &ConfigTarget,
|
||||
) -> Result<Vec<PlaylistGroup>, Vec<TuliproxError>> {
|
||||
let Some(trakt_config) = target.get_xtream_output().and_then(|output| output.trakt.as_ref()) else {
|
||||
debug!("No Trakt configuration found for target {}", target.name);
|
||||
return Ok(vec![]);
|
||||
};
|
||||
|
||||
let processor = TraktCategoriesProcessor::new(http_client, trakt_config);
|
||||
processor.process_trakt_categories(playlist, target, trakt_config).await
|
||||
}
|
||||
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
pub fn test_quality() {
|
||||
let quality = extract_quality("Hello HD UHD 720p");
|
||||
assert_eq!(true, quality.is_some());
|
||||
assert_eq!("UHD", quality.unwrap());
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,119 @@
|
||||
use crate::tuliprox_error::{info_err, notify_err};
|
||||
use crate::tuliprox_error::{str_to_io_error, TuliproxError, TuliproxErrorKind};
|
||||
use crate::model::{Config, ConfigInput};
|
||||
use crate::model::{FetchedPlaylist, PlaylistItem};
|
||||
use shared::model::{PlaylistEntry, PlaylistItemType, XtreamCluster};
|
||||
use crate::model::normalize_release_date;
|
||||
use crate::repository::storage::get_input_storage_path;
|
||||
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;
|
||||
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> {
|
||||
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 {
|
||||
tokio::time::sleep(std::time::Duration::new(u64::from(resolve_delay), 0)).await;
|
||||
}
|
||||
result
|
||||
}
|
||||
|
||||
pub(in crate::processing) fn normalize_json_content(content: String) -> String {
|
||||
match from_str::<Value>(&content) {
|
||||
Ok(mut json_value) => {
|
||||
if let Some(info) = json_value.get_mut("info").and_then(Value::as_object_mut) {
|
||||
normalize_release_date(info);
|
||||
}
|
||||
to_string(&json_value).unwrap_or(content)
|
||||
},
|
||||
Err(_) => content,
|
||||
}
|
||||
}
|
||||
|
||||
pub(in crate::processing) fn create_resolve_episode_wal_files(cfg: &Config, input: &ConfigInput) -> Option<(File, PathBuf)> {
|
||||
match get_input_storage_path(&input.name, &cfg.working_dir) {
|
||||
Ok(storage_path) => {
|
||||
let info_path = storage_path.join(format!("{}.{}", crate::model::XC_FILE_SERIES_EPISODE_RECORD, storage_const::FILE_SUFFIX_WAL));
|
||||
let info_file = utils::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.name, &cfg.working_dir) {
|
||||
Ok(storage_path) => {
|
||||
if let Some(file_prefix) = match cluster {
|
||||
XtreamCluster::Live => None,
|
||||
XtreamCluster::Video => Some(crate::model::XC_FILE_VOD_INFO),
|
||||
XtreamCluster::Series => Some(crate::model::XC_FILE_SERIES_INFO)
|
||||
} {
|
||||
let content_path = storage_path.join(format!("{file_prefix}_content.{}", storage_const::FILE_SUFFIX_WAL));
|
||||
let info_path = storage_path.join(format!("{file_prefix}_record.{}", storage_const::FILE_SUFFIX_WAL));
|
||||
let content_file = utils::append_or_crate_file(&content_path).ok()?;
|
||||
let info_file = utils::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: &mut PlaylistItem, processed_provider_ids: &HashMap<u32, u64>, field: &str) -> (bool, u32, u64) {
|
||||
let Some(provider_id) = pli.header.get_provider_id() else { return (false, 0, 0) };
|
||||
let last_modified = pli.header.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<TuliproxError>, 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 fpl_name = &fpl.input.name;
|
||||
let file_path = match get_input_storage_path(fpl_name, &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) => {
|
||||
errors.push(notify_err!(format!("Could not create storage path for input {fpl_name}: {err}")));
|
||||
return processed_info_ids;
|
||||
}
|
||||
};
|
||||
|
||||
{
|
||||
let file_lock = cfg.file_locks.read_lock(&file_path);
|
||||
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);
|
||||
}
|
||||
processed_info_ids
|
||||
}
|
||||
@@ -0,0 +1,237 @@
|
||||
use crate::tuliprox_error::{TuliproxError, TuliproxErrorKind};
|
||||
use crate::model::{Config, ConfigTarget, InputType};
|
||||
use crate::model::{FetchedPlaylist, PlaylistGroup, PlaylistItem};
|
||||
use shared::model::{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};
|
||||
use crate::repository::storage::get_input_storage_path;
|
||||
use crate::repository::xtream_repository::{write_series_info_to_wal_file, 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::tuliprox_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::{XtreamSeriesEpisode, XtreamSeriesInfoEpisode};
|
||||
use crate::utils;
|
||||
use crate::utils::bincode_serialize;
|
||||
use crate::processing::processor::xtream::normalize_json_content;
|
||||
|
||||
create_resolve_options_function_for_xtream_target!(series);
|
||||
|
||||
async fn read_processed_series_info_ids(cfg: &Config, errors: &mut Vec<TuliproxError>, fpl: &FetchedPlaylist<'_>) -> HashMap<u32, u64> {
|
||||
read_processed_info_ids(cfg, errors, fpl, PlaylistItemType::SeriesInfo, |ts: &u64| *ts).await
|
||||
}
|
||||
|
||||
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: &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(client: Arc<reqwest::Client>, cfg: &Config, 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;
|
||||
// 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 = utils::file_writer(&wal_content_file);
|
||||
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 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 {
|
||||
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;
|
||||
}
|
||||
}
|
||||
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}"))));
|
||||
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 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_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));
|
||||
}
|
||||
|
||||
// TODO better approach for transactional updates is multiplexed WAL file.
|
||||
// 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<TuliproxError>,
|
||||
) -> Vec<PlaylistGroup> {
|
||||
let mut result: Vec<PlaylistGroup> = vec![];
|
||||
let input = fpl.input;
|
||||
|
||||
let Ok(Some((info_path, idx_path))) = get_input_storage_path(&input.name, &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 _file_lock = cfg.file_locks.read_lock(&info_path);
|
||||
|
||||
// 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 = utils::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_mut()
|
||||
.filter(|pli| pli.header.item_type == PlaylistItemType::SeriesInfo)
|
||||
{
|
||||
let Some(provider_id) = pli.header.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;
|
||||
(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(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}"))));
|
||||
}
|
||||
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<TuliproxError>,
|
||||
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);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,148 @@
|
||||
use crate::tuliprox_error::{TuliproxError, TuliproxErrorKind};
|
||||
use crate::model::{Config, ConfigTarget, InputType};
|
||||
use crate::model::{FetchedPlaylist, PlaylistItem};
|
||||
use shared::model::{PlaylistItemType, XtreamCluster};
|
||||
use crate::processing::processor::xtream::{create_resolve_info_wal_files, playlist_resolve_download_playlist_item, read_processed_info_ids, should_update_info};
|
||||
use crate::repository::xtream_repository::{write_vod_info_to_wal_file, xtream_update_input_info_file, xtream_update_input_vod_record_from_wal_file, InputVodInfoRecord};
|
||||
use crate::tuliprox_error::{notify_err};
|
||||
use crate::processing::processor::{handle_error, handle_error_and_return, create_resolve_options_function_for_xtream_target};
|
||||
use crate::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;
|
||||
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;
|
||||
|
||||
create_resolve_options_function_for_xtream_target!(vod);
|
||||
|
||||
async fn read_processed_vod_info_ids(cfg: &Config, errors: &mut Vec<TuliproxError>, 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(crate::model::XC_TAG_VOD_INFO_MOVIE_DATA)?.as_object()?;
|
||||
let provider_id = get_u32_from_serde_value(
|
||||
movie_data.get(crate::model::XC_TAG_VOD_INFO_STREAM_ID)?,
|
||||
)?;
|
||||
|
||||
let added = movie_data
|
||||
.get(crate::model::XC_TAG_VOD_INFO_ADDED)
|
||||
.and_then(get_u64_from_serde_value)
|
||||
.unwrap_or(0);
|
||||
|
||||
let info_section = doc.get(crate::model::XC_TAG_VOD_INFO_INFO)?.as_object()?;
|
||||
|
||||
let tmdb_id = info_section
|
||||
.get(crate::model::XC_TAG_VOD_INFO_TMDB_ID)
|
||||
.and_then(get_u32_from_serde_value)
|
||||
.unwrap_or(0);
|
||||
|
||||
let release_date = info_section
|
||||
.get(crate::model::XC_TAG_VOD_INFO_RELEASEDATE)
|
||||
.and_then(get_string_from_serde_value);
|
||||
|
||||
Some((provider_id, InputVodInfoRecord {
|
||||
tmdb_id,
|
||||
ts: added,
|
||||
release_date,
|
||||
}))
|
||||
}
|
||||
|
||||
fn should_update_vod_info(pli: &mut PlaylistItem, processed_provider_ids: &HashMap<u32, u64>) -> (bool, u32, u64) {
|
||||
should_update_info(pli, processed_provider_ids, crate::model::XC_TAG_VOD_INFO_ADDED)
|
||||
}
|
||||
|
||||
pub async fn playlist_resolve_vod(client: Arc<reqwest::Client>, cfg: &Config, target: &ConfigTarget, errors: &mut Vec<TuliproxError>, fpl: &mut 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 = utils::file_writer(&wal_content_file);
|
||||
let mut record_writer = utils::file_writer(&wal_record_file);
|
||||
let mut content_updated = false;
|
||||
|
||||
// TODO merge both filters to one
|
||||
let vod_info_count = fpl.playlistgroups.iter()
|
||||
.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 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 {
|
||||
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;
|
||||
}
|
||||
}
|
||||
}
|
||||
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 {
|
||||
// TODO better approach for transactional updates is multiplexed WAL file.
|
||||
|
||||
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(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));
|
||||
}
|
||||
|
||||
// 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(cfg, fpl.input, provider_id, XtreamCluster::Video) {
|
||||
pli.header.additional_properties = from_str::<Map<String, Value>>(&content).ok().and_then(|info_doc| info_doc.get("info").cloned());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user