2025-02-05 20:01:49 +01:00
|
|
|
use crate::Config;
|
2025-04-12 14:14:06 +02:00
|
|
|
use crate::model::config::{ConfigInput, ConfigRename, EpgConfig, EpgSmartMatchConfig};
|
2025-02-05 20:01:49 +01:00
|
|
|
use crate::utils::network::epg;
|
|
|
|
|
use crate::utils::network::m3u;
|
|
|
|
|
use crate::utils::network::xtream;
|
2024-05-05 19:15:54 +02:00
|
|
|
use core::cmp::Ordering;
|
2025-04-10 18:08:11 +02:00
|
|
|
use std::borrow::Cow;
|
2023-10-27 19:03:29 +02:00
|
|
|
use std::collections::{HashMap, HashSet};
|
2025-01-06 09:17:05 +01:00
|
|
|
use std::path::PathBuf;
|
2025-03-11 14:56:10 +01:00
|
|
|
use std::sync::{Arc};
|
|
|
|
|
use tokio::sync::Mutex;
|
2023-02-13 12:19:58 +01:00
|
|
|
use std::thread;
|
2023-10-19 18:56:08 +02:00
|
|
|
|
2025-01-06 09:17:05 +01:00
|
|
|
use log::{debug, error, info, log_enabled, trace, warn, Level};
|
2024-11-03 14:52:56 +01:00
|
|
|
use std::time::Instant;
|
2025-04-10 20:26:22 +02:00
|
|
|
use deunicode::deunicode;
|
2025-04-11 20:25:09 +02:00
|
|
|
use rphonetic::{Encoder, Metaphone};
|
2025-02-05 20:01:49 +01:00
|
|
|
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};
|
2024-11-04 18:46:56 +01:00
|
|
|
use crate::messaging::{send_message, MsgKind};
|
2024-05-10 10:13:50 +02:00
|
|
|
use crate::model::config::{ConfigSortChannel, ConfigSortGroup, ConfigTarget, InputType,
|
2024-11-04 18:46:56 +01:00
|
|
|
ItemField, ProcessTargets, ProcessingOrder, SortOrder::{Asc, Desc}};
|
2024-09-11 19:41:50 +02:00
|
|
|
use crate::model::mapping::{CounterModifier, Mapping, MappingValueProcessor};
|
2025-01-06 10:11:14 +01:00
|
|
|
use crate::model::playlist::{FetchedPlaylist, FieldGetAccessor, FieldSetAccessor, PlaylistEntry, PlaylistGroup, PlaylistItem, UUIDType, XtreamCluster};
|
2025-01-06 09:17:05 +01:00
|
|
|
use crate::model::stats::{InputStats, PlaylistStats, SourceStats, TargetStats};
|
2025-02-05 20:01:49 +01:00
|
|
|
use crate::processing::processor::affix::apply_affixes;
|
2023-10-13 18:13:16 +02:00
|
|
|
use crate::processing::playlist_watch::process_group_watch;
|
2025-04-09 20:51:25 +02:00
|
|
|
use crate::processing::parser::xmltv::{flatten_tvguide, normalize_channel_name};
|
2025-02-05 20:01:49 +01:00
|
|
|
use crate::processing::processor::xtream_series::playlist_resolve_series;
|
|
|
|
|
use crate::processing::processor::xtream_vod::playlist_resolve_vod;
|
2024-05-04 20:30:20 +02:00
|
|
|
use crate::repository::playlist_repository::persist_playlist;
|
2024-05-10 12:01:48 +02:00
|
|
|
use crate::utils::default_utils::default_as_default;
|
2025-02-05 20:01:49 +01:00
|
|
|
use crate::utils::{debug_if_enabled};
|
2024-05-04 20:30:20 +02:00
|
|
|
|
2025-04-12 14:14:06 +02:00
|
|
|
use crate::model::xmltv::{Epg, XmlTag, EPG_ATTRIB_ID};
|
2025-04-03 18:03:59 +03:00
|
|
|
|
2024-05-04 20:30:20 +02:00
|
|
|
fn is_valid(pli: &PlaylistItem, target: &ConfigTarget) -> bool {
|
2025-03-11 14:56:10 +01:00
|
|
|
let provider = ValueProvider { pli };
|
2024-05-04 20:30:20 +02:00
|
|
|
target.filter(&provider)
|
|
|
|
|
}
|
2021-10-15 16:59:20 +02:00
|
|
|
|
2024-05-10 12:01:48 +02:00
|
|
|
#[allow(clippy::unnecessary_wraps)]
|
2023-10-08 19:54:59 +02:00
|
|
|
fn filter_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option<Vec<PlaylistGroup>> {
|
2023-10-28 22:20:57 +02:00
|
|
|
debug!("Filtering {} groups", playlist.len());
|
2025-01-09 13:31:05 +01:00
|
|
|
let mut new_playlist = Vec::with_capacity(128);
|
2024-05-10 12:01:48 +02:00
|
|
|
for pg in playlist.iter_mut() {
|
2024-05-04 20:30:20 +02:00
|
|
|
let channels = pg.channels.iter()
|
|
|
|
|
.filter(|&pli| is_valid(pli, target)).cloned().collect::<Vec<PlaylistItem>>();
|
2024-11-03 14:52:56 +01:00
|
|
|
trace!("Filtered group {} has now {}/{} items", pg.title, channels.len(), pg.channels.len());
|
2023-05-03 17:46:39 +02:00
|
|
|
if !channels.is_empty() {
|
2023-01-13 15:56:03 +01:00
|
|
|
new_playlist.push(PlaylistGroup {
|
2023-10-06 19:24:06 +02:00
|
|
|
id: pg.id,
|
2023-01-13 15:56:03 +01:00
|
|
|
title: pg.title.clone(),
|
|
|
|
|
channels,
|
2024-09-10 20:05:29 +02:00
|
|
|
xtream_cluster: pg.xtream_cluster,
|
2023-01-13 15:56:03 +01:00
|
|
|
});
|
|
|
|
|
}
|
2024-05-10 12:01:48 +02:00
|
|
|
}
|
2023-01-13 15:56:03 +01:00
|
|
|
Some(new_playlist)
|
|
|
|
|
}
|
|
|
|
|
|
2024-05-05 19:15:54 +02:00
|
|
|
fn playlistgroup_comparator(a: &PlaylistGroup, b: &PlaylistGroup, group_sort: &ConfigSortGroup, match_as_ascii: bool) -> Ordering {
|
2025-04-10 20:26:22 +02:00
|
|
|
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() };
|
2024-05-05 19:15:54 +02:00
|
|
|
let ordering = value_a.partial_cmp(&value_b).unwrap();
|
|
|
|
|
match group_sort.order {
|
|
|
|
|
Asc => ordering,
|
|
|
|
|
Desc => ordering.reverse()
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
fn playlistitem_comparator(a: &PlaylistItem, b: &PlaylistItem, channel_sort: &ConfigSortChannel, match_as_ascii: bool) -> Ordering {
|
|
|
|
|
let raw_value_a = get_field_value(a, &channel_sort.field);
|
|
|
|
|
let raw_value_b = get_field_value(b, &channel_sort.field);
|
2025-04-10 20:26:22 +02:00
|
|
|
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 };
|
2024-11-04 18:46:56 +01:00
|
|
|
channel_sort.sequence.as_ref().map_or_else(|| {
|
|
|
|
|
let ordering = value_a.partial_cmp(&value_b).unwrap();
|
|
|
|
|
match channel_sort.order {
|
|
|
|
|
Asc => ordering,
|
|
|
|
|
Desc => ordering.reverse()
|
|
|
|
|
}
|
|
|
|
|
}, |custom_order| {
|
2024-10-30 11:14:55 +01:00
|
|
|
// Check indices in the custom order vector
|
2025-03-11 14:56:10 +01:00
|
|
|
let index_a = custom_order.iter().position(|s| s == &value_a);
|
|
|
|
|
let index_b = custom_order.iter().position(|s| s == &value_b);
|
2024-10-30 11:14:55 +01:00
|
|
|
|
|
|
|
|
match (index_a, index_b) {
|
|
|
|
|
(Some(idx_a), Some(idx_b)) => {
|
|
|
|
|
// Both items found in custom order, compare indices
|
|
|
|
|
idx_a.cmp(&idx_b)
|
|
|
|
|
}
|
|
|
|
|
(Some(_), None) => {
|
|
|
|
|
// Only 'a' found in custom order, it comes first
|
|
|
|
|
Ordering::Less
|
|
|
|
|
}
|
|
|
|
|
(None, Some(_)) => {
|
|
|
|
|
// Only 'b' found in custom order, it comes first
|
|
|
|
|
Ordering::Greater
|
|
|
|
|
}
|
|
|
|
|
(None, None) => {
|
|
|
|
|
// Neither found, fall back to default ordering
|
|
|
|
|
let ordering = value_a.partial_cmp(&value_b).unwrap();
|
|
|
|
|
match channel_sort.order {
|
|
|
|
|
Asc => ordering,
|
|
|
|
|
Desc => ordering.reverse(),
|
2024-09-27 14:45:52 +02:00
|
|
|
}
|
|
|
|
|
}
|
2024-10-25 00:39:22 +02:00
|
|
|
}
|
2024-11-04 18:46:56 +01:00
|
|
|
})
|
2024-05-05 19:15:54 +02:00
|
|
|
}
|
|
|
|
|
|
2023-05-03 17:46:39 +02:00
|
|
|
fn sort_playlist(target: &ConfigTarget, new_playlist: &mut [PlaylistGroup]) {
|
2022-04-05 12:23:30 +02:00
|
|
|
if let Some(sort) = &target.sort {
|
2024-05-05 19:15:54 +02:00
|
|
|
let match_as_ascii = sort.match_as_ascii;
|
2023-09-07 11:39:37 +02:00
|
|
|
if let Some(group_sort) = &sort.groups {
|
2024-05-05 19:15:54 +02:00
|
|
|
new_playlist.sort_by(|a, b| playlistgroup_comparator(a, b, group_sort, match_as_ascii));
|
2023-09-07 11:39:37 +02:00
|
|
|
}
|
|
|
|
|
if let Some(channel_sorts) = &sort.channels {
|
2024-05-10 12:01:48 +02:00
|
|
|
for channel_sort in channel_sorts {
|
2023-10-06 19:24:06 +02:00
|
|
|
let regexp = channel_sort.re.as_ref().unwrap();
|
2024-05-10 12:01:48 +02:00
|
|
|
for group in new_playlist.iter_mut() {
|
2025-04-10 20:26:22 +02:00
|
|
|
let group_title = if match_as_ascii { deunicode(&group.title) } else { group.title.to_string() };
|
2024-01-17 13:50:49 +01:00
|
|
|
if regexp.is_match(group_title.as_str()) {
|
2024-05-05 19:15:54 +02:00
|
|
|
group.channels.sort_by(|chan1, chan2| playlistitem_comparator(chan1, chan2, channel_sort, match_as_ascii));
|
2023-09-07 11:39:37 +02:00
|
|
|
}
|
2024-05-10 12:01:48 +02:00
|
|
|
}
|
|
|
|
|
}
|
2023-09-07 11:39:37 +02:00
|
|
|
}
|
2022-04-05 12:23:30 +02:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2025-03-11 14:56:10 +01:00
|
|
|
fn channel_no_playlist(new_playlist: &mut [PlaylistGroup]) {
|
2025-04-12 18:20:46 +02:00
|
|
|
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()).map(str::parse::<u32>)
|
|
|
|
|
.flatten().collect();
|
2025-01-29 17:33:06 +01:00
|
|
|
let mut chno = 1;
|
|
|
|
|
for group in new_playlist {
|
2025-03-11 14:56:10 +01:00
|
|
|
for chan in &mut group.channels {
|
2025-04-12 18:20:46 +02:00
|
|
|
if chan.header.chno.is_empty() {
|
|
|
|
|
while assigned_chnos.contains(&chno) {
|
|
|
|
|
chno += 1;
|
|
|
|
|
}
|
|
|
|
|
chan.header.chno = chno.to_string();
|
|
|
|
|
chno += 1;
|
|
|
|
|
}
|
2025-01-29 17:33:06 +01:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2025-03-11 14:56:10 +01:00
|
|
|
fn exec_rename(pli: &mut PlaylistItem, rename: Option<&Vec<ConfigRename>>) {
|
2023-05-03 17:46:39 +02:00
|
|
|
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);
|
2025-01-01 22:15:49 +01:00
|
|
|
if log::log_enabled!(log::Level::Debug) && *value != cap {
|
|
|
|
|
debug_if_enabled!("Renamed {}={} to {}", &r.field, value, cap);
|
2025-01-01 10:59:51 +01:00
|
|
|
}
|
2023-05-03 17:46:39 +02:00
|
|
|
let value = cap.into_owned();
|
2025-03-11 14:56:10 +01:00
|
|
|
set_field_value(result, &r.field, value);
|
2023-01-13 15:56:03 +01:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2023-10-08 19:54:59 +02:00
|
|
|
fn rename_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option<Vec<PlaylistGroup>> {
|
2023-01-13 15:56:03 +01:00
|
|
|
match &target.rename {
|
|
|
|
|
Some(renames) => {
|
2023-05-03 17:46:39 +02:00
|
|
|
if !renames.is_empty() {
|
2025-01-09 13:31:05 +01:00
|
|
|
let mut new_playlist: Vec<PlaylistGroup> = Vec::with_capacity(playlist.len());
|
2023-01-13 15:56:03 +01:00
|
|
|
for g in playlist {
|
|
|
|
|
let mut grp = g.clone();
|
2023-01-12 15:59:40 +01:00
|
|
|
for r in renames {
|
2024-11-04 18:46:56 +01:00
|
|
|
if matches!(r.field, ItemField::Group) {
|
2023-05-03 17:46:39 +02:00
|
|
|
let cap = r.re.as_ref().unwrap().replace_all(&grp.title, &r.new_name);
|
2024-12-09 19:26:39 +01:00
|
|
|
debug_if_enabled!("Renamed group {} to {} for {}", &grp.title, cap, target.name);
|
2025-03-11 14:56:10 +01:00
|
|
|
grp.title = cap.into_owned();
|
2023-01-12 15:59:40 +01:00
|
|
|
}
|
|
|
|
|
}
|
2023-01-13 15:56:03 +01:00
|
|
|
|
2024-12-02 17:24:06 +01:00
|
|
|
grp.channels.iter_mut().for_each(|pli| exec_rename(pli, target.rename.as_ref()));
|
2023-01-13 15:56:03 +01:00
|
|
|
new_playlist.push(grp);
|
2023-01-12 15:59:40 +01:00
|
|
|
}
|
2023-02-13 12:19:58 +01:00
|
|
|
return Some(new_playlist);
|
2023-01-12 15:59:40 +01:00
|
|
|
}
|
2023-01-13 15:56:03 +01:00
|
|
|
None
|
2023-01-12 15:59:40 +01:00
|
|
|
}
|
2023-01-13 15:56:03 +01:00
|
|
|
_ => None
|
2023-01-12 15:59:40 +01:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2023-04-26 12:52:52 +02:00
|
|
|
macro_rules! apply_pattern {
|
2023-10-08 19:54:59 +02:00
|
|
|
($pattern:expr, $provider:expr, $processor:expr) => {{
|
2024-03-26 12:34:57 +01:00
|
|
|
if let Some(ptrn) = $pattern {
|
|
|
|
|
ptrn.filter($provider, $processor);
|
2023-04-26 12:52:52 +02:00
|
|
|
};
|
|
|
|
|
}};
|
|
|
|
|
}
|
|
|
|
|
|
2025-03-11 14:56:10 +01:00
|
|
|
fn map_channel(mut channel: PlaylistItem, mapping: &Mapping) -> PlaylistItem {
|
2023-05-03 17:46:39 +02:00
|
|
|
if !mapping.mapper.is_empty() {
|
2025-03-11 14:56:10 +01:00
|
|
|
let header = &channel.header;
|
2025-04-10 20:26:22 +02:00
|
|
|
let channel_name = if mapping.match_as_ascii { deunicode(&header.name) } else { header.name.to_string() };
|
2025-04-10 00:06:48 +02:00
|
|
|
if mapping.match_as_ascii && log_enabled!(Level::Trace) { trace!("Decoded {} for matching to {}", &header.name, &channel_name); }
|
2025-03-11 14:56:10 +01:00
|
|
|
// let ref_chan = &mut channel;
|
|
|
|
|
let ref_chan = &mut channel;
|
2023-04-26 12:52:52 +02:00
|
|
|
let mut mock_processor = MockValueProcessor {};
|
2023-01-12 15:59:40 +01:00
|
|
|
for m in &mapping.mapper {
|
2025-03-11 14:56:10 +01:00
|
|
|
let provider = ValueProvider { pli: &ref_chan.clone() };
|
|
|
|
|
let mut processor = MappingValueProcessor { pli: ref_chan, mapper: m };
|
2024-05-10 12:01:48 +02:00
|
|
|
match &m.t_filter {
|
2023-02-12 01:36:21 +01:00
|
|
|
Some(filter) => {
|
2025-03-11 14:56:10 +01:00
|
|
|
|
2023-10-08 19:54:59 +02:00
|
|
|
if filter.filter(&provider, &mut mock_processor) {
|
2024-05-10 12:01:48 +02:00
|
|
|
apply_pattern!(&m.t_pattern, &provider, &mut processor);
|
2023-04-26 12:52:52 +02:00
|
|
|
}
|
2022-04-05 12:23:30 +02:00
|
|
|
}
|
2023-04-26 12:52:52 +02:00
|
|
|
_ => {
|
2024-05-10 12:01:48 +02:00
|
|
|
apply_pattern!(&m.t_pattern, &provider, &mut processor);
|
2023-04-26 12:52:52 +02:00
|
|
|
}
|
2025-04-10 00:06:48 +02:00
|
|
|
}
|
2023-01-12 15:59:40 +01:00
|
|
|
}
|
|
|
|
|
}
|
2023-10-23 18:12:24 +02:00
|
|
|
channel
|
2023-01-12 15:59:40 +01:00
|
|
|
}
|
|
|
|
|
|
2023-10-08 19:54:59 +02:00
|
|
|
fn map_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option<Vec<PlaylistGroup>> {
|
2024-05-10 12:01:48 +02:00
|
|
|
if target.t_mapping.is_some() {
|
2023-05-03 17:46:39 +02:00
|
|
|
let new_playlist: Vec<PlaylistGroup> = playlist.iter().map(|playlist_group| {
|
2023-02-14 18:18:19 +01:00
|
|
|
let mut grp = playlist_group.clone();
|
2024-05-10 12:01:48 +02:00
|
|
|
let mappings = target.t_mapping.as_ref().unwrap();
|
2024-05-04 20:30:20 +02:00
|
|
|
mappings.iter().filter(|&mapping| !mapping.mapper.is_empty()).for_each(|mapping|
|
2024-11-04 18:46:56 +01:00
|
|
|
grp.channels = grp.channels.drain(..).map(|chan| map_channel(chan, mapping)).collect());
|
2023-05-03 17:46:39 +02:00
|
|
|
grp
|
|
|
|
|
}).collect();
|
|
|
|
|
|
2023-04-26 14:30:43 +02:00
|
|
|
// if the group names are changed, restructure channels to the right groups
|
|
|
|
|
// we use
|
2025-01-09 13:31:05 +01:00
|
|
|
let mut new_groups: Vec<PlaylistGroup> = Vec::with_capacity(128);
|
2024-03-28 16:27:39 +01:00
|
|
|
let mut grp_id: u32 = 0;
|
2023-04-26 14:30:43 +02:00
|
|
|
for playlist_group in new_playlist {
|
|
|
|
|
for channel in &playlist_group.channels {
|
2025-03-11 14:56:10 +01:00
|
|
|
let cluster = &channel.header.xtream_cluster;
|
|
|
|
|
let title = &channel.header.group;
|
2024-05-10 12:01:48 +02:00
|
|
|
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,
|
2025-03-11 14:56:10 +01:00
|
|
|
title: title.to_string(),
|
2024-05-10 12:01:48 +02:00
|
|
|
channels: vec![channel.clone()],
|
2024-09-10 20:05:29 +02:00
|
|
|
xtream_cluster: *cluster,
|
2024-05-10 12:01:48 +02:00
|
|
|
});
|
2023-04-26 14:30:43 +02:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
Some(new_groups)
|
2023-01-12 17:15:50 +01:00
|
|
|
} else {
|
|
|
|
|
None
|
2022-04-05 12:23:30 +02:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2025-03-11 14:56:10 +01:00
|
|
|
fn map_playlist_counter(target: &ConfigTarget, playlist: &mut [PlaylistGroup]) {
|
2024-09-11 19:41:50 +02:00
|
|
|
if target.t_mapping.is_some() {
|
|
|
|
|
let mut mock_processor = MockValueProcessor {};
|
|
|
|
|
let mappings = target.t_mapping.as_ref().unwrap();
|
|
|
|
|
for mapping in mappings {
|
|
|
|
|
if let Some(counter_list) = &mapping.t_counter {
|
|
|
|
|
for counter in counter_list {
|
2025-03-11 14:56:10 +01:00
|
|
|
for plg in &mut *playlist {
|
|
|
|
|
for channel in &mut plg.channels {
|
|
|
|
|
let provider = ValueProvider { pli: channel };
|
2024-09-11 19:41:50 +02:00
|
|
|
if counter.filter.filter(&provider, &mut mock_processor) {
|
2025-03-21 15:56:30 +01:00
|
|
|
let cntval = counter.value.load(core::sync::atomic::Ordering::SeqCst);
|
2024-10-25 00:39:22 +02:00
|
|
|
let new_value = if counter.modifier == CounterModifier::Assign {
|
2024-09-11 19:41:50 +02:00
|
|
|
cntval.to_string()
|
|
|
|
|
} else {
|
2025-03-11 14:56:10 +01:00
|
|
|
let value = channel.header.get_field(&counter.field).map_or_else(String::new, |field_value| field_value.to_string());
|
2024-09-11 19:41:50 +02:00
|
|
|
if counter.modifier == CounterModifier::Suffix {
|
2024-10-30 11:14:55 +01:00
|
|
|
format!("{value}{}{cntval}", counter.concat)
|
2024-10-25 00:39:22 +02:00
|
|
|
} else {
|
2024-10-30 11:14:55 +01:00
|
|
|
format!("{cntval}{}{value}", counter.concat)
|
2024-09-11 19:41:50 +02:00
|
|
|
}
|
|
|
|
|
};
|
2025-03-11 14:56:10 +01:00
|
|
|
channel.header.set_field(&counter.field, new_value.as_str());
|
2025-03-21 15:56:30 +01:00
|
|
|
counter.value.fetch_add(1, core::sync::atomic::Ordering::SeqCst);
|
2024-09-11 19:41:50 +02:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2023-10-12 11:08:37 +02:00
|
|
|
// 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.
|
2025-02-14 17:51:33 +01:00
|
|
|
fn is_input_enabled(enabled_inputs: usize, input: &ConfigInput, user_targets: &ProcessTargets) -> bool {
|
|
|
|
|
let input_enabled = input.enabled;
|
|
|
|
|
let input_id = input.id;
|
2023-10-12 11:08:37 +02:00
|
|
|
if enabled_inputs == 0 {
|
|
|
|
|
return user_targets.enabled && user_targets.has_input(input_id);
|
|
|
|
|
}
|
|
|
|
|
input_enabled
|
|
|
|
|
}
|
|
|
|
|
|
2023-12-08 19:08:04 +01:00
|
|
|
fn is_target_enabled(target: &ConfigTarget, user_targets: &ProcessTargets) -> bool {
|
|
|
|
|
(!user_targets.enabled && target.enabled) || (user_targets.enabled && user_targets.has_target(target.id))
|
|
|
|
|
}
|
|
|
|
|
|
2025-01-09 15:58:38 +01:00
|
|
|
async fn process_source(client: Arc<reqwest::Client>, cfg: Arc<Config>, source_idx: usize, user_targets: Arc<ProcessTargets>) -> (Vec<InputStats>, Vec<TargetStats>, Vec<M3uFilterError>) {
|
2023-02-13 12:19:58 +01:00
|
|
|
let source = cfg.sources.get(source_idx).unwrap();
|
2023-10-13 13:59:22 +02:00
|
|
|
let mut errors = vec![];
|
2025-01-27 15:58:57 +01:00
|
|
|
let mut input_stats = HashMap::<String, InputStats>::new();
|
2025-01-05 13:29:28 +01:00
|
|
|
let mut target_stats = Vec::<TargetStats>::new();
|
2025-01-09 13:31:05 +01:00
|
|
|
let mut source_playlists = Vec::with_capacity(128);
|
2025-02-14 17:51:33 +01:00
|
|
|
let enabled_inputs = source.inputs.iter().filter(|&input| input.enabled).count();
|
|
|
|
|
// Download the sources
|
2023-10-12 11:08:37 +02:00
|
|
|
for input in &source.inputs {
|
2025-02-14 17:51:33 +01:00
|
|
|
if is_input_enabled(enabled_inputs, input, &user_targets) {
|
2025-01-17 18:32:17 +01:00
|
|
|
let start_time = Instant::now();
|
2024-09-10 20:05:29 +02:00
|
|
|
let (mut playlistgroups, mut error_list) = match input.input_type {
|
2025-02-05 20:01:49 +01:00
|
|
|
InputType::M3u => m3u::get_m3u_playlist(Arc::clone(&client), &cfg, input, &cfg.working_dir).await,
|
|
|
|
|
InputType::Xtream => xtream::get_xtream_playlist(Arc::clone(&client), input, &cfg.working_dir).await,
|
2025-03-21 18:48:13 +01:00
|
|
|
InputType::M3uBatch | InputType::XtreamBatch => (vec![], vec![])
|
2023-10-12 11:08:37 +02:00
|
|
|
};
|
2023-10-27 19:03:29 +02:00
|
|
|
let (tvguide, mut tvguide_errors) = if error_list.is_empty() {
|
2025-02-05 20:01:49 +01:00
|
|
|
epg::get_xmltv(Arc::clone(&client), &cfg, input, &cfg.working_dir).await
|
2023-12-08 19:08:04 +01:00
|
|
|
} else {
|
|
|
|
|
(None, vec![])
|
2023-10-27 19:03:29 +02:00
|
|
|
};
|
2024-11-02 16:44:38 +01:00
|
|
|
errors.append(&mut error_list);
|
|
|
|
|
errors.append(&mut tvguide_errors);
|
2024-09-10 20:05:29 +02:00
|
|
|
let group_count = playlistgroups.len();
|
|
|
|
|
let channel_count = playlistgroups.iter()
|
2023-10-19 18:56:08 +02:00
|
|
|
.map(|group| group.channels.len())
|
|
|
|
|
.sum();
|
2025-01-27 15:58:57 +01:00
|
|
|
let input_name = &input.name;
|
2024-09-10 20:05:29 +02:00
|
|
|
if playlistgroups.is_empty() {
|
2025-01-17 17:16:08 +01:00
|
|
|
info!("Source is empty {input_name}");
|
|
|
|
|
errors.push(notify_err!(format!("Source is empty {input_name}")));
|
2023-10-19 18:56:08 +02:00
|
|
|
} else {
|
2024-09-10 20:05:29 +02:00
|
|
|
playlistgroups.iter_mut().for_each(PlaylistGroup::on_load);
|
|
|
|
|
source_playlists.push(
|
2023-10-13 13:59:22 +02:00
|
|
|
FetchedPlaylist {
|
2023-10-13 15:56:00 +02:00
|
|
|
input,
|
2024-09-10 20:05:29 +02:00
|
|
|
playlistgroups,
|
2023-12-08 19:08:04 +01:00
|
|
|
epg: tvguide,
|
2023-10-13 13:59:22 +02:00
|
|
|
}
|
|
|
|
|
);
|
2022-03-24 14:08:25 +01:00
|
|
|
}
|
2024-11-03 14:52:56 +01:00
|
|
|
let elapsed = start_time.elapsed().as_secs();
|
2025-01-27 15:58:57 +01:00
|
|
|
input_stats.insert(input_name.to_string(), create_input_stat(group_count, channel_count, error_list.len(),
|
2025-04-10 00:06:48 +02:00
|
|
|
input.input_type, input_name, elapsed));
|
2022-03-24 14:08:25 +01:00
|
|
|
}
|
2023-02-13 12:19:58 +01:00
|
|
|
}
|
2024-09-10 20:05:29 +02:00
|
|
|
if source_playlists.is_empty() {
|
2024-12-09 19:26:39 +01:00
|
|
|
debug!("Source at index {source_idx} is empty");
|
2024-12-23 11:30:33 +01:00
|
|
|
errors.push(notify_err!(format!("Source at {source_idx} is empty")));
|
2023-10-12 11:08:37 +02:00
|
|
|
} else {
|
2024-12-09 19:26:39 +01:00
|
|
|
debug_if_enabled!("Source has {} groups", source_playlists.iter().map(|fpl| fpl.playlistgroups.len()).sum::<usize>());
|
2023-12-08 19:08:04 +01:00
|
|
|
for target in &source.targets {
|
|
|
|
|
if is_target_enabled(target, &user_targets) {
|
2025-01-09 15:58:38 +01:00
|
|
|
match process_playlist_for_target(Arc::clone(&client), &mut source_playlists, target, &cfg, &mut input_stats, &mut errors).await {
|
2025-01-05 13:29:28 +01:00
|
|
|
Ok(()) => {
|
|
|
|
|
target_stats.push(TargetStats::success(&target.name));
|
|
|
|
|
}
|
|
|
|
|
Err(mut err) => {
|
|
|
|
|
target_stats.push(TargetStats::failure(&target.name));
|
|
|
|
|
errors.append(&mut err);
|
|
|
|
|
}
|
2023-10-12 11:08:37 +02:00
|
|
|
}
|
|
|
|
|
}
|
2023-12-08 19:08:04 +01:00
|
|
|
}
|
2023-10-12 11:08:37 +02:00
|
|
|
}
|
2025-01-05 13:29:28 +01:00
|
|
|
(input_stats.into_values().collect(), target_stats, errors)
|
2023-02-13 12:19:58 +01:00
|
|
|
}
|
|
|
|
|
|
2024-11-03 14:52:56 +01:00
|
|
|
fn create_input_stat(group_count: usize, channel_count: usize, error_count: usize, input_type: InputType, input_name: &str, secs_took: u64) -> InputStats {
|
2024-10-25 00:39:22 +02:00
|
|
|
InputStats {
|
|
|
|
|
name: input_name.to_string(),
|
2024-10-30 11:14:55 +01:00
|
|
|
input_type,
|
|
|
|
|
error_count,
|
2024-10-25 00:39:22 +02:00
|
|
|
raw_stats: PlaylistStats {
|
|
|
|
|
group_count,
|
|
|
|
|
channel_count,
|
|
|
|
|
},
|
|
|
|
|
processed_stats: PlaylistStats {
|
|
|
|
|
group_count: 0,
|
|
|
|
|
channel_count: 0,
|
|
|
|
|
},
|
2024-11-04 18:46:56 +01:00
|
|
|
secs_took,
|
2024-10-25 00:39:22 +02:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2025-01-09 15:58:38 +01:00
|
|
|
async fn process_sources(client: Arc<reqwest::Client>, config: Arc<Config>, user_targets: Arc<ProcessTargets>) -> (Vec<SourceStats>, Vec<M3uFilterError>) {
|
2023-02-25 16:23:26 +01:00
|
|
|
let mut handle_list = vec![];
|
2023-09-26 19:53:45 +02:00
|
|
|
let thread_num = config.threads;
|
|
|
|
|
let process_parallel = thread_num > 1 && config.sources.len() > 1;
|
2024-03-26 16:07:31 +01:00
|
|
|
if process_parallel && log_enabled!(Level::Debug) {
|
2025-04-10 00:06:48 +02:00
|
|
|
debug!("Using {thread_num} threads");
|
2023-11-03 18:05:13 +01:00
|
|
|
}
|
2023-10-13 13:59:22 +02:00
|
|
|
let errors = Arc::new(Mutex::<Vec<M3uFilterError>>::new(vec![]));
|
2025-01-05 13:29:28 +01:00
|
|
|
let stats = Arc::new(Mutex::<Vec<SourceStats>>::new(vec![]));
|
2023-09-26 19:53:45 +02:00
|
|
|
for (index, _) in config.sources.iter().enumerate() {
|
2025-01-06 09:17:05 +01:00
|
|
|
// We're using the file lock this way on purpose
|
|
|
|
|
let source_lock_path = PathBuf::from(format!("source_{index}"));
|
2025-03-11 14:56:10 +01:00
|
|
|
let Ok(update_lock) = config.file_locks.try_write_lock(&source_lock_path).await else {
|
2025-01-06 09:17:05 +01:00
|
|
|
warn!("The update operation for the source at index {index} was skipped because an update is already in progress.");
|
|
|
|
|
continue;
|
|
|
|
|
};
|
|
|
|
|
|
2023-10-13 13:59:22 +02:00
|
|
|
let shared_errors = errors.clone();
|
2023-10-19 18:56:08 +02:00
|
|
|
let shared_stats = stats.clone();
|
|
|
|
|
let cfg = config.clone();
|
|
|
|
|
let usr_trgts = user_targets.clone();
|
2023-02-13 12:19:58 +01:00
|
|
|
if process_parallel {
|
2025-01-09 15:58:38 +01:00
|
|
|
let http_client = Arc::clone(&client);
|
2023-02-13 19:08:28 +01:00
|
|
|
let handles = &mut handle_list;
|
2023-11-03 18:05:13 +01:00
|
|
|
let process = move || {
|
2025-03-11 14:56:10 +01:00
|
|
|
// TODO better way ?
|
|
|
|
|
let rt = tokio::runtime::Runtime::new().unwrap();
|
|
|
|
|
rt.block_on(async {
|
2025-01-09 15:58:38 +01:00
|
|
|
let (input_stats, target_stats, mut res_errors) = process_source(Arc::clone(&http_client), cfg, index, usr_trgts).await;
|
2025-03-11 14:56:10 +01:00
|
|
|
shared_errors.lock().await.append(&mut res_errors);
|
2025-01-05 13:29:28 +01:00
|
|
|
let process_stats = SourceStats::new(input_stats, target_stats);
|
2025-03-11 14:56:10 +01:00
|
|
|
shared_stats.lock().await.push(process_stats);
|
2023-11-03 18:05:13 +01:00
|
|
|
});
|
|
|
|
|
};
|
2023-05-03 17:46:39 +02:00
|
|
|
handles.push(thread::spawn(process));
|
2024-05-10 12:01:48 +02:00
|
|
|
if handles.len() >= thread_num as usize {
|
2023-10-13 13:59:22 +02:00
|
|
|
handles.drain(..).for_each(|handle| { let _ = handle.join(); });
|
2023-02-13 19:08:28 +01:00
|
|
|
}
|
2023-02-13 12:19:58 +01:00
|
|
|
} else {
|
2025-01-09 15:58:38 +01:00
|
|
|
let (input_stats, target_stats, mut res_errors) = process_source(Arc::clone(&client), cfg, index, usr_trgts).await;
|
2025-03-11 14:56:10 +01:00
|
|
|
shared_errors.lock().await.append(&mut res_errors);
|
2025-01-05 13:29:28 +01:00
|
|
|
let process_stats = SourceStats::new(input_stats, target_stats);
|
2025-03-11 14:56:10 +01:00
|
|
|
shared_stats.lock().await.push(process_stats);
|
2023-02-13 12:19:58 +01:00
|
|
|
}
|
2025-01-06 09:17:05 +01:00
|
|
|
drop(update_lock);
|
2023-02-13 12:19:58 +01:00
|
|
|
}
|
2023-02-13 19:08:28 +01:00
|
|
|
for handle in handle_list {
|
2023-02-25 16:23:26 +01:00
|
|
|
let _ = handle.join();
|
2022-03-24 14:08:25 +01:00
|
|
|
}
|
2024-12-05 19:54:34 +01:00
|
|
|
(Arc::try_unwrap(stats).unwrap().into_inner(), Arc::try_unwrap(errors).unwrap().into_inner())
|
2022-03-24 14:08:25 +01:00
|
|
|
}
|
2023-09-29 16:03:20 +02:00
|
|
|
|
2024-05-04 20:30:20 +02:00
|
|
|
pub type ProcessingPipe = Vec<fn(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option<Vec<PlaylistGroup>>>;
|
2023-09-29 16:03:20 +02:00
|
|
|
|
2023-10-13 13:59:22 +02:00
|
|
|
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]
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2025-01-06 10:11:14 +01:00
|
|
|
fn duplicate_hash(item: &PlaylistItem) -> UUIDType {
|
2025-01-29 17:33:06 +01:00
|
|
|
item.get_uuid()
|
2025-01-06 10:11:14 +01:00
|
|
|
}
|
2024-09-10 20:05:29 +02:00
|
|
|
|
2025-01-06 10:11:14 +01:00
|
|
|
fn execute_pipe<'a>(target: &ConfigTarget, pipe: &ProcessingPipe, fpl: &FetchedPlaylist<'a>, duplicates: &mut HashSet<UUIDType>) -> FetchedPlaylist<'a> {
|
2024-09-10 20:05:29 +02:00
|
|
|
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(),
|
|
|
|
|
};
|
2025-01-06 10:11:14 +01:00
|
|
|
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)));
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2024-09-10 20:05:29 +02:00
|
|
|
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.
|
2024-11-04 18:46:56 +01:00
|
|
|
fn flatten_groups(playlistgroups: Vec<PlaylistGroup>) -> Vec<PlaylistGroup> {
|
2024-09-11 12:13:05 +02:00
|
|
|
let mut sort_order: Vec<PlaylistGroup> = vec![];
|
|
|
|
|
let mut idx: usize = 0;
|
2025-03-11 14:56:10 +01:00
|
|
|
let mut group_map: HashMap<(String, XtreamCluster), usize> = HashMap::new();
|
2024-11-04 18:46:56 +01:00
|
|
|
for group in playlistgroups {
|
2025-03-11 14:56:10 +01:00
|
|
|
let key = (group.title.to_string(), group.xtream_cluster);
|
2024-09-10 20:05:29 +02:00
|
|
|
match group_map.entry(key) {
|
2024-09-11 12:13:05 +02:00
|
|
|
std::collections::hash_map::Entry::Vacant(v) => {
|
|
|
|
|
v.insert(idx);
|
|
|
|
|
idx += 1;
|
|
|
|
|
sort_order.push(group);
|
|
|
|
|
}
|
2024-09-10 20:05:29 +02:00
|
|
|
std::collections::hash_map::Entry::Occupied(o) => {
|
2024-09-11 12:13:05 +02:00
|
|
|
sort_order.get_mut(*o.get()).unwrap().channels.extend(group.channels);
|
2024-09-10 20:05:29 +02:00
|
|
|
}
|
2025-04-10 00:06:48 +02:00
|
|
|
}
|
2024-11-04 18:46:56 +01:00
|
|
|
}
|
2024-09-11 12:13:05 +02:00
|
|
|
sort_order
|
2024-09-10 20:05:29 +02:00
|
|
|
}
|
|
|
|
|
|
2025-04-12 14:14:06 +02:00
|
|
|
type PhoneticCodeAndMaybeEpgId<'a> = (Option<Cow<'a, str>>, Option<Cow<'a, str>>);
|
2025-04-11 10:04:24 +02:00
|
|
|
pub struct EpgIdCache<'a > {
|
|
|
|
|
pub channel: HashSet<Cow<'a, str>>,
|
2025-04-12 14:14:06 +02:00
|
|
|
pub normalized: HashMap<Cow<'a, str>, PhoneticCodeAndMaybeEpgId<'a>>,
|
|
|
|
|
pub normalized_phonetic: HashMap<Cow<'a, str>, Cow<'a, str>>,
|
2025-04-11 10:04:24 +02:00
|
|
|
pub processed: HashSet<String>,
|
2025-04-12 14:14:06 +02:00
|
|
|
pub smart_match_config: EpgSmartMatchConfig,
|
|
|
|
|
pub metaphone: Metaphone,
|
2025-04-11 10:04:24 +02:00
|
|
|
}
|
|
|
|
|
|
2025-04-11 20:25:09 +02:00
|
|
|
impl EpgIdCache<'_> {
|
|
|
|
|
pub fn new(epg_config: Option<&EpgConfig>) -> Self {
|
2025-04-12 14:14:06 +02:00
|
|
|
let normalize_config = epg_config.map_or_else(EpgSmartMatchConfig::default, |epg_config| epg_config.t_smart_match.clone());
|
2025-04-11 10:04:24 +02:00
|
|
|
EpgIdCache {
|
|
|
|
|
channel: HashSet::new(),
|
|
|
|
|
normalized: HashMap::new(),
|
2025-04-11 20:25:09 +02:00
|
|
|
normalized_phonetic: HashMap::new(),
|
2025-04-11 10:04:24 +02:00
|
|
|
processed: HashSet::new(),
|
2025-04-12 14:14:06 +02:00
|
|
|
smart_match_config: normalize_config,
|
|
|
|
|
metaphone: Metaphone::default(),
|
2025-04-11 10:04:24 +02:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2025-04-12 14:14:06 +02:00
|
|
|
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));
|
|
|
|
|
}
|
|
|
|
|
|
2025-04-11 10:04:24 +02:00
|
|
|
pub fn normalize(&self, name: &str) -> String {
|
2025-04-12 14:14:06 +02:00
|
|
|
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, normalize_enabled: bool, fuzzy_matching: bool) {
|
|
|
|
|
for channel in fp.playlistgroups.iter().flat_map(|g| &g.channels) {
|
|
|
|
|
let epg_id = channel.header.epg_channel_id.as_deref();
|
|
|
|
|
let name = &channel.header.name;
|
|
|
|
|
|
|
|
|
|
// 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 = normalize_enabled && (fuzzy_matching || epg_id.is_none_or(str::is_empty));
|
|
|
|
|
|
|
|
|
|
if needs_normalization {
|
|
|
|
|
id_cache.normalize_and_store(name);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2025-04-12 17:58:16 +02:00
|
|
|
fn assign_channel_epg(new_epg: &mut Vec<Epg>, fp: &mut FetchedPlaylist, smart_match_enabled: bool, id_cache: &mut EpgIdCache) {
|
2025-04-12 14:14:06 +02:00
|
|
|
if let Some(tv_guide) = &fp.epg {
|
|
|
|
|
if let Some(epg) = tv_guide.filter(id_cache) {
|
|
|
|
|
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();
|
|
|
|
|
fp.playlistgroups.iter_mut()
|
|
|
|
|
.flat_map(|g| &mut g.channels)
|
|
|
|
|
.filter(|c| c.header.xtream_cluster == XtreamCluster::Live)
|
|
|
|
|
.filter(|c| c.header.epg_channel_id.is_none() || c.header.logo.is_empty() || c.header.logo_small.is_empty())
|
|
|
|
|
.for_each(|c| {
|
2025-04-12 17:58:16 +02:00
|
|
|
if 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 c.header.epg_channel_id.is_none() || !id_cache.processed.contains(c.header.epg_channel_id.as_ref().unwrap()) {
|
|
|
|
|
let normalized = id_cache.normalize(&c.header.name);
|
|
|
|
|
if let Some((_, Some(epg_id))) = id_cache.normalized.get(&Cow::Borrowed(normalized.as_str())) {
|
|
|
|
|
c.header.epg_channel_id = Some(epg_id.to_string());
|
|
|
|
|
}
|
2025-04-12 14:14:06 +02:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
if c.header.epg_channel_id.is_some() && (c.header.logo.is_empty() || c.header.logo_small.is_empty()) {
|
|
|
|
|
if let Some(icon_tag) = icon_tags.get(c.header.epg_channel_id.as_ref().unwrap()) {
|
|
|
|
|
if let Some(icon) = icon_tag.icon.as_ref() {
|
|
|
|
|
if c.header.logo.is_empty() {
|
|
|
|
|
c.header.logo = (*icon).to_string();
|
|
|
|
|
}
|
|
|
|
|
if c.header.logo_small.is_empty() {
|
|
|
|
|
c.header.logo = (*icon).to_string();
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
});
|
|
|
|
|
new_epg.push(epg);
|
|
|
|
|
}
|
2025-04-11 10:04:24 +02:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2025-01-09 15:58:38 +01:00
|
|
|
async fn process_playlist_for_target(client: Arc<reqwest::Client>,
|
|
|
|
|
playlists: &mut [FetchedPlaylist<'_>],
|
2025-01-06 10:11:14 +01:00
|
|
|
target: &ConfigTarget,
|
|
|
|
|
cfg: &Config,
|
2025-01-27 15:58:57 +01:00
|
|
|
stats: &mut HashMap<String, InputStats>,
|
2025-01-06 10:11:14 +01:00
|
|
|
errors: &mut Vec<M3uFilterError>) -> Result<(), Vec<M3uFilterError>> {
|
2023-10-13 13:59:22 +02:00
|
|
|
let pipe = get_processing_pipe(target);
|
2024-12-09 19:26:39 +01:00
|
|
|
debug_if_enabled!("Processing order is {}", &target.processing_order);
|
2023-09-29 16:03:20 +02:00
|
|
|
|
2025-01-06 10:11:14 +01:00
|
|
|
let mut duplicates: HashSet<UUIDType> = HashSet::new();
|
2024-12-10 23:16:04 +01:00
|
|
|
let mut processed_fetched_playlists: Vec<FetchedPlaylist> = vec![];
|
|
|
|
|
for provider_fpl in playlists.iter_mut() {
|
2025-01-06 10:11:14 +01:00
|
|
|
let mut processed_fpl = execute_pipe(target, &pipe, provider_fpl, &mut duplicates);
|
2025-01-09 15:58:38 +01:00
|
|
|
playlist_resolve_series(Arc::clone(&client), cfg, target, errors, &pipe, provider_fpl, &mut processed_fpl).await;
|
2025-03-11 14:56:10 +01:00
|
|
|
playlist_resolve_vod(Arc::clone(&client), cfg, target, errors, &mut processed_fpl).await;
|
2023-10-19 18:56:08 +02:00
|
|
|
// stats
|
2025-01-27 15:58:57 +01:00
|
|
|
let input_stats = stats.get_mut(&processed_fpl.input.name);
|
2023-10-19 18:56:08 +02:00
|
|
|
if let Some(stat) = input_stats {
|
2024-12-10 23:16:04 +01:00
|
|
|
stat.processed_stats.group_count = processed_fpl.playlistgroups.len();
|
|
|
|
|
stat.processed_stats.channel_count = processed_fpl.playlistgroups.iter()
|
2023-10-19 18:56:08 +02:00
|
|
|
.map(|group| group.channels.len())
|
|
|
|
|
.sum();
|
|
|
|
|
}
|
2024-12-10 23:16:04 +01:00
|
|
|
processed_fetched_playlists.push(processed_fpl);
|
2023-12-08 19:08:04 +01:00
|
|
|
}
|
2023-09-29 16:03:20 +02:00
|
|
|
|
2024-12-10 23:16:04 +01:00
|
|
|
apply_affixes(&mut processed_fetched_playlists);
|
2024-03-15 18:34:39 +01:00
|
|
|
|
2023-10-13 15:56:00 +02:00
|
|
|
let mut new_playlist = vec![];
|
2023-10-27 19:03:29 +02:00
|
|
|
let mut new_epg = vec![];
|
2024-05-05 19:15:54 +02:00
|
|
|
|
2024-11-02 16:44:38 +01:00
|
|
|
// each fetched playlist can have its own epgl url.
|
|
|
|
|
// we need to process each input epg.
|
2024-12-10 23:16:04 +01:00
|
|
|
for mut fp in processed_fetched_playlists {
|
2024-10-25 00:39:22 +02:00
|
|
|
// collect all epg_channel ids
|
2025-04-11 20:25:09 +02:00
|
|
|
let mut id_cache = EpgIdCache::new(fp.input.epg.as_ref());
|
2025-04-12 17:58:16 +02:00
|
|
|
let smart_match_enabled = id_cache.smart_match_config.enabled;
|
|
|
|
|
let fuzzy_matching = smart_match_enabled && id_cache.smart_match_config.fuzzy_matching;
|
|
|
|
|
prepare_epg_id_cache(&mut fp, &mut id_cache, smart_match_enabled, fuzzy_matching);
|
2025-04-09 20:51:25 +02:00
|
|
|
// let epg_channel_ids: HashSet<_> = fp.playlistgroups.iter().flat_map(|g| &g.channels)
|
|
|
|
|
// .filter_map(|c| c.header.epg_channel_id.as_ref()).map(|a| a.as_str()).collect();
|
2025-04-11 10:04:24 +02:00
|
|
|
if id_cache.channel.is_empty() && id_cache.normalized.is_empty() {
|
2025-04-09 20:51:25 +02:00
|
|
|
debug!("channel ids are empty");
|
2025-04-12 14:14:06 +02:00
|
|
|
} else {
|
2025-04-09 20:51:25 +02:00
|
|
|
debug_if_enabled!("found epg information for {}", &target.name);
|
2025-04-12 17:58:16 +02:00
|
|
|
assign_channel_epg(&mut new_epg, &mut fp, smart_match_enabled, &mut id_cache);
|
2023-10-27 19:03:29 +02:00
|
|
|
}
|
2025-04-06 16:50:43 +02:00
|
|
|
new_playlist.append(&mut fp.playlistgroups);
|
2024-11-04 18:46:56 +01:00
|
|
|
}
|
2023-09-29 16:03:20 +02:00
|
|
|
|
2024-05-10 12:01:48 +02:00
|
|
|
if new_playlist.is_empty() {
|
|
|
|
|
info!("Playlist is empty: {}", &target.name);
|
|
|
|
|
Ok(())
|
|
|
|
|
} else {
|
2024-09-10 20:05:29 +02:00
|
|
|
let mut flat_new_playlist = flatten_groups(new_playlist);
|
2024-09-11 08:17:24 +02:00
|
|
|
sort_playlist(target, &mut flat_new_playlist);
|
2025-03-11 14:56:10 +01:00
|
|
|
channel_no_playlist(&mut flat_new_playlist);
|
|
|
|
|
map_playlist_counter(target, &mut flat_new_playlist);
|
2024-09-11 12:13:05 +02:00
|
|
|
process_watch(target, cfg, &flat_new_playlist);
|
2024-12-05 19:54:34 +01:00
|
|
|
persist_playlist(&mut flat_new_playlist, flatten_tvguide(&new_epg).as_ref(), target, cfg).await
|
2024-09-10 20:05:29 +02:00
|
|
|
}
|
|
|
|
|
}
|
2023-10-13 18:13:16 +02:00
|
|
|
|
2024-09-10 20:05:29 +02:00
|
|
|
fn process_watch(target: &ConfigTarget, cfg: &Config, new_playlist: &Vec<PlaylistGroup>) {
|
|
|
|
|
if target.t_watch_re.is_some() {
|
|
|
|
|
if default_as_default().eq_ignore_ascii_case(&target.name) {
|
|
|
|
|
error!("cant watch a target with no unique name");
|
|
|
|
|
} else {
|
|
|
|
|
let watch_re = target.t_watch_re.as_ref().unwrap();
|
|
|
|
|
for pl in new_playlist {
|
|
|
|
|
if watch_re.iter().any(|r| r.is_match(&pl.title)) {
|
|
|
|
|
process_group_watch(cfg, &target.name, pl);
|
2024-05-10 12:01:48 +02:00
|
|
|
}
|
2023-10-13 18:13:16 +02:00
|
|
|
}
|
|
|
|
|
}
|
2023-10-06 19:24:06 +02:00
|
|
|
}
|
2023-09-29 16:03:20 +02:00
|
|
|
}
|
2023-10-19 18:56:08 +02:00
|
|
|
|
2025-01-09 15:58:38 +01:00
|
|
|
pub async fn exec_processing(client: Arc<reqwest::Client>, cfg: Arc<Config>, targets: Arc<ProcessTargets>) {
|
2024-12-25 14:12:30 +01:00
|
|
|
let start_time = Instant::now();
|
2025-01-09 15:58:38 +01:00
|
|
|
let (stats, errors) = process_sources(client, cfg.clone(), targets.clone()).await;
|
2024-12-13 17:51:47 +01:00
|
|
|
// log errors
|
2024-12-28 20:01:28 +01:00
|
|
|
for err in &errors {
|
|
|
|
|
error!("{}", err.message);
|
|
|
|
|
}
|
2024-12-25 14:12:30 +01:00
|
|
|
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
|
2025-04-10 00:06:48 +02:00
|
|
|
info!("{stats_msg}");
|
2024-12-25 14:12:30 +01:00
|
|
|
// send stats
|
|
|
|
|
send_message(&MsgKind::Stats, cfg.messaging.as_ref(), stats_msg.as_str());
|
2024-12-13 17:51:47 +01:00
|
|
|
}
|
2023-10-20 08:23:25 +02:00
|
|
|
// send errors
|
2023-10-19 18:56:08 +02:00
|
|
|
if let Some(message) = get_errors_notify_message!(errors, 255) {
|
2024-12-25 14:12:30 +01:00
|
|
|
if let Ok(error_msg) = serde_json::to_string(&serde_json::Value::Object(serde_json::map::Map::from_iter([("errors".to_string(), serde_json::Value::String(message))]))) {
|
|
|
|
|
send_message(&MsgKind::Error, cfg.messaging.as_ref(), error_msg.as_str());
|
|
|
|
|
}
|
2023-10-19 18:56:08 +02:00
|
|
|
}
|
2024-12-25 14:12:30 +01:00
|
|
|
let elapsed = start_time.elapsed().as_secs();
|
|
|
|
|
info!("Update process finished! Took {elapsed} secs.");
|
2025-02-01 21:50:13 -06:00
|
|
|
}
|
2025-04-11 20:25:09 +02:00
|
|
|
|
|
|
|
|
#[cfg(test)]
|
|
|
|
|
mod tests {
|
|
|
|
|
#[test]
|
|
|
|
|
fn test() {
|
|
|
|
|
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));
|
|
|
|
|
}
|
|
|
|
|
}
|