mirror of
https://github.com/euzu/tuliprox.git
synced 2026-10-08 00:42:32 +02:00
Staged filter and safe yaml update (#846)
New Features
Added separate processing and persistence filters, with empty/non-empty checks and expanded filter fields.
Added persistence-stage filter controls in target settings.
Added clear_invalid_epg_ids, preserving playlist entries while clearing unresolved EPG IDs; the legacy setting remains supported.
Added safer, atomic source configuration updates that preserve comments and formatting.
Bug Fixes
Null-like Xtream and Stalker values are now handled correctly.
Duplicate credential warnings provide clearer details.
This commit is contained in:
@@ -2,7 +2,7 @@ use crate::{fetched_playlist::FetchedPlaylist, parser::xmltv::normalize_channel_
|
||||
use log::{debug, trace, warn};
|
||||
use rphonetic::{DoubleMetaphone, Encoder};
|
||||
use shared::{
|
||||
model::{EpgNamePrefix, EpgSmartMatchConfigDto, PlaylistItem, XtreamCluster},
|
||||
model::{EpgNamePrefix, EpgSmartMatchConfigDto, PlaylistGroup, PlaylistItem, XtreamCluster},
|
||||
utils::{Internable, CONSTANTS},
|
||||
};
|
||||
use std::{
|
||||
@@ -743,13 +743,16 @@ fn assign_live_channel_epg(
|
||||
icon_override_channels: &HashSet<Arc<str>>,
|
||||
icon_assigned: &mut HashSet<Arc<str>>,
|
||||
stats: &mut EpgAssignmentStats,
|
||||
) -> bool {
|
||||
clear_invalid_epg_ids: bool,
|
||||
) {
|
||||
if id_cache.smart_match_enabled {
|
||||
stats.record(assign_smart_epg_id(channel, id_cache));
|
||||
}
|
||||
let has_epg = has_processed_epg(channel, id_cache);
|
||||
assign_epg_icon(channel, icon_tags, icon_override_channels, icon_assigned);
|
||||
has_epg
|
||||
if clear_invalid_epg_ids && !has_epg {
|
||||
channel.header.epg_channel_id = None;
|
||||
}
|
||||
}
|
||||
|
||||
fn referenced_live_epg_ids(fp: &mut FetchedPlaylist<'_>) -> HashSet<Arc<str>> {
|
||||
@@ -761,13 +764,27 @@ fn referenced_live_epg_ids(fp: &mut FetchedPlaylist<'_>) -> HashSet<Arc<str>> {
|
||||
.collect()
|
||||
}
|
||||
|
||||
pub(crate) fn retain_live_items_with_processed_epg(fp: &mut FetchedPlaylist<'_>, epg: &[Epg]) {
|
||||
pub(crate) fn retain_epg_referenced_by_groups(groups: &[PlaylistGroup], epg: &mut [Epg]) {
|
||||
let referenced_epg_ids = groups
|
||||
.iter()
|
||||
.flat_map(|group| &group.channels)
|
||||
.filter(|channel| is_live_epg_item(channel))
|
||||
.filter_map(|channel| {
|
||||
channel.header.epg_channel_id.as_ref().map(|id| with_folded_epg_id(id, |folded| folded.intern()))
|
||||
})
|
||||
.collect::<HashSet<_>>();
|
||||
for source in epg {
|
||||
source.children.retain(|channel| with_folded_epg_id(&channel.id, |folded| referenced_epg_ids.contains(folded)));
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn clear_invalid_live_epg_ids(fp: &mut FetchedPlaylist<'_>, epg: &[Epg]) {
|
||||
let processed_epg_ids = epg
|
||||
.iter()
|
||||
.flat_map(|source| &source.children)
|
||||
.map(|channel| with_folded_epg_id(&channel.id, |folded| folded.intern()))
|
||||
.collect::<HashSet<_>>();
|
||||
let mut removed = 0usize;
|
||||
let mut cleared = 0usize;
|
||||
fp.source.retain_memory_items_mut(|item| {
|
||||
if !is_live_epg_item(item) {
|
||||
return true;
|
||||
@@ -777,18 +794,18 @@ pub(crate) fn retain_live_items_with_processed_epg(fp: &mut FetchedPlaylist<'_>,
|
||||
.epg_channel_id
|
||||
.as_deref()
|
||||
.is_some_and(|id| with_folded_epg_id(id, |folded| processed_epg_ids.contains(folded)));
|
||||
if !has_epg {
|
||||
removed += 1;
|
||||
if !has_epg && item.header.epg_channel_id.take().is_some() {
|
||||
cleared += 1;
|
||||
}
|
||||
has_epg
|
||||
true
|
||||
});
|
||||
if removed > 0 {
|
||||
debug!("Removed {removed} live channels invalidated by after-EPG mappings from input '{}'", fp.input.name);
|
||||
if cleared > 0 {
|
||||
debug!("Cleared {cleared} unmatched EPG IDs from input '{}'", fp.input.name);
|
||||
}
|
||||
}
|
||||
|
||||
/// Assigns EPG IDs and logos to live playlist channels by matching them with EPG data.
|
||||
/// When EPG data is required by the target, unmatched live entries are removed in the same pass.
|
||||
/// Invalid EPG IDs are optionally cleared without removing playlist entries.
|
||||
///
|
||||
/// For each live channel in the playlist missing an EPG ID, attempts to assign one using normalized name matching if smart matching is enabled. If a channel has an EPG ID but lacks logos, assigns logos from the corresponding EPG icon tags. Adds the matched EPG data to the provided vector.
|
||||
///
|
||||
@@ -804,18 +821,14 @@ async fn assign_channel_epg(
|
||||
new_epg: &mut Vec<Epg>,
|
||||
fp: &mut FetchedPlaylist<'_>,
|
||||
id_cache: &mut EpgIdCache,
|
||||
required_epg: bool,
|
||||
clear_invalid_epg_ids: bool,
|
||||
) {
|
||||
let Some(tv_guide) = &fp.epg else {
|
||||
return;
|
||||
};
|
||||
let mut merged_epg = tv_guide.filter_merged_with_icon_overrides(id_cache).await;
|
||||
if merged_epg.is_none() && !required_epg {
|
||||
return;
|
||||
}
|
||||
|
||||
let mut stats = EpgAssignmentStats::default();
|
||||
let mut removed = 0usize;
|
||||
if fp.is_memory() {
|
||||
let icon_tags = merged_epg
|
||||
.as_ref()
|
||||
@@ -838,30 +851,20 @@ async fn assign_channel_epg(
|
||||
|
||||
let mut process_channel = |channel: &mut PlaylistItem| {
|
||||
if !is_live_epg_item(channel) {
|
||||
return true;
|
||||
return;
|
||||
}
|
||||
let has_epg = assign_live_channel_epg(
|
||||
assign_live_channel_epg(
|
||||
channel,
|
||||
id_cache,
|
||||
&icon_tags,
|
||||
&icon_override_channels,
|
||||
&mut icon_assigned,
|
||||
&mut stats,
|
||||
clear_invalid_epg_ids,
|
||||
);
|
||||
if required_epg && !has_epg {
|
||||
removed += 1;
|
||||
return false;
|
||||
}
|
||||
true
|
||||
};
|
||||
|
||||
if required_epg {
|
||||
fp.source.retain_memory_items_mut(&mut process_channel);
|
||||
} else {
|
||||
fp.items_mut().for_each(|channel| {
|
||||
process_channel(channel);
|
||||
});
|
||||
}
|
||||
fp.items_mut().for_each(&mut process_channel);
|
||||
} else {
|
||||
warn!("Disk based playlist modification is not supported!");
|
||||
}
|
||||
@@ -872,10 +875,6 @@ async fn assign_channel_epg(
|
||||
fp.input.name, stats.live, stats.existing, stats.exact, stats.fuzzy, stats.corrected, stats.unresolved
|
||||
);
|
||||
}
|
||||
if removed > 0 {
|
||||
debug!("Removed {removed} live channels without EPG from input '{}'", fp.input.name);
|
||||
}
|
||||
|
||||
if let Some((mut epg_source, _)) = merged_epg.take() {
|
||||
let referenced_epg_ids = referenced_live_epg_ids(fp);
|
||||
epg_source
|
||||
@@ -896,7 +895,7 @@ async fn assign_channel_epg(
|
||||
/// let mut epg_data = Vec::new();
|
||||
/// process_playlist_epg(&mut playlist, &mut epg_data, false);
|
||||
/// ```
|
||||
pub async fn process_playlist_epg(fp: &mut FetchedPlaylist<'_>, epg: &mut Vec<Epg>, required_epg: bool) {
|
||||
pub async fn process_playlist_epg(fp: &mut FetchedPlaylist<'_>, epg: &mut Vec<Epg>, clear_invalid_epg_ids: bool) {
|
||||
if fp.input.epg.is_none() {
|
||||
return;
|
||||
}
|
||||
@@ -904,10 +903,10 @@ pub async fn process_playlist_epg(fp: &mut FetchedPlaylist<'_>, epg: &mut Vec<Ep
|
||||
let mut id_cache = EpgIdCache::new(fp.input.epg.as_ref());
|
||||
id_cache.collect_epg_id(fp);
|
||||
|
||||
if id_cache.is_empty() && !id_cache.smart_match_enabled && !required_epg {
|
||||
if id_cache.is_empty() && !id_cache.smart_match_enabled && !clear_invalid_epg_ids {
|
||||
debug!("No epg ids found for input {}", fp.input.name);
|
||||
} else {
|
||||
assign_channel_epg(epg, fp, &mut id_cache, required_epg).await;
|
||||
assign_channel_epg(epg, fp, &mut id_cache, clear_invalid_epg_ids).await;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -920,7 +919,11 @@ mod tests {
|
||||
model::{ConfigInputDto, PlaylistGroup, PlaylistItem, PlaylistItemHeader, PlaylistItemType},
|
||||
utils::Internable,
|
||||
};
|
||||
use std::{collections::HashSet, fs, sync::Arc};
|
||||
use std::{
|
||||
collections::{HashMap, HashSet},
|
||||
fs,
|
||||
sync::Arc,
|
||||
};
|
||||
use tempfile::tempdir;
|
||||
use tokio::time::Instant;
|
||||
use tuliprox_core::model::{
|
||||
@@ -1009,7 +1012,7 @@ mod tests {
|
||||
xmltv: &str,
|
||||
channels: &[(&str, Option<&str>)],
|
||||
smart_matching: bool,
|
||||
required_epg: bool,
|
||||
clear_invalid_epg_ids: bool,
|
||||
) -> (Vec<Option<Arc<str>>>, Vec<tuliprox_core::model::Epg>) {
|
||||
let dir = tempdir().unwrap();
|
||||
let epg_path = dir.path().join("smart-match.xml");
|
||||
@@ -1045,7 +1048,7 @@ mod tests {
|
||||
};
|
||||
let mut epg = Vec::new();
|
||||
|
||||
super::process_playlist_epg(&mut playlist, &mut epg, required_epg).await;
|
||||
super::process_playlist_epg(&mut playlist, &mut epg, clear_invalid_epg_ids).await;
|
||||
let assigned_ids = playlist.items_mut().map(|item| item.header.epg_channel_id.clone()).collect();
|
||||
(assigned_ids, epg)
|
||||
}
|
||||
@@ -1270,7 +1273,7 @@ mod tests {
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn required_epg_removes_only_unmatched_live_items() {
|
||||
fn clear_invalid_epg_ids_preserves_items_and_clears_only_unmatched_live_ids() {
|
||||
let runtime = tokio::runtime::Runtime::new().unwrap();
|
||||
runtime.block_on(async move {
|
||||
let dir = tempdir().unwrap();
|
||||
@@ -1334,18 +1337,42 @@ mod tests {
|
||||
|
||||
super::process_playlist_epg(&mut playlist, &mut epg, true).await;
|
||||
|
||||
let names = playlist.items_mut().map(|item| item.header.name.clone()).collect::<HashSet<_>>();
|
||||
assert!(names.contains("Matched Live"));
|
||||
assert!(!names.contains("Unmatched Live"));
|
||||
assert!(names.contains("VOD"));
|
||||
assert!(names.contains("Series"));
|
||||
assert!(names.contains("Local VOD"));
|
||||
assert!(names.contains("Local Series"));
|
||||
let items = playlist
|
||||
.items_mut()
|
||||
.map(|item| (item.header.name.clone(), item.header.epg_channel_id.clone()))
|
||||
.collect::<HashMap<_, _>>();
|
||||
assert_eq!(items.get("Matched Live").and_then(Option::as_deref), Some("matched.live"));
|
||||
assert_eq!(items.get("Unmatched Live"), Some(&None));
|
||||
assert!(items.contains_key("VOD"));
|
||||
assert!(items.contains_key("Series"));
|
||||
assert!(items.contains_key("Local VOD"));
|
||||
assert!(items.contains_key("Local Series"));
|
||||
});
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn required_epg_is_ignored_without_a_materialized_epg_source() {
|
||||
fn default_epg_processing_preserves_unmatched_existing_id() {
|
||||
let runtime = tokio::runtime::Runtime::new().unwrap();
|
||||
runtime.block_on(async move {
|
||||
let (assigned_ids, _) = run_xmltv_matches(
|
||||
r#"<tv>
|
||||
<channel id="matched.live"><display-name>Matched Live</display-name></channel>
|
||||
<programme start="20260425000000 +0000" stop="20260425010000 +0000" channel="matched.live">
|
||||
<title>Programme</title>
|
||||
</programme>
|
||||
</tv>"#,
|
||||
&[("Unmatched Live", Some("missing.live"))],
|
||||
false,
|
||||
false,
|
||||
)
|
||||
.await;
|
||||
|
||||
assert_eq!(assigned_ids, vec![Some("missing.live".intern())]);
|
||||
});
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn clear_invalid_epg_ids_is_ignored_without_a_materialized_epg_source() {
|
||||
let runtime = tokio::runtime::Runtime::new().unwrap();
|
||||
runtime.block_on(async move {
|
||||
let mut input = ConfigInput::from(ConfigInputDto::default());
|
||||
@@ -1362,6 +1389,10 @@ mod tests {
|
||||
super::process_playlist_epg(&mut playlist, &mut Vec::new(), true).await;
|
||||
|
||||
assert_eq!(playlist.items_mut().count(), 1);
|
||||
assert_eq!(
|
||||
playlist.items_mut().next().and_then(|item| item.header.epg_channel_id.as_deref()),
|
||||
Some("missing.live")
|
||||
);
|
||||
});
|
||||
}
|
||||
|
||||
@@ -1417,7 +1448,7 @@ mod tests {
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn required_epg_removes_live_items_without_ids_when_smart_matching_is_disabled() {
|
||||
fn clear_invalid_epg_ids_preserves_live_items_without_ids() {
|
||||
let runtime = tokio::runtime::Runtime::new().unwrap();
|
||||
runtime.block_on(async move {
|
||||
let dir = tempdir().unwrap();
|
||||
@@ -1446,13 +1477,13 @@ mod tests {
|
||||
|
||||
super::process_playlist_epg(&mut playlist, &mut Vec::new(), true).await;
|
||||
|
||||
assert_eq!(playlist.items_mut().count(), 0);
|
||||
assert_eq!(playlist.get_group_count(), 0);
|
||||
assert_eq!(playlist.items_mut().count(), 1);
|
||||
assert_eq!(playlist.get_group_count(), 1);
|
||||
});
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn required_epg_keeps_live_items_assigned_by_smart_matching() {
|
||||
fn clear_invalid_epg_ids_keeps_ids_assigned_by_smart_matching() {
|
||||
let runtime = tokio::runtime::Runtime::new().unwrap();
|
||||
runtime.block_on(async move {
|
||||
let (assigned_ids, epg) = run_xmltv_matches(
|
||||
|
||||
@@ -7,7 +7,7 @@ use crate::{
|
||||
parser::xmltv::{flatten_tvguide, merge_epg_trees, EpgMergeAccumulator, TVGuide},
|
||||
playlist_watch::{process_group_watch, process_target_groups_watch},
|
||||
processor::{
|
||||
epg::{process_playlist_epg, retain_live_items_with_processed_epg},
|
||||
epg::{clear_invalid_live_epg_ids, process_playlist_epg, retain_epg_referenced_by_groups},
|
||||
sort::sort_playlist,
|
||||
trakt::process_trakt_categories_for_target,
|
||||
xtream_series::playlist_resolve_series,
|
||||
@@ -45,9 +45,10 @@ use tokio::{
|
||||
};
|
||||
use tuliprox_core::{
|
||||
model::{
|
||||
is_valid, AppConfig, CompiledMapping, ConfigFavourites, ConfigInput, ConfigInputFlags, ConfigInputOptions,
|
||||
ConfigRename, ConfigTarget, Epg, MappingProgram, ProcessTargets, ProviderIdType, ResolveReason,
|
||||
ReverseProxyDisabledHeaderConfig, TransformStage, UpdateGuard, UpdateTask,
|
||||
is_valid, retain_filtered_playlist, AppConfig, CompiledMapping, ConfigFavourites, ConfigInput,
|
||||
ConfigInputFlags, ConfigInputOptions, ConfigRename, ConfigTarget, Epg, FilterOutcome, MappingProgram,
|
||||
ProcessTargets, ProviderIdType, ResolveReason, ReverseProxyDisabledHeaderConfig, TransformStage, UpdateGuard,
|
||||
UpdateTask,
|
||||
},
|
||||
utils::{debug_if_enabled, log_memory_snapshot, trace_if_enabled, StepMeasure, StepMeasureCallback},
|
||||
};
|
||||
@@ -116,13 +117,6 @@ fn stalker_checkpoint_message(input: &str) -> String {
|
||||
format!("Input '{input}': Stalker refresh checkpoint saved; active snapshot remains in service")
|
||||
}
|
||||
|
||||
#[derive(Debug, Default, PartialEq, Eq)]
|
||||
pub struct FilterOutcome {
|
||||
pub inspected: usize,
|
||||
pub retained: usize,
|
||||
pub removed: usize,
|
||||
}
|
||||
|
||||
fn retain_playlist_items(
|
||||
source: &mut PlaylistSource,
|
||||
mut keep: impl FnMut(&PlaylistItem) -> bool,
|
||||
@@ -130,9 +124,7 @@ fn retain_playlist_items(
|
||||
let mut groups: IndexMap<CategoryKey, PlaylistGroup> = IndexMap::new();
|
||||
let mut outcome = FilterOutcome::default();
|
||||
for pli in source.into_items() {
|
||||
outcome.inspected += 1;
|
||||
if keep(&pli) {
|
||||
outcome.retained += 1;
|
||||
if outcome.record(keep(&pli)) {
|
||||
let group_title = pli.header.group.clone();
|
||||
let cluster = pli.header.xtream_cluster;
|
||||
let cat_id = pli.header.category_id;
|
||||
@@ -148,8 +140,6 @@ fn retain_playlist_items(
|
||||
})
|
||||
.channels
|
||||
.push(pli);
|
||||
} else {
|
||||
outcome.removed += 1;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1471,20 +1461,14 @@ impl TransformBuffer {
|
||||
|
||||
fn apply_filter(&mut self, target: &ConfigTarget) -> FilterOutcome {
|
||||
let mut outcome = FilterOutcome::default();
|
||||
self.items.retain(|item| {
|
||||
outcome.inspected += 1;
|
||||
let provider = ValueProvider { pli: item, match_as_ascii: false };
|
||||
if target.filter(&provider) {
|
||||
outcome.retained += 1;
|
||||
true
|
||||
} else {
|
||||
outcome.removed += 1;
|
||||
false
|
||||
}
|
||||
});
|
||||
self.items.retain(|item| outcome.record(target.filter(&ValueProvider { pli: item, match_as_ascii: false })));
|
||||
self.normalize_filter_grouping();
|
||||
outcome
|
||||
}
|
||||
|
||||
fn normalize_filter_grouping(&mut self) {
|
||||
self.grouping = GroupingPolicy::NormalizedCategory;
|
||||
self.reorder_for_grouping();
|
||||
outcome
|
||||
}
|
||||
|
||||
fn apply_rename(&mut self, target: &ConfigTarget) -> Option<RenameOutcome> {
|
||||
@@ -1567,7 +1551,13 @@ fn execute_pipeline_on_items(
|
||||
let mut outcome = PipelineOutcome::default();
|
||||
for stage in pipe {
|
||||
match stage {
|
||||
TransformStage::Filter => outcome.filter = Some(buffer.apply_filter(target)),
|
||||
TransformStage::Filter => {
|
||||
if target.filter.processing.is_some() {
|
||||
outcome.filter = Some(buffer.apply_filter(target));
|
||||
} else {
|
||||
buffer.normalize_filter_grouping();
|
||||
}
|
||||
}
|
||||
TransformStage::Rename => outcome.rename = buffer.apply_rename(target),
|
||||
TransformStage::Map => outcome.mapping = buffer.apply_mapping(target, MappingStage::Processing),
|
||||
}
|
||||
@@ -1575,6 +1565,14 @@ fn execute_pipeline_on_items(
|
||||
(buffer.into_groups(), outcome)
|
||||
}
|
||||
|
||||
fn apply_persist_filter(target: &ConfigTarget, groups: &mut Vec<PlaylistGroup>) {
|
||||
let Some(filter) = target.filter.persist.as_ref() else {
|
||||
return;
|
||||
};
|
||||
let outcome = retain_filtered_playlist(groups, filter);
|
||||
debug!("Target '{}' persist filter outcome: {outcome:?}", target.name);
|
||||
}
|
||||
|
||||
pub(super) fn execute_pipeline_on_groups(
|
||||
groups: Vec<PlaylistGroup>,
|
||||
target: &ConfigTarget,
|
||||
@@ -1694,9 +1692,9 @@ async fn prepare_playlist_for_target<E: EventSink + Clone + 'static, M: Metadata
|
||||
log_memory_snapshot(
|
||||
format!("target '{}' input '{}' after_vod_resolve", target.name, provider_fpl.input.name).as_str(),
|
||||
);
|
||||
let required_epg = target.options.as_ref().is_some_and(ConfigTargetOptions::required_epg);
|
||||
let clear_invalid_epg_ids = target.options.as_ref().is_some_and(ConfigTargetOptions::clear_invalid_epg_ids);
|
||||
let input_epg_start = new_epg.len();
|
||||
process_playlist_epg(&mut processed_fpl, &mut new_epg, required_epg).await;
|
||||
process_playlist_epg(&mut processed_fpl, &mut new_epg, clear_invalid_epg_ids).await;
|
||||
log_memory_snapshot(
|
||||
format!("target '{}' input '{}' after_epg_apply", target.name, processed_fpl.input.name).as_str(),
|
||||
);
|
||||
@@ -1708,9 +1706,9 @@ async fn prepare_playlist_for_target<E: EventSink + Clone + 'static, M: Metadata
|
||||
deduplicate.then_some(&mut duplicates),
|
||||
) {
|
||||
processed_fpl.source = MemoryPlaylistSource::new(groups).into_source();
|
||||
if required_epg && new_epg.len() > input_epg_start {
|
||||
retain_live_items_with_processed_epg(&mut processed_fpl, &new_epg[input_epg_start..]);
|
||||
}
|
||||
}
|
||||
if clear_invalid_epg_ids && processed_fpl.epg.is_some() {
|
||||
clear_invalid_live_epg_ids(&mut processed_fpl, &new_epg[input_epg_start..]);
|
||||
}
|
||||
if let Some(stat) = stats.get_mut(&processed_fpl.input.name) {
|
||||
stat.processed_stats.group_count = processed_fpl.get_group_count();
|
||||
@@ -1778,7 +1776,7 @@ async fn finalize_prepared_target<E: EventSink + Clone + 'static, M: MetadataUpd
|
||||
) -> (Result<(), Vec<TuliproxError>>, Vec<TuliproxError>) {
|
||||
let target = &prepared.target;
|
||||
let mut new_playlist = prepared.playlist;
|
||||
let new_epg = prepared.epg;
|
||||
let mut new_epg = prepared.epg;
|
||||
let mut errors = Vec::new();
|
||||
let broadcast_step = create_broadcast_callback(&ctx.events);
|
||||
let mut step = StepMeasure::new(&target.name, broadcast_step);
|
||||
@@ -1836,6 +1834,9 @@ async fn finalize_prepared_target<E: EventSink + Clone + 'static, M: MetadataUpd
|
||||
}
|
||||
}
|
||||
|
||||
apply_persist_filter(target, &mut flat_new_playlist);
|
||||
retain_epg_referenced_by_groups(&flat_new_playlist, &mut new_epg);
|
||||
|
||||
if process_watch(&ctx.config, &ctx.events, target, &flat_new_playlist).await {
|
||||
step.tick("group watches");
|
||||
log_memory_snapshot(format!("target '{}' after_group_watches", target.name).as_str());
|
||||
@@ -3102,7 +3103,7 @@ mod tests {
|
||||
xtream_cluster: XtreamCluster::Live,
|
||||
}];
|
||||
let mut target = ConfigTarget::from(&ConfigTargetDto::default());
|
||||
target.filter = get_filter(r#"name ~ "Allowed""#, None).expect("filter should parse");
|
||||
target.filter = get_filter(r#"name ~ "Allowed""#, None).expect("filter should parse").into();
|
||||
|
||||
let (groups, outcome) = execute_pipeline_on_groups(groups, &target, &[TransformStage::Filter]);
|
||||
|
||||
@@ -3110,6 +3111,47 @@ mod tests {
|
||||
assert_eq!(outcome.filter, Some(FilterOutcome { inspected: 1, retained: 0, removed: 1 }));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn missing_processing_filter_skips_filter_stage() {
|
||||
let groups = vec![PlaylistGroup {
|
||||
id: 1,
|
||||
title: "Test Group".intern(),
|
||||
channels: vec![make_test_item("Allowed", PlaylistItemType::Live)],
|
||||
xtream_cluster: XtreamCluster::Live,
|
||||
}];
|
||||
let target = ConfigTarget::from(&ConfigTargetDto::default());
|
||||
|
||||
let (groups, outcome) = execute_pipeline_on_groups(groups, &target, &[TransformStage::Filter]);
|
||||
|
||||
assert_eq!(groups.len(), 1);
|
||||
assert_eq!(groups[0].channels.len(), 1);
|
||||
assert!(outcome.filter.is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn missing_processing_filter_preserves_filter_stage_group_normalization() {
|
||||
let mut first = make_test_item("One", PlaylistItemType::Live);
|
||||
first.header.group = "News".intern();
|
||||
let mut second = make_test_item("Two", PlaylistItemType::Live);
|
||||
second.header.group = "news".intern();
|
||||
let groups = vec![
|
||||
PlaylistGroup { id: 1, title: "News".intern(), channels: vec![first], xtream_cluster: XtreamCluster::Live },
|
||||
PlaylistGroup {
|
||||
id: 2,
|
||||
title: "news".intern(),
|
||||
channels: vec![second],
|
||||
xtream_cluster: XtreamCluster::Live,
|
||||
},
|
||||
];
|
||||
let target = ConfigTarget::from(&ConfigTargetDto::default());
|
||||
|
||||
let (groups, outcome) = execute_pipeline_on_groups(groups, &target, &[TransformStage::Filter]);
|
||||
|
||||
assert_eq!(groups.len(), 1);
|
||||
assert_eq!(groups[0].channels.len(), 2);
|
||||
assert!(outcome.filter.is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pipeline_reports_filter_and_rename_outcomes() {
|
||||
let groups = vec![PlaylistGroup {
|
||||
@@ -3122,7 +3164,7 @@ mod tests {
|
||||
xtream_cluster: XtreamCluster::Live,
|
||||
}];
|
||||
let mut target = ConfigTarget::from(&ConfigTargetDto::default());
|
||||
target.filter = get_filter(r#"name ~ "Allowed""#, None).expect("filter should parse");
|
||||
target.filter = get_filter(r#"name ~ "Allowed""#, None).expect("filter should parse").into();
|
||||
target.rename = Some(vec![ConfigRename::from(&ConfigRenameDto {
|
||||
field: ItemField::Name,
|
||||
pattern: "Allowed".to_string(),
|
||||
@@ -3473,6 +3515,58 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn persist_filter_runs_after_after_epg_mapping() {
|
||||
let runtime = Runtime::new().expect("runtime");
|
||||
runtime.block_on(async {
|
||||
let mut input = ConfigInput::from(ConfigInputDto::default());
|
||||
input.name = "input".intern();
|
||||
let groups = vec![PlaylistGroup {
|
||||
id: 1,
|
||||
title: "Live".intern(),
|
||||
channels: vec![PlaylistItem {
|
||||
header: PlaylistItemHeader {
|
||||
name: "Before".intern(),
|
||||
group: "Live".intern(),
|
||||
xtream_cluster: XtreamCluster::Live,
|
||||
item_type: PlaylistItemType::Live,
|
||||
..Default::default()
|
||||
},
|
||||
}],
|
||||
xtream_cluster: XtreamCluster::Live,
|
||||
}];
|
||||
let mut playlist = FetchedPlaylist {
|
||||
input: &input,
|
||||
source: MemoryPlaylistSource::new(groups).into_source(),
|
||||
epg: None,
|
||||
};
|
||||
let rename = build_mapping("rename", MappingStage::AfterEpg, r#"@Name = "After""#);
|
||||
let mut target = build_target(vec![rename], false);
|
||||
target.filter.persist = Some(get_filter(r#"Name = "After""#, None).expect("filter parses"));
|
||||
let mut stats = HashMap::from([(
|
||||
Arc::clone(&input.name),
|
||||
create_input_stat(1, 1, 0, input.input_type, &input.name, 0),
|
||||
)]);
|
||||
let mut errors = Vec::new();
|
||||
|
||||
let mut prepared = prepare_playlist_for_target(
|
||||
&processing_context(),
|
||||
std::slice::from_mut(&mut playlist),
|
||||
&target,
|
||||
&mut stats,
|
||||
&mut errors,
|
||||
false,
|
||||
)
|
||||
.await
|
||||
.expect("target preparation");
|
||||
|
||||
assert!(errors.is_empty());
|
||||
apply_persist_filter(&target, &mut prepared.playlist);
|
||||
let item = &prepared.playlist[0].channels[0];
|
||||
assert_eq!(item.header.name.as_ref(), "After");
|
||||
});
|
||||
}
|
||||
|
||||
fn make_channel(name: &str) -> PlaylistItem {
|
||||
let mut item = PlaylistItem {
|
||||
header: PlaylistItemHeader {
|
||||
@@ -3624,7 +3718,7 @@ match {
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn required_epg_removes_ids_invalidated_by_after_epg_mapping() {
|
||||
fn clear_invalid_epg_ids_clears_ids_invalidated_by_after_epg_mapping() {
|
||||
let runtime = Runtime::new().expect("runtime");
|
||||
runtime.block_on(async {
|
||||
let dir = tempdir().expect("tempdir");
|
||||
@@ -3674,7 +3768,7 @@ match {
|
||||
build_mapping("rewrite", MappingStage::AfterEpg, r#"@epg_channel_id = "missing.epg""#);
|
||||
let add_virtual = build_mapping("virtual", MappingStage::AfterEpg, r#"add_favourite("Echo")"#);
|
||||
let mut target = build_target(vec![rewrite_epg, add_virtual], false);
|
||||
target.options = Some(ConfigTargetOptions { required_epg: true, ..Default::default() });
|
||||
target.options = Some(ConfigTargetOptions { clear_invalid_epg_ids: true, ..Default::default() });
|
||||
let mut stats = HashMap::from([(
|
||||
Arc::clone(&input.name),
|
||||
create_input_stat(1, 1, 0, input.input_type, &input.name, 0),
|
||||
@@ -3693,9 +3787,13 @@ match {
|
||||
.expect("target preparation");
|
||||
|
||||
assert!(errors.is_empty());
|
||||
assert!(prepared.playlist.is_empty());
|
||||
assert_eq!(stats[&input.name].processed_stats.group_count, 0);
|
||||
assert_eq!(stats[&input.name].processed_stats.channel_count, 0);
|
||||
assert!(!prepared.playlist.is_empty());
|
||||
assert!(prepared
|
||||
.playlist
|
||||
.iter()
|
||||
.flat_map(|group| &group.channels)
|
||||
.all(|channel| channel.header.epg_channel_id.is_none()));
|
||||
assert_eq!(stats[&input.name].processed_stats.channel_count, 2);
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
@@ -1,10 +1,13 @@
|
||||
use crate::fetched_playlist::FetchedPlaylist;
|
||||
use log::{debug, warn};
|
||||
use parking_lot::Mutex;
|
||||
use shared::{
|
||||
error::TuliproxError,
|
||||
model::{LiveStreamProperties, StreamProperties, XtreamCluster, XtreamPlaylistItem},
|
||||
model::{
|
||||
LiveStreamProperties, PlaylistEntry, PlaylistItemType, StreamProperties, XtreamCluster, XtreamPlaylistItem,
|
||||
},
|
||||
};
|
||||
use std::sync::Arc;
|
||||
use std::{collections::HashMap, sync::Arc};
|
||||
use tuliprox_core::{
|
||||
model::{AppConfig, ConfigInput, ConfigInputFlags, ProviderHandle, ProviderIdType},
|
||||
utils::{
|
||||
@@ -206,6 +209,55 @@ pub async fn update_live_stream_metadata(
|
||||
Ok(Some(properties))
|
||||
}
|
||||
|
||||
pub(crate) fn sync_resolved_xtream_properties<T: Clone>(
|
||||
provider_fpl: &mut FetchedPlaylist<'_>,
|
||||
processed_fpl: &mut FetchedPlaylist<'_>,
|
||||
cluster: XtreamCluster,
|
||||
item_type: PlaylistItemType,
|
||||
extract: impl Fn(&StreamProperties) -> Option<&T>,
|
||||
wrap: impl Fn(Box<T>) -> StreamProperties,
|
||||
) {
|
||||
let mut resolved_by_provider_id: HashMap<u32, T> = HashMap::new();
|
||||
|
||||
for pli in processed_fpl.items() {
|
||||
if pli.header.xtream_cluster != cluster || pli.header.item_type != item_type {
|
||||
continue;
|
||||
}
|
||||
|
||||
let Some(provider_id) = pli.get_provider_id() else {
|
||||
continue;
|
||||
};
|
||||
if provider_id == 0 {
|
||||
continue;
|
||||
}
|
||||
|
||||
if let Some(props) = pli.header.additional_properties.as_ref().and_then(&extract) {
|
||||
resolved_by_provider_id.entry(provider_id).or_insert_with(|| props.clone());
|
||||
}
|
||||
}
|
||||
|
||||
if resolved_by_provider_id.is_empty() {
|
||||
return;
|
||||
}
|
||||
|
||||
for source_pli in provider_fpl.items_mut() {
|
||||
if source_pli.header.xtream_cluster != cluster || source_pli.header.item_type != item_type {
|
||||
continue;
|
||||
}
|
||||
|
||||
let Some(provider_id) = source_pli.get_provider_id() else {
|
||||
continue;
|
||||
};
|
||||
if provider_id == 0 {
|
||||
continue;
|
||||
}
|
||||
|
||||
if let Some(resolved) = resolved_by_provider_id.get(&provider_id) {
|
||||
source_pli.header.additional_properties = Some(wrap(Box::new(resolved.clone())));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn apply_live_probe_success(
|
||||
properties: &mut LiveStreamProperties,
|
||||
raw_video: Option<serde_json::Value>,
|
||||
|
||||
@@ -22,7 +22,7 @@ use shared::{
|
||||
},
|
||||
};
|
||||
use std::{
|
||||
collections::{HashMap, HashSet},
|
||||
collections::HashSet,
|
||||
sync::{
|
||||
atomic::{AtomicBool, Ordering},
|
||||
Arc,
|
||||
@@ -148,47 +148,17 @@ async fn playlist_resolve_series_info<E: EventSink + Clone + 'static, M: Metadat
|
||||
}
|
||||
|
||||
fn sync_resolved_series_properties(provider_fpl: &mut FetchedPlaylist<'_>, processed_fpl: &mut FetchedPlaylist<'_>) {
|
||||
let mut resolved_series_by_provider_id: HashMap<u32, SeriesStreamProperties> = HashMap::new();
|
||||
|
||||
for pli in processed_fpl.items() {
|
||||
if pli.header.xtream_cluster != XtreamCluster::Series || pli.header.item_type != PlaylistItemType::SeriesInfo {
|
||||
continue;
|
||||
}
|
||||
|
||||
let Some(provider_id) = pli.get_provider_id() else {
|
||||
continue;
|
||||
};
|
||||
if provider_id == 0 {
|
||||
continue;
|
||||
}
|
||||
|
||||
if let Some(StreamProperties::Series(properties)) = pli.header.additional_properties.as_ref() {
|
||||
resolved_series_by_provider_id.entry(provider_id).or_insert_with(|| properties.as_ref().clone());
|
||||
}
|
||||
}
|
||||
|
||||
if resolved_series_by_provider_id.is_empty() {
|
||||
return;
|
||||
}
|
||||
|
||||
for source_pli in provider_fpl.items_mut() {
|
||||
if source_pli.header.xtream_cluster != XtreamCluster::Series
|
||||
|| source_pli.header.item_type != PlaylistItemType::SeriesInfo
|
||||
{
|
||||
continue;
|
||||
}
|
||||
|
||||
let Some(provider_id) = source_pli.get_provider_id() else {
|
||||
continue;
|
||||
};
|
||||
if provider_id == 0 {
|
||||
continue;
|
||||
}
|
||||
|
||||
if let Some(resolved) = resolved_series_by_provider_id.get(&provider_id) {
|
||||
source_pli.header.additional_properties = Some(StreamProperties::Series(Box::new(resolved.clone())));
|
||||
}
|
||||
}
|
||||
crate::processor::xtream::sync_resolved_xtream_properties(
|
||||
provider_fpl,
|
||||
processed_fpl,
|
||||
XtreamCluster::Series,
|
||||
PlaylistItemType::SeriesInfo,
|
||||
|props| match props {
|
||||
StreamProperties::Series(s) => Some(s.as_ref()),
|
||||
_ => None,
|
||||
},
|
||||
StreamProperties::Series,
|
||||
);
|
||||
}
|
||||
|
||||
fn queue_background_series_info<E: EventSink + Clone + 'static, M: MetadataUpdateSink>(
|
||||
|
||||
@@ -17,12 +17,12 @@ use shared::{
|
||||
error::TuliproxError,
|
||||
foundation::ValueProvider,
|
||||
model::{
|
||||
EventSink, MediaQuality, PlaylistEntry, PlaylistItem, PlaylistItemType, StreamProperties,
|
||||
VideoStreamDetailProperties, VideoStreamProperties, XtreamCluster, XtreamPlaylistItem, XtreamVideoInfo,
|
||||
EventSink, MediaQuality, PlaylistItem, PlaylistItemType, StreamProperties, VideoStreamDetailProperties,
|
||||
VideoStreamProperties, XtreamCluster, XtreamPlaylistItem, XtreamVideoInfo,
|
||||
},
|
||||
};
|
||||
use std::{
|
||||
collections::{HashMap, HashSet},
|
||||
collections::HashSet,
|
||||
sync::{
|
||||
atomic::{AtomicBool, Ordering},
|
||||
Arc,
|
||||
@@ -96,47 +96,17 @@ pub async fn playlist_resolve_vod<E: EventSink + Clone + 'static, M: MetadataUpd
|
||||
}
|
||||
|
||||
fn sync_resolved_vod_properties(provider_fpl: &mut FetchedPlaylist<'_>, processed_fpl: &mut FetchedPlaylist<'_>) {
|
||||
let mut resolved_vod_by_provider_id: HashMap<u32, VideoStreamProperties> = HashMap::new();
|
||||
|
||||
for pli in processed_fpl.items() {
|
||||
if pli.header.xtream_cluster != XtreamCluster::Video || pli.header.item_type != PlaylistItemType::Video {
|
||||
continue;
|
||||
}
|
||||
|
||||
let Some(provider_id) = pli.get_provider_id() else {
|
||||
continue;
|
||||
};
|
||||
if provider_id == 0 {
|
||||
continue;
|
||||
}
|
||||
|
||||
if let Some(StreamProperties::Video(properties)) = pli.header.additional_properties.as_ref() {
|
||||
resolved_vod_by_provider_id.entry(provider_id).or_insert_with(|| properties.as_ref().clone());
|
||||
}
|
||||
}
|
||||
|
||||
if resolved_vod_by_provider_id.is_empty() {
|
||||
return;
|
||||
}
|
||||
|
||||
for source_pli in provider_fpl.items_mut() {
|
||||
if source_pli.header.xtream_cluster != XtreamCluster::Video
|
||||
|| source_pli.header.item_type != PlaylistItemType::Video
|
||||
{
|
||||
continue;
|
||||
}
|
||||
|
||||
let Some(provider_id) = source_pli.get_provider_id() else {
|
||||
continue;
|
||||
};
|
||||
if provider_id == 0 {
|
||||
continue;
|
||||
}
|
||||
|
||||
if let Some(resolved) = resolved_vod_by_provider_id.get(&provider_id) {
|
||||
source_pli.header.additional_properties = Some(StreamProperties::Video(Box::new(resolved.clone())));
|
||||
}
|
||||
}
|
||||
crate::processor::xtream::sync_resolved_xtream_properties(
|
||||
provider_fpl,
|
||||
processed_fpl,
|
||||
XtreamCluster::Video,
|
||||
PlaylistItemType::Video,
|
||||
|props| match props {
|
||||
StreamProperties::Video(v) => Some(v.as_ref()),
|
||||
_ => None,
|
||||
},
|
||||
StreamProperties::Video,
|
||||
);
|
||||
}
|
||||
|
||||
#[allow(clippy::too_many_lines, clippy::too_many_arguments)]
|
||||
|
||||
Reference in New Issue
Block a user