smart epg refactoring

This commit is contained in:
euzu
2025-04-23 20:21:36 +02:00
parent 82b53b271a
commit e2990221b7
8 changed files with 284 additions and 205 deletions
+182
View File
@@ -0,0 +1,182 @@
use crate::model::config::{EpgConfig, EpgSmartMatchConfig};
use crate::model::playlist::{FetchedPlaylist, PlaylistItem, XtreamCluster};
use crate::model::xmltv::Epg;
use crate::processing::parser::xmltv::normalize_channel_name;
use log::debug;
use rphonetic::{DoubleMetaphone, Encoder};
use std::borrow::Cow;
use std::collections::{HashMap, HashSet};
pub struct EpgIdCache<'a> {
pub channel_epg_id: HashSet<Cow<'a, str>>,
pub normalized: HashMap<String, Option<String>>,
pub processed: HashSet<String>,
pub smart_match_config: EpgSmartMatchConfig,
pub metaphone: DoubleMetaphone,
pub smart_match_enabled: bool,
pub fuzzy_match_enabled: bool,
}
impl EpgIdCache<'_> {
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(),
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()
}
fn normalize_and_store(&mut self, name: &str, epg_id: Option<&String>) {
let normalized_name = self.normalize(name);
self.normalized.insert(normalized_name, epg_id.map(|v| v.to_string()));
}
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
}
}
fn assign_channel_epg(new_epg: &mut Vec<Epg>, fp: &mut FetchedPlaylist, id_cache: &mut EpgIdCache) {
if let Some(tv_guide) = &fp.epg {
if let Some(epg) = tv_guide.filter(id_cache) {
// // icon tags
// let icon_tags: HashMap<&String, &XmlTag> = epg.children.iter()
// .filter(|tag| tag.icon.is_some() && tag.get_attribute_value(EPG_ATTRIB_ID).is_some())
// .map(|t| (t.get_attribute_value(EPG_ATTRIB_ID).unwrap(), t)).collect();
let filter_live = |c: &&mut PlaylistItem| c.header.xtream_cluster == XtreamCluster::Live;
// let filter_missing_epg_id = |chan: &mut PlaylistItem| chan.header.epg_channel_id.is_none() || chan.header.logo.is_empty() || chan.header.logo_small.is_empty();
let filter_missing_epg_id = |chan: &&mut PlaylistItem| chan.header.epg_channel_id.is_none();
let assign_values = |chan: &mut PlaylistItem| {
if id_cache.smart_match_enabled {
// 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) {
chan.header.epg_channel_id = epg_id.clone();
}
}
}
// if chan.header.epg_channel_id.is_some() && (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 Some(icon) = icon_tag.icon.as_ref() {
// if chan.header.logo.is_empty() {
// chan.header.logo = (*icon).to_string();
// }
// if chan.header.logo_small.is_empty() {
// chan.header.logo = (*icon).to_string();
// }
// }
// }
// }
};
fp.playlistgroups.iter_mut()
.flat_map(|g| &mut g.channels)
.filter(filter_live)
.filter(filter_missing_epg_id)
.for_each(assign_values);
new_epg.push(epg);
}
}
}
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() { // TODO && not fuzzy matching {
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_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());
}
}
+1
View File
@@ -3,6 +3,7 @@ mod xtream;
mod affix;
mod xtream_vod;
mod xtream_series;
pub mod epg;
#[macro_export]
macro_rules! handle_error {
+23 -155
View File
@@ -1,10 +1,9 @@
use crate::Config;
use crate::model::config::{ConfigInput, ConfigRename, EpgConfig, EpgSmartMatchConfig};
use crate::{Config};
use crate::model::config::{ConfigInput, ConfigRename};
use crate::utils::network::epg;
use crate::utils::network::m3u;
use crate::utils::network::xtream;
use core::cmp::Ordering;
use std::borrow::Cow;
use std::collections::{HashMap, HashSet};
use std::path::PathBuf;
use std::sync::{Arc};
@@ -14,7 +13,6 @@ use std::thread;
use log::{debug, error, info, log_enabled, trace, warn, Level};
use std::time::Instant;
use deunicode::deunicode;
use rphonetic::{Encoder, Metaphone};
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};
@@ -25,14 +23,16 @@ use crate::model::playlist::{FetchedPlaylist, FieldGetAccessor, FieldSetAccessor
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, normalize_channel_name};
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};
use crate::model::xmltv::{Epg, XmlTag, EPG_ATTRIB_ID};
use crate::model::xmltv::{Epg};
use crate::processing::parser::xmltv::flatten_tvguide;
use crate::processing::processor::epg::process_playlist_epg;
use crate::utils::step_measure::StepMeasure;
fn is_valid(pli: &PlaylistItem, target: &ConfigTarget) -> bool {
let provider = ValueProvider { pli };
@@ -510,144 +510,6 @@ fn flatten_groups(playlistgroups: Vec<PlaylistGroup>) -> Vec<PlaylistGroup> {
sort_order
}
type PhoneticCodeAndMaybeEpgId<'a> = (Option<Cow<'a, str>>, Option<Cow<'a, str>>);
pub struct EpgIdCache<'a > {
pub channel: HashSet<Cow<'a, str>>,
pub normalized: HashMap<Cow<'a, str>, PhoneticCodeAndMaybeEpgId<'a>>,
pub normalized_phonetic: HashMap<Cow<'a, str>, Cow<'a, str>>,
pub processed: HashSet<String>,
pub smart_match_config: EpgSmartMatchConfig,
pub metaphone: Metaphone,
pub smart_match_enabled: bool,
pub fuzzy_match_enabled: bool,
}
impl EpgIdCache<'_> {
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: HashSet::new(),
normalized: HashMap::new(),
normalized_phonetic: HashMap::new(),
processed: HashSet::new(),
metaphone: Metaphone::default(),
smart_match_enabled: normalize_config.enabled,
fuzzy_match_enabled: normalize_config.enabled && normalize_config.fuzzy_matching,
smart_match_config: normalize_config,
}
}
pub fn is_empty(&self) -> bool {
self.channel.is_empty() && self.normalized.is_empty()
}
pub fn normalize_and_store(&mut self, name: &str) {
let normalized_name = self.normalize(name);
let normalized_key: Cow<str> = Cow::Owned(normalized_name.clone());
let phonetic_code = if self.smart_match_config.fuzzy_matching {
let code: Cow<str> = Cow::Owned(self.metaphone.encode(&normalized_name));
self.normalized_phonetic.insert(normalized_key.clone(), code.clone());
Some(code)
} else {
None
};
self.normalized.insert(normalized_key, (phonetic_code, None));
}
pub fn normalize(&self, name: &str) -> String {
normalize_channel_name(name, &self.smart_match_config)
}
pub fn phonetic(&self, name: &str) -> Cow<str> {
self.normalized_phonetic.get(&Cow::Owned(name.to_string())).map_or_else(|| Cow::Owned(self.metaphone.encode(name)), std::clone::Clone::clone)
// match self.normalized_phonetic.entry(name.clone()) {
// Entry::Occupied(entry) => {
// entry.get().clone()
// }
// Entry::Vacant(entry) => {
//
// let code: Cow<str> = Cow::Owned(self.metaphone.encode(&name));
// // entry.insert(code.clone());
// code
// }
// }
}
}
fn prepare_epg_id_cache(fp: &mut FetchedPlaylist, id_cache: &mut EpgIdCache) {
let smart_match_enabled = id_cache.smart_match_enabled;
let fuzzy_matching = id_cache.fuzzy_match_enabled;
for channel in fp.playlistgroups.iter().flat_map(|g| &g.channels) {
let epg_id = channel.header.epg_channel_id.as_deref();
// insert epg_id to known channel epg_ids
if let Some(id) = epg_id {
if !id.is_empty() {
id_cache.channel.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 || epg_id.is_none_or(str::is_empty));
if needs_normalization {
let name = &channel.header.name;
id_cache.normalize_and_store(name);
}
}
}
fn assign_channel_epg(new_epg: &mut Vec<Epg>, fp: &mut FetchedPlaylist, id_cache: &mut EpgIdCache) {
if let Some(tv_guide) = &fp.epg {
if let Some(epg) = tv_guide.filter(id_cache) {
// icon tags
let icon_tags : HashMap<&String, &XmlTag> = epg.children.iter()
.filter(|tag| tag.icon.is_some() && tag.get_attribute_value(EPG_ATTRIB_ID).is_some())
.map(|t| (t.get_attribute_value(EPG_ATTRIB_ID).unwrap(), t)).collect();
let filter_live = |c: &&mut PlaylistItem| c.header.xtream_cluster == XtreamCluster::Live;
let filter_missing_epg_id = |chan: &&mut PlaylistItem | chan.header.epg_channel_id.is_none() || chan.header.logo.is_empty() || chan.header.logo_small.is_empty();
let assign_values = |chan: &mut PlaylistItem| {
if id_cache.smart_match_enabled {
// 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
if chan.header.epg_channel_id.is_none() || !id_cache.processed.contains(chan.header.epg_channel_id.as_ref().unwrap()) {
let normalized = id_cache.normalize(&chan.header.name);
if let Some((_, Some(epg_id))) = id_cache.normalized.get(&Cow::Borrowed(normalized.as_str())) {
chan.header.epg_channel_id = Some(epg_id.to_string());
}
}
}
if chan.header.epg_channel_id.is_some() && (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 Some(icon) = icon_tag.icon.as_ref() {
if chan.header.logo.is_empty() {
chan.header.logo = (*icon).to_string();
}
if chan.header.logo_small.is_empty() {
chan.header.logo = (*icon).to_string();
}
}
}
}
};
fp.playlistgroups.iter_mut()
.flat_map(|g| &mut g.channels)
.filter(filter_live)
.filter(filter_missing_epg_id)
.for_each(assign_values);
new_epg.push(epg);
}
}
}
async fn process_playlist_for_target(client: Arc<reqwest::Client>,
playlists: &mut [FetchedPlaylist<'_>],
target: &ConfigTarget,
@@ -659,6 +521,10 @@ async fn process_playlist_for_target(client: Arc<reqwest::Client>,
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;
@@ -674,40 +540,42 @@ async fn process_playlist_for_target(client: Arc<reqwest::Client>,
processed_fetched_playlists.push(processed_fpl);
}
step.tick("Processed affixes");
apply_affixes(&mut processed_fetched_playlists);
step.tick("Processed epg");
let (new_epg, new_playlist) = process_epg(&mut processed_fetched_playlists);
if new_playlist.is_empty() {
info!("Playlist is empty: {}", &target.name);
Ok(())
} else {
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");
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(target, cfg, &flat_new_playlist);
persist_playlist(&mut flat_new_playlist, flatten_tvguide(&new_epg).as_ref(), target, cfg).await
step.tick("Persisting playlists");
let result = persist_playlist(&mut flat_new_playlist, flatten_tvguide(&new_epg).as_ref(), target, cfg).await;
step.tick("");
result
}
}
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 {
// collect all epg_channel ids
let mut id_cache = EpgIdCache::new(fp.input.epg.as_ref());
prepare_epg_id_cache(fp, &mut id_cache);
if id_cache.is_empty() {
debug!("channel ids are empty");
} else {
//debug_if_enabled!("found epg information for {}", &target.name);
assign_channel_epg(&mut new_epg, fp, &mut id_cache);
}
process_playlist_epg(fp, &mut new_epg);
new_playlist.append(&mut fp.playlistgroups);
}
(new_epg, new_playlist)