config refactor

This commit is contained in:
euzu
2025-06-24 11:25:42 +02:00
parent f7acbe3bca
commit c292f0b73e
145 changed files with 4954 additions and 4061 deletions
-59
View File
@@ -1,59 +0,0 @@
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);
});
}
}
}
}
+6 -3
View File
@@ -1,12 +1,12 @@
use crate::model::{Epg, TVGuide, XmlTag, XmlTagIcon, EPG_ATTRIB_ID};
use crate::model::{EpgConfig, EpgSmartMatchConfig};
use crate::model::{FetchedPlaylist, PlaylistItem};
use crate::model::{FetchedPlaylist};
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;
use shared::model::{EpgSmartMatchConfigDto, PlaylistItem, XtreamCluster};
pub struct EpgIdCache<'a> {
pub channel_epg_id: HashSet<Cow<'a, str>>,
@@ -31,7 +31,10 @@ impl EpgIdCache<'_> {
/// 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());
let normalize_config: EpgSmartMatchConfig = epg_config
.and_then(|cfg| cfg.smart_match.clone())
.unwrap_or_else(|| EpgSmartMatchConfigDto::default().into());
EpgIdCache {
channel_epg_id: HashSet::new(), // contains the epg_ids collected from playlist channels
normalized: HashMap::new(),
+54 -48
View File
@@ -1,4 +1,4 @@
use crate::model::{ConfigInput, ConfigRename};
use crate::model::{AppConfig, ConfigInput, ConfigRename};
use crate::utils::epg;
use crate::utils::m3u;
use crate::utils::xtream;
@@ -9,30 +9,30 @@ 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};
use crate::model::{ConfigTarget, InputType, ProcessTargets};
use crate::model::{CounterModifier, Mapping};
use crate::model::{FetchedPlaylist, PlaylistGroup, PlaylistItem};
use shared::model::{FieldGetAccessor, FieldSetAccessor, ItemField, MsgKind, PlaylistEntry, ProcessingOrder, UUIDType, XtreamCluster};
use crate::messaging::send_message;
use crate::model::Epg;
use crate::model::{ConfigTarget, ProcessTargets};
use crate::model::{Mapping};
use crate::model::FetchedPlaylist;
use crate::model::{InputStats, PlaylistStats, SourceStats, TargetStats};
use crate::processing::parser::xmltv::flatten_tvguide;
use crate::processing::playlist_watch::process_group_watch;
use crate::processing::processor::xtream_series::playlist_resolve_series;
use crate::processing::processor::epg::process_playlist_epg;
use crate::processing::processor::sort::sort_playlist;
use crate::processing::processor::trakt::process_trakt_categories_for_target;
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 shared::error::{get_errors_notify_message, notify_err, TuliproxError, TuliproxErrorKind};
use crate::utils::debug_if_enabled;
use shared::utils::default_as_default;
use crate::utils::StepMeasure;
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;
use shared::error::{get_errors_notify_message, notify_err, TuliproxError, TuliproxErrorKind};
use shared::foundation::filter::{get_field_value, set_field_value, ValueAccessor, ValueProvider};
use shared::model::{CounterModifier, FieldGetAccessor, FieldSetAccessor, InputType, ItemField, MsgKind, PlaylistEntry, PlaylistGroup, PlaylistItem, ProcessingOrder, UUIDType, XtreamCluster};
use shared::utils::default_as_default;
use std::time::Instant;
fn is_valid(pli: &PlaylistItem, target: &ConfigTarget) -> bool {
let provider = ValueProvider { pli };
@@ -85,7 +85,7 @@ fn exec_rename(pli: &mut PlaylistItem, rename: Option<&Vec<ConfigRename>>) {
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);
let cap = r.pattern.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);
}
@@ -105,7 +105,7 @@ fn rename_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Opt
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);
let cap = r.pattern.replace_all(&grp.title, &r.new_name);
debug_if_enabled!("Renamed group {} to {} for {}", &grp.title, cap, target.name);
grp.title = cap.into_owned();
}
@@ -147,7 +147,7 @@ fn map_channel(mut channel: PlaylistItem, mapping: &Mapping) -> PlaylistItem {
}
fn map_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option<Vec<PlaylistGroup>> {
if let Some(mappings) = target.t_mapping.load().as_ref() {
if let Some(mappings) = target.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()))
@@ -184,10 +184,9 @@ fn map_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option
}
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(guard) = &*target.mapping.load() {
let mappings = guard.as_ref();
for mapping in mappings {
if let Some(counter_list) = &mapping.t_counter {
for counter in counter_list {
for plg in &mut *playlist {
@@ -233,8 +232,9 @@ fn is_target_enabled(target: &ConfigTarget, user_targets: &ProcessTargets) -> bo
(!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();
async fn process_source(client: Arc<reqwest::Client>, cfg: Arc<AppConfig>, source_idx: usize, user_targets: Arc<ProcessTargets>) -> (Vec<InputStats>, Vec<TargetStats>, Vec<TuliproxError>) {
let sources = cfg.sources.load();
let source = 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();
@@ -243,15 +243,18 @@ async fn process_source(client: Arc<reqwest::Client>, cfg: Arc<Config>, source_i
let mut source_downloaded = false;
for input in &source.inputs {
if is_input_enabled(input, &user_targets) {
let config = cfg.config.load();
let working_dir = &config.working_dir;
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::M3u => m3u::get_m3u_playlist(Arc::clone(&client), &config, input, working_dir).await,
InputType::Xtream => xtream::get_xtream_playlist(&config, Arc::clone(&client), input, 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
epg::get_xmltv(Arc::clone(&client), input, working_dir).await
} else {
(None, vec![])
};
@@ -288,7 +291,7 @@ async fn process_source(client: Arc<reqwest::Client>, cfg: Arc<Config>, source_i
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 {
match process_playlist_for_target(&cfg, Arc::clone(&client), &mut source_playlists, target, &mut input_stats, &mut errors).await {
Ok(()) => {
target_stats.push(TargetStats::success(&target.name));
}
@@ -321,16 +324,17 @@ fn create_input_stat(group_count: usize, channel_count: usize, error_count: usiz
}
}
async fn process_sources(client: Arc<reqwest::Client>, config: Arc<Config>, user_targets: Arc<ProcessTargets>) -> (Vec<SourceStats>, Vec<TuliproxError>) {
async fn process_sources(client: Arc<reqwest::Client>, config: &Arc<AppConfig>, 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;
let thread_num = config.config.load().threads;
let sources = config.sources.load();
let process_parallel = thread_num > 1 && 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() {
for (index, _) in 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 {
@@ -433,10 +437,10 @@ fn flatten_groups(playlistgroups: Vec<PlaylistGroup>) -> Vec<PlaylistGroup> {
sort_order
}
async fn process_playlist_for_target(client: Arc<reqwest::Client>,
async fn process_playlist_for_target(app_config: &AppConfig,
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);
@@ -450,8 +454,8 @@ async fn process_playlist_for_target(client: Arc<reqwest::Client>,
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;
playlist_resolve_series(app_config, Arc::clone(&client), target, errors, &pipe, provider_fpl, &mut processed_fpl).await;
playlist_resolve_vod(app_config, Arc::clone(&client), target, errors, &mut processed_fpl).await;
// stats
let input_stats = stats.get_mut(&processed_fpl.input.name);
if let Some(stat) = input_stats {
@@ -486,9 +490,10 @@ async fn process_playlist_for_target(client: Arc<reqwest::Client>,
map_playlist_counter(target, &mut flat_new_playlist);
step.tick("Processed group watches");
process_watch(&client, target, cfg, &flat_new_playlist);
let config = app_config.config.load();
process_watch(&config, &client, target, &flat_new_playlist);
step.tick("Persisting playlists");
let result = persist_playlist(&mut flat_new_playlist, flatten_tvguide(&new_epg).as_ref(), target, cfg).await;
let result = persist_playlist(app_config, &mut flat_new_playlist, flatten_tvguide(&new_epg).as_ref(), target).await;
step.stop();
result
}
@@ -522,14 +527,13 @@ fn process_epg(processed_fetched_playlists: &mut Vec<FetchedPlaylist>) -> (Vec<E
(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() {
fn process_watch(cfg: &Config, client: &Arc<reqwest::Client>, target: &ConfigTarget, new_playlist: &Vec<PlaylistGroup>) {
if let Some(watches) = &target.watch {
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)) {
if watches.iter().any(|r| r.is_match(&pl.title)) {
process_group_watch(client, cfg, &target.name, pl);
}
}
@@ -537,23 +541,25 @@ fn process_watch(client: &Arc<reqwest::Client>, target: &ConfigTarget, cfg: &Con
}
}
pub async fn exec_processing(client: Arc<reqwest::Client>, cfg: Arc<Config>, targets: Arc<ProcessTargets>) {
pub async fn exec_processing(client: Arc<reqwest::Client>, cfg: Arc<AppConfig>, targets: Arc<ProcessTargets>) {
let start_time = Instant::now();
let (stats, errors) = process_sources(Arc::clone(&client), cfg.clone(), targets.clone()).await;
let (stats, errors) = process_sources(Arc::clone(&client), &cfg, targets.clone()).await;
// log errors
for err in &errors {
error!("{}", err.message);
}
let config = cfg.config.load();
let messaging = config.messaging.as_ref();
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_message(&client, &MsgKind::Stats, messaging, 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());
send_message(&client, &MsgKind::Error, messaging, error_msg.as_str());
}
}
let elapsed = start_time.elapsed().as_secs();
+12 -16
View File
@@ -1,8 +1,8 @@
use crate::foundation::filter::get_field_value;
use crate::model::{ConfigSortChannel, ConfigSortGroup, ConfigTarget, SortOrder};
use crate::model::{PlaylistGroup, PlaylistItem};
use crate::model::{ConfigSortChannel, ConfigSortGroup, ConfigTarget};
use deunicode::deunicode;
use std::cmp::Ordering;
use shared::foundation::filter::get_field_value;
use shared::model::{PlaylistGroup, PlaylistItem, SortOrder};
fn playlist_comparator(
sequence: Option<&Vec<regex::Regex>>,
@@ -98,7 +98,7 @@ fn playlistgroup_comparator(a: &PlaylistGroup, b: &PlaylistGroup, group_sort: &C
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)
playlist_comparator(group_sort.sequence.as_ref(), group_sort.order, &value_a, &value_b)
}
fn playlistitem_comparator(
@@ -112,7 +112,7 @@ fn playlistitem_comparator(
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)
playlist_comparator(channel_sort.sequence.as_ref(), channel_sort.order, &value_a, &value_b)
}
pub(in crate::processing::processor) fn sort_playlist(target: &ConfigTarget, new_playlist: &mut [PlaylistGroup]) {
@@ -123,7 +123,7 @@ pub(in crate::processing::processor) fn sort_playlist(target: &ConfigTarget, new
}
if let Some(channel_sorts) = &sort.channels {
for channel_sort in channel_sorts {
let regexp = channel_sort.t_re_group_pattern.as_ref().unwrap();
let regexp = &channel_sort.group_pattern;
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()) {
@@ -137,10 +137,10 @@ pub(in crate::processing::processor) fn sort_playlist(target: &ConfigTarget, new
#[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;
use shared::model::{ItemField, PlaylistItem, PlaylistItemHeader, SortOrder};
use crate::model::ConfigSortChannel;
#[test]
fn test_sort() {
@@ -151,15 +151,13 @@ mod tests {
let channel_sort = ConfigSortChannel {
field: ItemField::Caption,
group_pattern: ".*".to_string(),
order: SortOrder::Asc,
sequence: None,
t_re_sequence: Some(vec![
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()),
group_pattern: Regex::new(".*").unwrap(),
};
channels.sort_by(|chan1, chan2| playlistitem_comparator(chan1, chan2, &channel_sort, true));
@@ -199,10 +197,8 @@ mod tests {
let channel_sort = ConfigSortChannel {
field: ItemField::Caption,
group_pattern: ".*US.*".to_string(),
order: SortOrder::Asc,
sequence: None,
t_re_sequence: Some(vec![
sequence: Some(vec![
Regex::new(r"^US\| EAST.*?\[\bUHD\b\](?P<c1>.*)").unwrap(),
Regex::new(r"^US\| EAST.*?\[\bFHD\b\](?P<c1>.*)").unwrap(),
Regex::new(r"^US\| EAST.*?\[\bHD\b\](?P<c1>.*)").unwrap(),
@@ -212,7 +208,7 @@ mod tests {
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()),
group_pattern: Regex::new(".*").unwrap(),
};
channels.sort_by(|chan1, chan2| playlistitem_comparator(chan1, chan2, &channel_sort, true));
+9 -9
View File
@@ -1,10 +1,10 @@
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::model::{ConfigTarget, TraktListItem, TraktMatchItem};
use shared::model::{FieldGetAccessor, FieldSetAccessor, PlaylistGroup, PlaylistItem, TraktContentType, XtreamCluster};
use crate::model::{TraktConfig, TraktListConfig, TraktMatchResult};
use shared::error::TuliproxError;
use crate::utils::{TraktClient, extract_year_from_title, normalize_title_for_matching};
use crate::utils::{get_u32_from_serde_value};
use shared::utils::{get_u32_from_serde_value};
use shared::utils::{CONSTANTS};
use crate::utils::{trace_if_enabled, with};
use log::{debug, info, trace, warn};
@@ -22,7 +22,7 @@ fn extract_quality(value: &str) -> Option<&str> {
/// Utility functions for content type compatibility
fn should_include_item(item: &TraktListItem, content_type: &TraktContentType) -> bool {
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,
@@ -30,7 +30,7 @@ fn should_include_item(item: &TraktListItem, content_type: &TraktContentType) ->
}
}
fn is_compatible_content_type(cluster: XtreamCluster, content_type: &TraktContentType) -> bool {
fn is_compatible_content_type(cluster: XtreamCluster, content_type: TraktContentType) -> bool {
match content_type {
TraktContentType::Vod => cluster == XtreamCluster::Video,
TraktContentType::Series => cluster == XtreamCluster::Series,
@@ -216,7 +216,7 @@ fn match_trakt_items_with_playlist<'a>(
) -> Option<PlaylistGroup> {
let trakt_match_items: Vec<TraktMatchItem<'a>> = trakt_items
.iter()
.filter(|item| should_include_item(item, &list_config.content_type))
.filter(|item| should_include_item(item, list_config.content_type))
.filter_map(TraktMatchItem::from_trakt_list_item)
.collect();
@@ -225,7 +225,7 @@ fn match_trakt_items_with_playlist<'a>(
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) {
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);
@@ -314,7 +314,7 @@ mod tests {
#[test]
pub fn test_quality() {
let quality = extract_quality("Hello HD UHD 720p");
assert_eq!(true, quality.is_some());
assert!(quality.is_some());
assert_eq!("UHD", quality.unwrap());
}
}
+5 -5
View File
@@ -1,8 +1,8 @@
use shared::error::{info_err, notify_err};
use shared::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::{AppConfig, Config, ConfigInput};
use crate::model::{FetchedPlaylist};
use shared::model::{PlaylistEntry, PlaylistItem, PlaylistItemType, XtreamCluster};
use crate::model::normalize_release_date;
use crate::repository::storage::get_input_storage_path;
use serde::{Deserialize, Serialize};
@@ -87,7 +87,7 @@ pub(in crate::processing) fn should_update_info(pli: &mut PlaylistItem, processe
|| *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<'_>,
pub(in crate::processing) async fn read_processed_info_ids<V, F>(cfg: &AppConfig, errors: &mut Vec<TuliproxError>, fpl: &FetchedPlaylist<'_>,
item_type: PlaylistItemType, extract_ts: F) -> HashMap<u32, u64>
where
F: Fn(&V) -> u64,
@@ -96,7 +96,7 @@ where
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)
let file_path = match get_input_storage_path(fpl_name, &cfg.config.load().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,
@@ -1,7 +1,8 @@
use shared::model::InputType;
use shared::error::{TuliproxError, TuliproxErrorKind};
use crate::model::{Config, ConfigTarget, InputType};
use crate::model::{FetchedPlaylist, PlaylistGroup, PlaylistItem};
use shared::model::{PlaylistItemType, XtreamCluster};
use crate::model::{AppConfig, ConfigTarget};
use crate::model::{FetchedPlaylist};
use shared::model::{PlaylistGroup, PlaylistItem, 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};
@@ -23,7 +24,7 @@ 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> {
async fn read_processed_series_info_ids(cfg: &AppConfig, errors: &mut Vec<TuliproxError>, fpl: &FetchedPlaylist<'_>) -> HashMap<u32, u64> {
read_processed_info_ids(cfg, errors, fpl, PlaylistItemType::SeriesInfo, |ts: &u64| *ts).await
}
@@ -46,13 +47,13 @@ fn should_update_series_info(pli: &mut PlaylistItem, processed_provider_ids: &Ha
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>,
async fn playlist_resolve_series_info(cfg: &AppConfig, client: Arc<reqwest::Client>, 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)
let Some((wal_content_file, wal_record_file, wal_content_path, wal_record_path)) = create_resolve_info_wal_files(&cfg.config.load(), fpl.input, XtreamCluster::Series)
else { return !processed_info_ids.is_empty(); };
let mut content_writer = utils::file_writer(&wal_content_file);
@@ -124,26 +125,26 @@ async fn playlist_resolve_series_info(client: Arc<reqwest::Client>, cfg: &Config
!processed_info_ids.is_empty()
}
async fn process_series_info(
cfg: &Config,
app_config: &AppConfig,
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)
let config = app_config.config.load();
let Ok(Some((info_path, idx_path))) = get_input_storage_path(&input.name, &config.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);
let _file_lock = app_config.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 {
let Some((wal_file, wal_path)) = create_resolve_episode_wal_files(&config, input) else {
errors.push(notify_err!("Could not create wal file for series episodes record".to_string()));
return result;
};
@@ -201,13 +202,15 @@ async fn process_series_info(
|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,
handle_error!(xtream_update_input_series_episodes_record_from_wal_file(app_config, input, &wal_path).await,
|err| errors.push(err));
result
}
pub async fn playlist_resolve_series(client: Arc<reqwest::Client>, cfg: &Config, target: &ConfigTarget,
pub async fn playlist_resolve_series(cfg: &AppConfig,
client: Arc<reqwest::Client>,
target: &ConfigTarget,
errors: &mut Vec<TuliproxError>,
pipe: &ProcessingPipe,
provider_fpl: &mut FetchedPlaylist<'_>,
@@ -216,7 +219,7 @@ pub async fn playlist_resolve_series(client: Arc<reqwest::Client>, cfg: &Config,
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; }
if !playlist_resolve_series_info(cfg, client, 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
+13 -11
View File
@@ -1,12 +1,13 @@
use shared::model::InputType;
use shared::error::{TuliproxError, TuliproxErrorKind};
use crate::model::{Config, ConfigTarget, InputType};
use crate::model::{FetchedPlaylist, PlaylistItem};
use shared::model::{PlaylistItemType, XtreamCluster};
use crate::model::{AppConfig, ConfigTarget};
use crate::model::{FetchedPlaylist};
use shared::model::{PlaylistItem, 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 shared::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 shared::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;
@@ -19,7 +20,7 @@ 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> {
async fn read_processed_vod_info_ids(cfg: &AppConfig, errors: &mut Vec<TuliproxError>, fpl: &FetchedPlaylist<'_>) -> HashMap<u32, u64> {
read_processed_info_ids(cfg, errors, fpl, PlaylistItemType::Video, |record: &InputVodInfoRecord| record.ts).await
}
@@ -58,17 +59,18 @@ fn should_update_vod_info(pli: &mut PlaylistItem, processed_provider_ids: &HashM
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<'_>) {
pub async fn playlist_resolve_vod(app_config: &AppConfig, client: Arc<reqwest::Client>, 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)
let config = app_config.config.load();
let Some((wal_content_file, wal_record_file, wal_content_path, wal_record_path)) = create_resolve_info_wal_files(&config, fpl.input, XtreamCluster::Video)
else { return; };
let mut processed_info_ids = read_processed_vod_info_ids(cfg, errors, fpl).await;
let mut processed_info_ids = read_processed_vod_info_ids(app_config, 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;
@@ -126,9 +128,9 @@ pub async fn playlist_resolve_vod(client: Arc<reqwest::Client>, cfg: &Config, ta
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,
handle_error!(xtream_update_input_info_file(app_config, 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,
handle_error!(xtream_update_input_vod_record_from_wal_file(app_config, fpl.input, &wal_record_path).await,
|err| errors.push(err));
}
@@ -140,7 +142,7 @@ pub async fn playlist_resolve_vod(client: Arc<reqwest::Client>, cfg: &Config, ta
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) {
if let Some(content) = xtream_get_input_info(app_config, 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());
}
}