mirror of
https://github.com/euzu/tuliprox.git
synced 2026-10-04 23:12:27 +02:00
Merge branch 'feature/multi_xtream_input' into feature/multi_xtream
This commit is contained in:
@@ -5,109 +5,54 @@ use std::collections::{HashMap, HashSet};
|
||||
use std::rc::Rc;
|
||||
use std::sync::{Arc, Mutex};
|
||||
use std::thread;
|
||||
use actix_rt::System;
|
||||
|
||||
use actix_rt::System;
|
||||
use log::{debug, error, info, Level, log_enabled};
|
||||
use unidecode::unidecode;
|
||||
|
||||
use crate::{Config, get_errors_notify_message, model::config, valid_property};
|
||||
use crate::{Config, get_errors_notify_message, model::config};
|
||||
use crate::filter::{get_field_value, MockValueProcessor, set_field_value, ValueProvider};
|
||||
use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind};
|
||||
use crate::messaging::{MsgKind, send_message};
|
||||
use crate::model::config::{ConfigTarget, default_as_default, InputAffix, InputType, ProcessTargets};
|
||||
use crate::model::config::{ConfigTarget, default_as_default, InputType, ProcessTargets};
|
||||
use crate::model::config::{ItemField, ProcessingOrder, SortOrder::{Asc, Desc}, TargetType};
|
||||
use crate::model::mapping::{Mapping, MappingValueProcessor};
|
||||
use crate::model::config::{AFFIX_FIELDS, ItemField, ProcessingOrder, SortOrder::{Asc, Desc}, TargetType};
|
||||
use crate::model::playlist::{FetchedPlaylist, FieldAccessor, PlaylistGroup, PlaylistItem, PlaylistItemHeader};
|
||||
|
||||
use crate::model::playlist::{FetchedPlaylist, PlaylistGroup, PlaylistItem};
|
||||
use crate::model::stats::{InputStats, PlaylistStats};
|
||||
use crate::model::xmltv::{Epg};
|
||||
use crate::model::xtream::MultiXtreamMapping;
|
||||
use crate::processing::playlist_watch::process_group_watch;
|
||||
use crate::processing::xmltv_parser::flatten_tvguide;
|
||||
use crate::repository::epg_repository::write_epg;
|
||||
use crate::repository::m3u_repository::{write_m3u_playlist, write_strm_playlist};
|
||||
use crate::repository::xtream_repository::write_xtream_playlist;
|
||||
use crate::processing::xtream_processor::playlist_resolve_series;
|
||||
use crate::repository::playlist_repository::persist_playlist;
|
||||
use crate::repository::xtream_repository::write_xtream_mapping;
|
||||
use crate::utils::download;
|
||||
use crate::processing::affix_processor::apply_affixes;
|
||||
|
||||
fn is_valid(pli: &PlaylistItem, target: &ConfigTarget) -> bool {
|
||||
let provider = ValueProvider { pli: RefCell::new(pli) };
|
||||
target.filter(&provider)
|
||||
}
|
||||
|
||||
fn filter_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option<Vec<PlaylistGroup>> {
|
||||
debug!("Filtering {} groups", playlist.len());
|
||||
let mut new_playlist = Vec::new();
|
||||
playlist.iter_mut().for_each(|pg| {
|
||||
if log_enabled!(Level::Debug) {
|
||||
debug!("Filtering group {} with {} items", pg.title, pg.channels.len());
|
||||
}
|
||||
let mut channels = Vec::new();
|
||||
pg.channels.iter_mut().for_each(|pli| {
|
||||
if is_valid(pli, target) {
|
||||
channels.push(pli.clone());
|
||||
}
|
||||
});
|
||||
if log_enabled!(Level::Debug) {
|
||||
debug!("Filtered group {} has now {} items", pg.title, channels.len());
|
||||
}
|
||||
let channels = pg.channels.iter()
|
||||
.filter(|&pli| is_valid(pli, target)).cloned().collect::<Vec<PlaylistItem>>();
|
||||
debug!("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.clone()
|
||||
xtream_cluster: pg.xtream_cluster.clone(),
|
||||
});
|
||||
}
|
||||
});
|
||||
Some(new_playlist)
|
||||
}
|
||||
|
||||
fn apply_affixes(fetched_playlists: &mut [FetchedPlaylist]) {
|
||||
fetched_playlists.iter_mut().for_each(|fetched_playlist| {
|
||||
let FetchedPlaylist { input, playlist, epg: _ } = fetched_playlist;
|
||||
if input.suffix.is_some() || input.prefix.is_some() {
|
||||
let validate_affix = |a: &Option<InputAffix>| match a {
|
||||
Some(affix) => {
|
||||
valid_property!(&affix.field.as_str(), AFFIX_FIELDS) && !affix.value.is_empty()
|
||||
}
|
||||
_ => false
|
||||
};
|
||||
|
||||
let apply_prefix = validate_affix(&input.prefix);
|
||||
let apply_suffix = validate_affix(&input.suffix);
|
||||
|
||||
if apply_prefix || apply_suffix {
|
||||
let get_affix_applied_value = |header: &mut PlaylistItemHeader, affix: &InputAffix, prefix: bool| {
|
||||
if let Some(field_value) = header.get_field(affix.field.as_str()) {
|
||||
return if prefix {
|
||||
format!("{}{}", &affix.value, field_value.as_str())
|
||||
} else {
|
||||
format!("{}{}", field_value.as_str(), &affix.value)
|
||||
};
|
||||
}
|
||||
String::from(&affix.value)
|
||||
};
|
||||
|
||||
playlist.iter_mut().for_each(|group| {
|
||||
group.channels.iter_mut().for_each(|channel| {
|
||||
if apply_suffix {
|
||||
if let Some(suffix) = &input.suffix {
|
||||
let value = get_affix_applied_value(&mut channel.header.borrow_mut(), suffix, false);
|
||||
if log_enabled!(Level::Debug) {
|
||||
debug!("Applying input suffix: {}={}", &suffix.field, &value);
|
||||
}
|
||||
channel.header.borrow_mut().set_field(&suffix.field, value.as_str());
|
||||
}
|
||||
}
|
||||
if apply_prefix {
|
||||
if let Some(prefix) = &input.prefix {
|
||||
let value = get_affix_applied_value(&mut channel.header.borrow_mut(), prefix, true);
|
||||
if log_enabled!(Level::Debug) {
|
||||
debug!("Applying input prefix: {}={}", &prefix.field, &value);
|
||||
}
|
||||
channel.header.borrow_mut().set_field(&prefix.field, value.as_str());
|
||||
}
|
||||
}
|
||||
});
|
||||
});
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
fn sort_playlist(target: &ConfigTarget, new_playlist: &mut [PlaylistGroup]) {
|
||||
if let Some(sort) = &target.sort {
|
||||
let match_as_ascii = &sort.match_as_ascii;
|
||||
@@ -146,12 +91,6 @@ fn sort_playlist(target: &ConfigTarget, new_playlist: &mut [PlaylistGroup]) {
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
fn is_valid(pli: &mut PlaylistItem, target: &ConfigTarget) -> bool {
|
||||
let provider = ValueProvider { pli: RefCell::new(pli) };
|
||||
target.filter(&provider)
|
||||
}
|
||||
|
||||
fn exec_rename(pli: &mut PlaylistItem, rename: &Option<Vec<config::ConfigRename>>) {
|
||||
if let Some(renames) = rename {
|
||||
if !renames.is_empty() {
|
||||
@@ -236,46 +175,33 @@ fn map_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option
|
||||
let new_playlist: Vec<PlaylistGroup> = playlist.iter().map(|playlist_group| {
|
||||
let mut grp = playlist_group.clone();
|
||||
let mappings = target._mapping.as_ref().unwrap();
|
||||
mappings.iter().filter(|mapping| !mapping.mapper.is_empty()).for_each(|mapping|
|
||||
mappings.iter().filter(|&mapping| !mapping.mapper.is_empty()).for_each(|mapping|
|
||||
grp.channels = grp.channels.drain(..).map(|chan| map_channel(chan, mapping)).collect());
|
||||
grp
|
||||
}).collect();
|
||||
|
||||
// if the group names are changed, restructure channels to the right groups
|
||||
// we use
|
||||
let mut max_group_id = 0;
|
||||
let mut new_groups: Vec<PlaylistGroup> = Vec::new();
|
||||
let mut grp_id: u32 = 0;
|
||||
for playlist_group in new_playlist {
|
||||
let mut group_id_used = false;
|
||||
for channel in &playlist_group.channels {
|
||||
let cluster = &channel.header.borrow().xtream_cluster;
|
||||
let title = &channel.header.borrow().group;
|
||||
match new_groups.iter_mut().find(|x| *x.title == **title) {
|
||||
Some(grp) => grp.channels.push(channel.clone()),
|
||||
_ => {
|
||||
let new_group_id = if group_id_used {
|
||||
0
|
||||
} else if *title == playlist_group.title {
|
||||
group_id_used = true;
|
||||
max_group_id = max_group_id.max(playlist_group.id);
|
||||
playlist_group.id
|
||||
} else {
|
||||
0
|
||||
};
|
||||
grp_id += 1;
|
||||
new_groups.push(PlaylistGroup {
|
||||
id: new_group_id,
|
||||
id: grp_id,
|
||||
title: Rc::clone(title),
|
||||
channels: vec![channel.clone()],
|
||||
xtream_cluster: cluster.clone()
|
||||
xtream_cluster: cluster.clone(),
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
new_groups.iter_mut().filter(|g| g.id == 0).for_each(|grp| {
|
||||
max_group_id += 1;
|
||||
grp.id = max_group_id;
|
||||
});
|
||||
Some(new_groups)
|
||||
} else {
|
||||
None
|
||||
@@ -372,7 +298,7 @@ async fn process_source(cfg: Arc<Config>, source_idx: usize, user_targets: Arc<P
|
||||
(stats.drain().map(|(_, v)| v).collect(), errors)
|
||||
}
|
||||
|
||||
pub(crate) async fn process_sources(config: Arc<Config>, user_targets: Arc<ProcessTargets>) -> (Vec<InputStats>, Vec<M3uFilterError>) {
|
||||
async fn process_sources(config: Arc<Config>, user_targets: Arc<ProcessTargets>) -> (Vec<InputStats>, Vec<M3uFilterError>) {
|
||||
let mut handle_list = vec![];
|
||||
let thread_num = config.threads;
|
||||
let process_parallel = thread_num > 1 && config.sources.len() > 1;
|
||||
@@ -415,8 +341,7 @@ pub(crate) async fn process_sources(config: Arc<Config>, user_targets: Arc<Proce
|
||||
(Arc::try_unwrap(stats).unwrap().into_inner().unwrap(), Arc::try_unwrap(errors).unwrap().into_inner().unwrap())
|
||||
}
|
||||
|
||||
|
||||
type ProcessingPipe = Vec<fn(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option<Vec<PlaylistGroup>>>;
|
||||
pub type ProcessingPipe = Vec<fn(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option<Vec<PlaylistGroup>>>;
|
||||
|
||||
fn get_processing_pipe(target: &ConfigTarget) -> ProcessingPipe {
|
||||
match &target.processing_order {
|
||||
@@ -429,10 +354,10 @@ fn get_processing_pipe(target: &ConfigTarget) -> ProcessingPipe {
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) async fn process_playlist<'a>(playlists: &mut [FetchedPlaylist<'a>],
|
||||
target: &ConfigTarget, cfg: &Config,
|
||||
stats: &mut HashMap<u16, InputStats>,
|
||||
errors: &mut Vec<M3uFilterError>) -> Result<(), Vec<M3uFilterError>> {
|
||||
async fn process_playlist<'a>(playlists: &mut [FetchedPlaylist<'a>],
|
||||
target: &ConfigTarget, cfg: &Config,
|
||||
stats: &mut HashMap<u16, InputStats>,
|
||||
errors: &mut Vec<M3uFilterError>) -> Result<(), Vec<M3uFilterError>> {
|
||||
let pipe = get_processing_pipe(target);
|
||||
if log_enabled!(Level::Debug) {
|
||||
debug!("Processing order is {}", &target.processing_order);
|
||||
@@ -452,32 +377,7 @@ pub(crate) async fn process_playlist<'a>(playlists: &mut [FetchedPlaylist<'a>],
|
||||
new_fpl.playlist = v;
|
||||
}
|
||||
}
|
||||
let (resolve_series, resolve_series_delay) =
|
||||
if let Some(options) = &target.options {
|
||||
(options.xtream_resolve_series && fpl.input.input_type == InputType::Xtream && target.has_output(&TargetType::M3u),
|
||||
options.xtream_resolve_series_delay)
|
||||
} else {
|
||||
(false, 0)
|
||||
};
|
||||
if resolve_series {
|
||||
let mut series_playlist = download::get_xtream_playlist_series(fpl, errors, resolve_series_delay).await;
|
||||
// original content saved into original list
|
||||
for plg in &series_playlist {
|
||||
fpl.update_playlist(plg);
|
||||
}
|
||||
// run processing pipe over new items
|
||||
for f in &pipe {
|
||||
let r = f(&mut series_playlist, target);
|
||||
if let Some(v) = r {
|
||||
series_playlist = v;
|
||||
}
|
||||
}
|
||||
// assign new items to the new playlist
|
||||
for plg in &series_playlist {
|
||||
new_fpl.update_playlist(plg);
|
||||
}
|
||||
}
|
||||
|
||||
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 {
|
||||
@@ -491,6 +391,49 @@ pub(crate) async fn process_playlist<'a>(playlists: &mut [FetchedPlaylist<'a>],
|
||||
}
|
||||
|
||||
apply_affixes(&mut new_fetched_playlists);
|
||||
|
||||
if target.is_multi_input() && target.has_output(&TargetType::Xtream) {
|
||||
let mut stream_id_mappings: Vec<MultiXtreamMapping> = Vec::new();
|
||||
let mut counter: u32 = 0;
|
||||
new_fetched_playlists.iter()
|
||||
.flat_map(|pl| {
|
||||
let input_id = &pl.input.id;
|
||||
pl.playlist.iter().map(move |plg| (input_id, &plg.channels))
|
||||
})
|
||||
.flat_map(|(input_id, channels)| channels.iter().map(move |chan| (input_id, chan)))
|
||||
.for_each(|(input_id, chan)| {
|
||||
let mut header = chan.header.borrow_mut();
|
||||
if header.stream_id.is_empty() {
|
||||
header.stream_id = Rc::clone(&header.id);
|
||||
}
|
||||
match header.stream_id.parse::<u32>() {
|
||||
Ok(stream_id) => {
|
||||
let xtream_mapping = MultiXtreamMapping {
|
||||
stream_id,
|
||||
input_id: *input_id,
|
||||
};
|
||||
stream_id_mappings.push(xtream_mapping);
|
||||
}
|
||||
Err(_) => {
|
||||
error!("Failed to parse stream_id: {}", &header.id)
|
||||
}
|
||||
}
|
||||
counter += 1;
|
||||
header.id = Rc::new(counter.to_string());
|
||||
});
|
||||
|
||||
match write_xtream_mapping(&stream_id_mappings, cfg, &target.name) {
|
||||
Ok(_) => {
|
||||
debug!("wrote multi xtream input mapping for {}", &target.name);
|
||||
}
|
||||
Err(err) => {
|
||||
return Err(vec![M3uFilterError::new(
|
||||
M3uFilterErrorKind::Notify,
|
||||
format!("Write multi xtream input mapping {} failed: {}", target.name, err))]);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
let mut new_playlist = vec![];
|
||||
let mut new_epg = vec![];
|
||||
let mut tv_guides = vec![];
|
||||
@@ -524,7 +467,7 @@ pub(crate) async fn process_playlist<'a>(playlists: &mut [FetchedPlaylist<'a>],
|
||||
let watch_re = target._watch_re.as_ref().unwrap();
|
||||
new_playlist.iter().for_each(|pl| {
|
||||
if watch_re.iter().any(|r| r.is_match(&pl.title)) {
|
||||
process_group_watch(cfg, &target.name, pl)
|
||||
process_group_watch(cfg, &target.name, pl)
|
||||
}
|
||||
});
|
||||
}
|
||||
@@ -537,30 +480,6 @@ pub(crate) async fn process_playlist<'a>(playlists: &mut [FetchedPlaylist<'a>],
|
||||
}
|
||||
}
|
||||
|
||||
fn persist_playlist(playlist: &mut [PlaylistGroup], epg: Option<Epg>,
|
||||
target: &ConfigTarget, cfg: &Config) -> Result<(), Vec<M3uFilterError>> {
|
||||
let mut errors = vec![];
|
||||
for output in &target.output {
|
||||
match match output.target {
|
||||
TargetType::M3u => write_m3u_playlist(target, cfg, playlist, &output.filename),
|
||||
TargetType::Strm => write_strm_playlist(target, cfg, playlist, &output.filename),
|
||||
TargetType::Xtream => write_xtream_playlist(target, cfg, playlist)
|
||||
} {
|
||||
Ok(_) => {
|
||||
if !playlist.is_empty() {
|
||||
match write_epg(target, cfg, &epg, output) {
|
||||
Ok(_) => {}
|
||||
Err(err) => errors.push(err)
|
||||
}
|
||||
}
|
||||
}
|
||||
Err(err) => errors.push(err)
|
||||
}
|
||||
}
|
||||
|
||||
if errors.is_empty() { Ok(()) } else { Err(errors) }
|
||||
}
|
||||
|
||||
pub(crate) async fn exec_processing(cfg: Arc<Config>, targets: Arc<ProcessTargets>) {
|
||||
let (stats, errors) = process_sources(cfg.to_owned(), targets.to_owned()).await;
|
||||
let stats_msg = format!("{{\"stats\": {}}}", stats.iter().map(|stat| stat.to_string()).collect::<Vec<String>>().join("\n"));
|
||||
@@ -572,7 +491,7 @@ pub(crate) async fn exec_processing(cfg: Arc<Config>, targets: Arc<ProcessTarget
|
||||
errors.iter().for_each(|err| error!("{}", err.message));
|
||||
// send errors
|
||||
if let Some(message) = get_errors_notify_message!(errors, 255) {
|
||||
let error_msg = format!("{{\"errors\": \"{}\"}}",message.as_str());
|
||||
let error_msg = format!("{{\"errors\": \"{}\"}}", message.as_str());
|
||||
send_message(&MsgKind::Error, &cfg.messaging, error_msg.as_str());
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user