mirror of
https://github.com/euzu/tuliprox.git
synced 2026-10-01 05:22:28 +02:00
Merged groups with same title and cluster together
after generating the complete playlist.
This commit is contained in:
@@ -18,7 +18,7 @@ use crate::messaging::{MsgKind, send_message};
|
||||
use crate::model::config::{ConfigSortChannel, ConfigSortGroup, ConfigTarget, InputType,
|
||||
ItemField, ProcessingOrder, ProcessTargets, SortOrder::{Asc, Desc}};
|
||||
use crate::model::mapping::{Mapping, MappingValueProcessor};
|
||||
use crate::model::playlist::{FetchedPlaylist, PlaylistGroup, PlaylistItem};
|
||||
use crate::model::playlist::{FetchedPlaylist, PlaylistGroup, PlaylistItem, XtreamCluster};
|
||||
use crate::model::stats::{InputStats, PlaylistStats};
|
||||
use crate::processing::affix_processor::apply_affixes;
|
||||
use crate::processing::playlist_watch::process_group_watch;
|
||||
@@ -46,7 +46,7 @@ fn filter_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Opt
|
||||
id: pg.id,
|
||||
title: pg.title.clone(),
|
||||
channels,
|
||||
xtream_cluster: pg.xtream_cluster.clone(),
|
||||
xtream_cluster: pg.xtream_cluster,
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -180,7 +180,7 @@ fn map_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option
|
||||
let mut grp = playlist_group.clone();
|
||||
let mappings = target.t_mapping.as_ref().unwrap();
|
||||
mappings.iter().filter(|&mapping| !mapping.mapper.is_empty()).for_each(|mapping|
|
||||
grp.channels = grp.channels.drain(..).map(|chan| map_channel(chan, mapping)).collect());
|
||||
grp.channels = grp.channels.drain(..).map(|chan| map_channel(chan, mapping)).collect());
|
||||
grp
|
||||
}).collect();
|
||||
|
||||
@@ -200,7 +200,7 @@ fn map_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option
|
||||
id: grp_id,
|
||||
title: Rc::clone(title),
|
||||
channels: vec![channel.clone()],
|
||||
xtream_cluster: cluster.clone(),
|
||||
xtream_cluster: *cluster,
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -227,14 +227,14 @@ fn is_target_enabled(target: &ConfigTarget, user_targets: &ProcessTargets) -> bo
|
||||
|
||||
async fn process_source(cfg: Arc<Config>, source_idx: usize, user_targets: Arc<ProcessTargets>) -> (Vec<InputStats>, Vec<M3uFilterError>) {
|
||||
let source = cfg.sources.get(source_idx).unwrap();
|
||||
let mut all_playlist = Vec::new();
|
||||
let mut source_playlists = Vec::new();
|
||||
let enabled_inputs = source.inputs.iter().filter(|item| item.enabled).count();
|
||||
let mut errors = vec![];
|
||||
let mut stats = HashMap::<u16, InputStats>::new();
|
||||
for input in &source.inputs {
|
||||
let input_id = input.id;
|
||||
if is_input_enabled(enabled_inputs, input.enabled, input_id, &user_targets) {
|
||||
let (playlist, mut error_list) = match input.input_type {
|
||||
let (mut playlistgroups, mut error_list) = match input.input_type {
|
||||
InputType::M3u => download::get_m3u_playlist(&cfg, input, &cfg.working_dir).await,
|
||||
InputType::Xtream => download::get_xtream_playlist(input, &cfg.working_dir).await,
|
||||
};
|
||||
@@ -249,18 +249,19 @@ async fn process_source(cfg: Arc<Config>, source_idx: usize, user_targets: Arc<P
|
||||
None => input.url.as_str(),
|
||||
Some(name_val) => name_val.as_str()
|
||||
};
|
||||
let group_count = playlist.len();
|
||||
let channel_count = playlist.iter()
|
||||
let group_count = playlistgroups.len();
|
||||
let channel_count = playlistgroups.iter()
|
||||
.map(|group| group.channels.len())
|
||||
.sum();
|
||||
if playlist.is_empty() {
|
||||
if playlistgroups.is_empty() {
|
||||
info!("source is empty {}", input.url);
|
||||
errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("source is empty {input_name}")));
|
||||
} else {
|
||||
all_playlist.push(
|
||||
playlistgroups.iter_mut().for_each(PlaylistGroup::on_load);
|
||||
source_playlists.push(
|
||||
FetchedPlaylist {
|
||||
input,
|
||||
playlist,
|
||||
playlistgroups,
|
||||
epg: tvguide,
|
||||
}
|
||||
);
|
||||
@@ -280,18 +281,18 @@ async fn process_source(cfg: Arc<Config>, source_idx: usize, user_targets: Arc<P
|
||||
});
|
||||
}
|
||||
}
|
||||
if all_playlist.is_empty() {
|
||||
if source_playlists.is_empty() {
|
||||
if log_enabled!(Level::Debug) {
|
||||
debug!("Source at index {} input is empty", source_idx);
|
||||
debug!("Source at index {source_idx} is empty");
|
||||
}
|
||||
errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Source at {source_idx} input is empty")));
|
||||
errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Source at {source_idx} is empty")));
|
||||
} else {
|
||||
if log_enabled!(Level::Debug) {
|
||||
debug!("Input has {} groups", all_playlist.len());
|
||||
debug!("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(&mut all_playlist, target, &cfg, &mut stats, &mut errors).await {
|
||||
match process_playlist(&mut source_playlists, target, &cfg, &mut stats, &mut errors).await {
|
||||
Ok(()) => {}
|
||||
Err(mut err) => err.drain(..).for_each(|e| errors.push(e))
|
||||
}
|
||||
@@ -357,6 +358,51 @@ fn get_processing_pipe(target: &ConfigTarget) -> ProcessingPipe {
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
fn execute_pipe<'a>(target: &ConfigTarget, pipe: &ProcessingPipe, fpl: &mut FetchedPlaylist<'a>) -> 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(),
|
||||
};
|
||||
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(mut playlistgroups: Vec<PlaylistGroup>) -> Vec<PlaylistGroup> {
|
||||
let mut group_map: HashMap<(Rc<String>, XtreamCluster), PlaylistGroup> = HashMap::new();
|
||||
let mut sort_order = vec![];
|
||||
playlistgroups.drain(..).for_each(|group| {
|
||||
let key = (Rc::clone(&group.title), group.xtream_cluster);
|
||||
match group_map.entry(key) {
|
||||
std::collections::hash_map::Entry::Occupied(o) => {
|
||||
// we loose the group id (category id) at this point, which should be available in
|
||||
// the playlist item header.
|
||||
o.into_mut().channels.extend(group.channels);
|
||||
}
|
||||
std::collections::hash_map::Entry::Vacant(v) => {
|
||||
sort_order.push(Rc::clone(&group.title));
|
||||
v.insert(group);
|
||||
},
|
||||
};
|
||||
});
|
||||
let mut flat_groups: Vec<PlaylistGroup> = group_map.into_values().collect();
|
||||
// apply the initial sort order
|
||||
let mut sort_iterator = sort_order.iter();
|
||||
flat_groups.sort_by(|f, s| {
|
||||
let i1 = sort_iterator.position(|r| **r == *f.title).unwrap();
|
||||
let i2 = sort_iterator.position(|r| **r == *s.title).unwrap();
|
||||
i1.cmp(&i2)
|
||||
});
|
||||
flat_groups
|
||||
}
|
||||
|
||||
async fn process_playlist<'a>(playlists: &mut [FetchedPlaylist<'a>],
|
||||
target: &ConfigTarget, cfg: &Config,
|
||||
stats: &mut HashMap<u16, InputStats>,
|
||||
@@ -368,28 +414,16 @@ async fn process_playlist<'a>(playlists: &mut [FetchedPlaylist<'a>],
|
||||
|
||||
let mut new_fetched_playlists: Vec<FetchedPlaylist> = vec![];
|
||||
for fpl in playlists.iter_mut() {
|
||||
let mut new_fpl = FetchedPlaylist {
|
||||
input: fpl.input,
|
||||
playlist: fpl.playlist.clone(), // we need to clone, because of multiple target definitions, we cant change the initial playlist.
|
||||
epg: fpl.epg.clone(),
|
||||
};
|
||||
for f in &pipe {
|
||||
let playlist = &mut new_fpl.playlist;
|
||||
let r = f(playlist, target);
|
||||
if let Some(v) = r {
|
||||
new_fpl.playlist = v;
|
||||
}
|
||||
}
|
||||
let mut new_fpl = execute_pipe(target, &pipe, fpl);
|
||||
playlist_resolve_series(target, errors, &pipe, fpl, &mut new_fpl).await;
|
||||
// stats
|
||||
let input_stats = stats.get_mut(&new_fpl.input.id);
|
||||
if let Some(stat) = input_stats {
|
||||
stat.processed_stats.group_count = new_fpl.playlist.len();
|
||||
stat.processed_stats.channel_count = new_fpl.playlist.iter()
|
||||
stat.processed_stats.group_count = new_fpl.playlistgroups.len();
|
||||
stat.processed_stats.channel_count = new_fpl.playlistgroups.iter()
|
||||
.map(|group| group.channels.len())
|
||||
.sum();
|
||||
}
|
||||
|
||||
new_fetched_playlists.push(new_fpl);
|
||||
}
|
||||
|
||||
@@ -399,9 +433,12 @@ async fn process_playlist<'a>(playlists: &mut [FetchedPlaylist<'a>],
|
||||
let mut new_epg = vec![];
|
||||
|
||||
new_fetched_playlists.drain(..).for_each(|mut fp| {
|
||||
let epg_channel_ids: HashSet<_> = fp.playlist.iter().flat_map(|g| &g.channels)
|
||||
let epg_channel_ids: HashSet<_> = fp.playlistgroups.iter().flat_map(|g| &g.channels)
|
||||
.filter_map(|c| c.header.borrow().epg_channel_id.clone()).collect();
|
||||
fp.playlist.drain(..).for_each(|group| new_playlist.push(group));
|
||||
|
||||
fp.playlistgroups.drain(..).for_each(|group| {
|
||||
new_playlist.push(group);
|
||||
});
|
||||
if !epg_channel_ids.is_empty() {
|
||||
if let Some(tv_guide) = fp.epg {
|
||||
debug!("found epg information for {}", &target.name);
|
||||
@@ -419,21 +456,24 @@ async fn process_playlist<'a>(playlists: &mut [FetchedPlaylist<'a>],
|
||||
Ok(())
|
||||
} else {
|
||||
sort_playlist(target, &mut new_playlist);
|
||||
process_watch(target, cfg, &new_playlist);
|
||||
let mut flat_new_playlist = flatten_groups(new_playlist);
|
||||
persist_playlist(&mut flat_new_playlist, flatten_tvguide(&new_epg).as_ref(), target, cfg)
|
||||
}
|
||||
}
|
||||
|
||||
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);
|
||||
}
|
||||
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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
persist_playlist(&mut new_playlist, flatten_tvguide(&new_epg).as_ref(), target, cfg)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user