From a683927ece19ddef50bb762b90d3f5a77b51010b Mon Sep 17 00:00:00 2001 From: euzu <33094714+euzu@users.noreply.github.com> Date: Sat, 29 Aug 2026 21:41:38 +0200 Subject: [PATCH] new attribute for target, required_epg to remove all live channels without epg, if at least one epg source is defined. (#841) New Features Added a target option to require valid EPG programme data for live entries. Entries without matching EPG data are filtered out when enabled; non-live entries remain unaffected. Added support for configuring this option in the target editor, with localized labels and explanations. Documentation Documented the new setting, including its default behavior and handling when no EPG source is available. --- CHANGELOG.md | 10 + backend/processing/src/processor/epg.rs | 415 +++++++++++++++--- backend/processing/src/processor/playlist.rs | 95 +++- backend/repository/src/playlist_source.rs | 51 +++ config/source.yml | 2 + docs/src/configuration/source.md | 3 + frontend/public/assets/i18n/ar.json | 2 + frontend/public/assets/i18n/en.json | 2 + frontend/public/assets/i18n/ru.json | 4 + .../components/source_editor/target_form.rs | 16 + shared/src/model/config/target.rs | 18 + 11 files changed, 545 insertions(+), 73 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 29bb811fd..50291744c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -106,6 +106,11 @@ ## 🌟 New Features +- **Targets can require usable EPG data for live channels**: setting `options.required_epg: true` keeps only live + entries whose EPG ID resolves to programme data from an available EPG source. The normal target filter and mappings + still run first, so EPG matching operates on the already reduced playlist. VOD, series, catch-up, and local-library + entries are unaffected, and inputs without a successfully materialized EPG source keep their live entries. + - **Ten events for the failures that used to be silent**: the registry described states nothing emitted, and several subsystems reported their start and their success but never their own failure. The taxonomy is now 42 events (up from 32) and gains two domains, `scheduled_task.*` and `notification.*`. @@ -1191,6 +1196,11 @@ ## ⚙️ New Settings +- **source.yml (target `options`)**: + - Added optional `required_epg` (`bool`, default `false`) to remove unmatched live entries after EPG matching. + The setting is evaluated independently for each target and does not affect VOD, series, catch-up, or local-library + entries. If an input has no successfully materialized EPG source, its live entries are left unchanged. + - **config.yml (`video.download.recording`)**: - Added `enabled` (`bool`, default `true`): master switch for the DVR. When `false` the REST routes answer `501 recording_disabled`, the rule scheduler and supervisors idle, the WebSocket serves no recording data, and the diff --git a/backend/processing/src/processor/epg.rs b/backend/processing/src/processor/epg.rs index 818561e95..48b5234c2 100644 --- a/backend/processing/src/processor/epg.rs +++ b/backend/processing/src/processor/epg.rs @@ -728,16 +728,67 @@ fn assign_epg_icon( }); } +fn is_live_epg_item(item: &PlaylistItem) -> bool { + item.header.xtream_cluster == XtreamCluster::Live && item.header.item_type.is_live() +} + +fn has_processed_epg(item: &PlaylistItem, id_cache: &EpgIdCache) -> bool { + item.header.epg_channel_id.as_deref().is_some_and(|id| id_cache.contains_processed_epg_id(id)) +} + +fn assign_live_channel_epg( + channel: &mut PlaylistItem, + id_cache: &EpgIdCache, + icon_tags: &HashMap, &Arc>, + icon_override_channels: &HashSet>, + icon_assigned: &mut HashSet>, + stats: &mut EpgAssignmentStats, +) -> 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 +} + fn referenced_live_epg_ids(fp: &mut FetchedPlaylist<'_>) -> HashSet> { fp.items() - .filter(|channel| channel.header.xtream_cluster == XtreamCluster::Live && channel.header.item_type.is_live()) + .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() } +pub(crate) fn retain_live_items_with_processed_epg(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::>(); + let mut removed = 0usize; + fp.source.retain_memory_items_mut(|item| { + if !is_live_epg_item(item) { + return true; + } + let has_epg = item + .header + .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; + } + has_epg + }); + if removed > 0 { + debug!("Removed {removed} live channels invalidated by after-EPG mappings 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. /// /// 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. /// @@ -747,68 +798,90 @@ fn referenced_live_epg_ids(fp: &mut FetchedPlaylist<'_>) -> HashSet> { /// let mut new_epg = Vec::new(); /// let mut playlist = FetchedPlaylist::default(); /// let mut id_cache = EpgIdCache::new(None); -/// assign_channel_epg(&mut new_epg, &mut playlist, &mut id_cache); +/// assign_channel_epg(&mut new_epg, &mut playlist, &mut id_cache, false); /// ``` -async fn assign_channel_epg(new_epg: &mut Vec, fp: &mut FetchedPlaylist<'_>, id_cache: &mut EpgIdCache) { - if let Some(tv_guide) = &fp.epg { - if let Some((mut epg_source, icon_override_channels)) = - tv_guide.filter_merged_with_icon_overrides(id_cache).await - { - let stats = { - let icon_tags = epg_source - .children - .iter() - .filter_map(|tag| { - tag.icon - .as_ref() - .filter(|icon| !icon.is_empty()) - .map(|icon| (with_folded_epg_id(&tag.id, |folded| folded.intern()), icon)) - }) - .collect::, &Arc>>(); - let icon_override_channels = icon_override_channels - .into_iter() - .map(|id| with_folded_epg_id(&id, |folded| folded.intern())) - .collect::>(); - let mut icon_assigned = HashSet::new(); - let mut stats = EpgAssignmentStats::default(); +async fn assign_channel_epg( + new_epg: &mut Vec, + fp: &mut FetchedPlaylist<'_>, + id_cache: &mut EpgIdCache, + required_epg: 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; + } - if fp.is_memory() { - fp.items_mut() - .filter(|channel| { - channel.header.xtream_cluster == XtreamCluster::Live && channel.header.item_type.is_live() - }) - .for_each(|channel| { - if id_cache.smart_match_enabled { - stats.record(assign_smart_epg_id(channel, id_cache)); - } - assign_epg_icon(channel, &icon_tags, &icon_override_channels, &mut icon_assigned); - }); - } else { - warn!("Disk based playlist modification is not supported!"); - } - stats - }; + let mut stats = EpgAssignmentStats::default(); + let mut removed = 0usize; + if fp.is_memory() { + let icon_tags = merged_epg + .as_ref() + .into_iter() + .flat_map(|(epg_source, _)| &epg_source.children) + .filter_map(|tag| { + tag.icon + .as_ref() + .filter(|icon| !icon.is_empty()) + .map(|icon| (with_folded_epg_id(&tag.id, |folded| folded.intern()), icon)) + }) + .collect::, &Arc>>(); + let icon_override_channels = merged_epg + .as_ref() + .into_iter() + .flat_map(|(_, channels)| channels) + .map(|id| with_folded_epg_id(id, |folded| folded.intern())) + .collect::>(); + let mut icon_assigned = HashSet::with_capacity(icon_tags.len()); - let referenced_epg_ids = referenced_live_epg_ids(fp); - epg_source - .children - .retain(|channel| with_folded_epg_id(&channel.id, |folded| referenced_epg_ids.contains(folded))); - - if id_cache.smart_match_enabled { - debug!( - "Smart EPG summary for input '{}': live={}, existing={}, exact={}, fuzzy={}, corrected={}, unresolved={}", - fp.input.name, - stats.live, - stats.existing, - stats.exact, - stats.fuzzy, - stats.corrected, - stats.unresolved - ); + let mut process_channel = |channel: &mut PlaylistItem| { + if !is_live_epg_item(channel) { + return true; } + let has_epg = assign_live_channel_epg( + channel, + id_cache, + &icon_tags, + &icon_override_channels, + &mut icon_assigned, + &mut stats, + ); + if required_epg && !has_epg { + removed += 1; + return false; + } + true + }; - new_epg.push(epg_source); + if required_epg { + fp.source.retain_memory_items_mut(&mut process_channel); + } else { + fp.items_mut().for_each(|channel| { + process_channel(channel); + }); } + } else { + warn!("Disk based playlist modification is not supported!"); + } + + if id_cache.smart_match_enabled { + debug!( + "Smart EPG summary for input '{}': live={}, existing={}, exact={}, fuzzy={}, corrected={}, unresolved={}", + 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 + .children + .retain(|channel| with_folded_epg_id(&channel.id, |folded| referenced_epg_ids.contains(folded))); + new_epg.push(epg_source); } } @@ -821,9 +894,9 @@ async fn assign_channel_epg(new_epg: &mut Vec, fp: &mut FetchedPlaylist<'_> /// ```text /// let mut playlist = FetchedPlaylist::default(); /// let mut epg_data = Vec::new(); -/// process_playlist_epg(&mut playlist, &mut epg_data); +/// process_playlist_epg(&mut playlist, &mut epg_data, false); /// ``` -pub async fn process_playlist_epg(fp: &mut FetchedPlaylist<'_>, epg: &mut Vec) { +pub async fn process_playlist_epg(fp: &mut FetchedPlaylist<'_>, epg: &mut Vec, required_epg: bool) { if fp.input.epg.is_none() { return; } @@ -831,10 +904,10 @@ pub async fn process_playlist_epg(fp: &mut FetchedPlaylist<'_>, epg: &mut Vec)], ) -> (Vec>>, Vec) { - run_xmltv_matches(xmltv, channels, true).await + run_xmltv_matches(xmltv, channels, true, false).await } async fn run_xmltv_matches( xmltv: &str, channels: &[(&str, Option<&str>)], smart_matching: bool, + required_epg: bool, ) -> (Vec>>, Vec) { let dir = tempdir().unwrap(); let epg_path = dir.path().join("smart-match.xml"); @@ -971,7 +1045,7 @@ mod tests { }; let mut epg = Vec::new(); - super::process_playlist_epg(&mut playlist, &mut epg).await; + super::process_playlist_epg(&mut playlist, &mut epg, required_epg).await; let assigned_ids = playlist.items_mut().map(|item| item.header.epg_channel_id.clone()).collect(); (assigned_ids, epg) } @@ -1185,7 +1259,7 @@ mod tests { }; let mut epg = Vec::new(); - super::process_playlist_epg(&mut playlist, &mut epg).await; + super::process_playlist_epg(&mut playlist, &mut epg, false).await; let updated = playlist.items_mut().next().unwrap(); assert_eq!(updated.header.epg_channel_id.as_deref(), Some("demo.channel")); @@ -1195,6 +1269,211 @@ mod tests { }); } + #[test] + fn required_epg_removes_only_unmatched_live_items() { + let runtime = tokio::runtime::Runtime::new().unwrap(); + runtime.block_on(async move { + let dir = tempdir().unwrap(); + let epg_path = dir.path().join("required-epg.xml"); + fs::write( + &epg_path, + r#" + Matched Live + + Programme + +"#, + ) + .unwrap(); + + let mut input = ConfigInput::from(ConfigInputDto::default()); + input.epg = Some(EpgConfig { sources: vec![], smart_match: None }); + let item = |name: &str, epg_id: Option<&str>, item_type: PlaylistItemType| PlaylistItem { + header: PlaylistItemHeader { + name: name.intern(), + epg_channel_id: epg_id.map(Internable::intern), + xtream_cluster: item_type.cluster(), + item_type, + ..PlaylistItemHeader::default() + }, + }; + let groups = vec![ + PlaylistGroup { + id: 1, + title: "Live".intern(), + channels: vec![ + item("Matched Live", Some("matched.live"), PlaylistItemType::Live), + item("Unmatched Live", Some("missing.live"), PlaylistItemType::LiveHls), + ], + xtream_cluster: super::XtreamCluster::Live, + }, + PlaylistGroup { + id: 2, + title: "Non-Live".intern(), + channels: vec![ + item("VOD", None, PlaylistItemType::Video), + item("Series", None, PlaylistItemType::Series), + item("Local VOD", None, PlaylistItemType::LocalVideo), + item("Local Series", None, PlaylistItemType::LocalSeries), + ], + xtream_cluster: super::XtreamCluster::Video, + }, + ]; + let tv_guide = TVGuide::new(vec![PersistedEpgSource { + file_path: epg_path, + priority: 0, + logo_override: false, + kind: PersistedEpgSourceKind::Xmltv, + }]); + let mut playlist = FetchedPlaylist { + input: &input, + source: MemoryPlaylistSource::new(groups).into_source(), + epg: Some(tv_guide), + }; + let mut epg = Vec::new(); + + super::process_playlist_epg(&mut playlist, &mut epg, true).await; + + let names = playlist.items_mut().map(|item| item.header.name.clone()).collect::>(); + 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")); + }); + } + + #[test] + fn required_epg_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()); + input.epg = Some(EpgConfig { sources: vec![], smart_match: None }); + let groups = vec![PlaylistGroup { + id: 1, + title: "Live".intern(), + channels: vec![live_playlist_item("Unmatched Live", Some("missing.live"))], + xtream_cluster: super::XtreamCluster::Live, + }]; + let mut playlist = + FetchedPlaylist { input: &input, source: MemoryPlaylistSource::new(groups).into_source(), epg: None }; + + super::process_playlist_epg(&mut playlist, &mut Vec::new(), true).await; + + assert_eq!(playlist.items_mut().count(), 1); + }); + } + + #[test] + fn optional_epg_preserves_existing_empty_groups() { + let runtime = tokio::runtime::Runtime::new().unwrap(); + runtime.block_on(async move { + let dir = tempdir().unwrap(); + let epg_path = dir.path().join("optional-epg-empty-group.xml"); + fs::write( + &epg_path, + r#" + Matched Live + + Programme + +"#, + ) + .unwrap(); + + let mut input = ConfigInput::from(ConfigInputDto::default()); + input.epg = Some(EpgConfig { sources: vec![], smart_match: None }); + let groups = vec![ + PlaylistGroup { + id: 1, + title: "Empty".intern(), + channels: vec![], + xtream_cluster: super::XtreamCluster::Live, + }, + PlaylistGroup { + id: 2, + title: "Live".intern(), + channels: vec![live_playlist_item("Matched Live", Some("matched.live"))], + xtream_cluster: super::XtreamCluster::Live, + }, + ]; + let tv_guide = TVGuide::new(vec![PersistedEpgSource { + file_path: epg_path, + priority: 0, + logo_override: false, + kind: PersistedEpgSourceKind::Xmltv, + }]); + let mut playlist = FetchedPlaylist { + input: &input, + source: MemoryPlaylistSource::new(groups).into_source(), + epg: Some(tv_guide), + }; + + super::process_playlist_epg(&mut playlist, &mut Vec::new(), false).await; + + assert_eq!(playlist.get_group_count(), 2); + }); + } + + #[test] + fn required_epg_removes_live_items_without_ids_when_smart_matching_is_disabled() { + let runtime = tokio::runtime::Runtime::new().unwrap(); + runtime.block_on(async move { + let dir = tempdir().unwrap(); + let epg_path = dir.path().join("required-epg-no-ids.xml"); + fs::write(&epg_path, "").unwrap(); + + let mut input = ConfigInput::from(ConfigInputDto::default()); + input.epg = Some(EpgConfig { sources: vec![], smart_match: None }); + let groups = vec![PlaylistGroup { + id: 1, + title: "Live".intern(), + channels: vec![live_playlist_item("Live without ID", None)], + xtream_cluster: super::XtreamCluster::Live, + }]; + let tv_guide = TVGuide::new(vec![PersistedEpgSource { + file_path: epg_path, + priority: 0, + logo_override: false, + kind: PersistedEpgSourceKind::Xmltv, + }]); + let mut playlist = FetchedPlaylist { + input: &input, + source: MemoryPlaylistSource::new(groups).into_source(), + epg: Some(tv_guide), + }; + + 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); + }); + } + + #[test] + fn required_epg_keeps_live_items_assigned_by_smart_matching() { + let runtime = tokio::runtime::Runtime::new().unwrap(); + runtime.block_on(async move { + let (assigned_ids, epg) = run_xmltv_matches( + r#" + TF1 + + Programme TF1 + +"#, + &[("TF1", None)], + true, + true, + ) + .await; + + assert_eq!(assigned_ids, vec![Some("tf1.fr".intern())]); + assert_eq!(epg[0].children.len(), 1); + assert_eq!(epg[0].children[0].id.as_ref(), "tf1.fr"); + }); + } + #[test] fn smart_match_replaces_an_existing_id_that_has_no_programmes() { let runtime = tokio::runtime::Runtime::new().unwrap(); @@ -1228,7 +1507,7 @@ mod tests { "#; for smart_matching in [false, true] { let (assigned_ids, epg) = - run_xmltv_matches(xmltv, &[("Empty channel", Some("empty.fr"))], smart_matching).await; + run_xmltv_matches(xmltv, &[("Empty channel", Some("empty.fr"))], smart_matching, false).await; assert_eq!(assigned_ids, vec![Some("empty.fr".intern())]); assert_eq!(epg[0].children.len(), 1); @@ -1308,7 +1587,7 @@ mod tests { }; let mut epg = Vec::new(); - super::process_playlist_epg(&mut playlist, &mut epg).await; + super::process_playlist_epg(&mut playlist, &mut epg, false).await; assert_eq!(epg.len(), 1); assert_eq!(epg[0].children[0].id.as_ref(), "f1.calendar"); @@ -1346,7 +1625,7 @@ mod tests { }; let mut epg = Vec::new(); - super::process_playlist_epg(&mut playlist, &mut epg).await; + super::process_playlist_epg(&mut playlist, &mut epg, false).await; let updated = playlist.items_mut().next().unwrap(); assert_eq!(updated.header.epg_channel_id.as_deref(), Some("f1.calendar")); diff --git a/backend/processing/src/processor/playlist.rs b/backend/processing/src/processor/playlist.rs index aea03bcff..820570e52 100644 --- a/backend/processing/src/processor/playlist.rs +++ b/backend/processing/src/processor/playlist.rs @@ -7,8 +7,12 @@ 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, sort::sort_playlist, trakt::process_trakt_categories_for_target, - xtream_series::playlist_resolve_series, xtream_vod::playlist_resolve_vod, StalkerRefreshMode, + epg::{process_playlist_epg, retain_live_items_with_processed_epg}, + sort::sort_playlist, + trakt::process_trakt_categories_for_target, + xtream_series::playlist_resolve_series, + xtream_vod::playlist_resolve_vod, + StalkerRefreshMode, }, }; use futures::{FutureExt, StreamExt}; @@ -21,8 +25,8 @@ use shared::{ error::{get_errors_notify_message, TuliproxError}, foundation::{get_field_value, set_field_value, Filter, ValueAccessor, ValueProvider}, model::{ - ClusterFlags, CounterModifier, EventMessage, EventSink, FieldGet, FieldSet, InputStats, InputType, - MappingStage, PipelineStats, PlaylistGroup, PlaylistItem, PlaylistItemType, PlaylistStats, + ClusterFlags, ConfigTargetOptions, CounterModifier, EventMessage, EventSink, FieldGet, FieldSet, InputStats, + InputType, MappingStage, PipelineStats, PlaylistGroup, PlaylistItem, PlaylistItemType, PlaylistStats, PlaylistUpdateProgressEvent, PlaylistUpdateSummary, ProviderFetchFailure, SourceStats, StreamProperties, TargetStats, UUIDType, WatchDisabled, WatchDisabledReason, WatchUnmatched, XtreamCluster, }, @@ -1690,7 +1694,9 @@ async fn prepare_playlist_for_target input_epg_start { + retain_live_items_with_processed_epg(&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(); @@ -3614,6 +3623,82 @@ match { }); } + #[test] + fn required_epg_removes_ids_invalidated_by_after_epg_mapping() { + let runtime = Runtime::new().expect("runtime"); + runtime.block_on(async { + let dir = tempdir().expect("tempdir"); + let ics_path = dir.path().join("bbc.ics"); + std::fs::write( + &ics_path, + "BEGIN:VCALENDAR\nBEGIN:VEVENT\nSUMMARY:News\nDTSTART:20260306T120000Z\nDTEND:20260306T130000Z\nEND:VEVENT\nEND:VCALENDAR", + ) + .expect("write ics"); + + let mut input = ConfigInput::from(ConfigInputDto::default()); + input.name = "input".intern(); + input.epg = Some(EpgConfig { sources: vec![], smart_match: None }); + let groups = vec![PlaylistGroup { + id: 1, + title: "Live".intern(), + channels: vec![PlaylistItem { + header: PlaylistItemHeader { + name: "BBC One".intern(), + epg_channel_id: Some("bbc.one".intern()), + group: "Live".intern(), + xtream_cluster: XtreamCluster::Live, + item_type: PlaylistItemType::Live, + ..Default::default() + }, + }], + xtream_cluster: XtreamCluster::Live, + }]; + let tv_guide = TVGuide::new(vec![PersistedEpgSource { + file_path: ics_path, + priority: 0, + logo_override: false, + kind: PersistedEpgSourceKind::Ics { + channel_id: "bbc.one".intern(), + channel_title: Some("BBC One".intern()), + match_names: vec![], + config: Box::new(IcsEpgSourceConfig::default()), + }, + }]); + let mut playlist = FetchedPlaylist { + input: &input, + source: MemoryPlaylistSource::new(groups).into_source(), + epg: Some(tv_guide), + }; + + let rewrite_epg = + 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() }); + 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 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()); + 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); + }); + } + fn live_item_for_epg(name: &str) -> PlaylistItem { PlaylistItem { header: PlaylistItemHeader { diff --git a/backend/repository/src/playlist_source.rs b/backend/repository/src/playlist_source.rs index 42da1260f..6717a63d7 100644 --- a/backend/repository/src/playlist_source.rs +++ b/backend/repository/src/playlist_source.rs @@ -261,6 +261,21 @@ impl PlaylistSource { ClusterFiltered::new(dispatch!(self.items()), skip_set) } + /// Retains visible items of an in-memory playlist without rebuilding its source wrapper. + /// + /// Returns `false` without modifying the source when it is not memory-backed. + pub fn retain_memory_items_mut(&mut self, predicate: F) -> bool + where + F: FnMut(&mut PlaylistItem) -> bool, + { + let skip_set = self.skip_set.as_deref(); + let PlaylistSourceKind::Memory(source) = &mut self.kind else { + return false; + }; + source.retain_items_mut(skip_set, predicate); + true + } + pub async fn update_playlist(&mut self, plg: &PlaylistGroup) { if self.skip_set.as_ref().is_some_and(|skip_set| skip_set.contains(&plg.xtream_cluster)) { return; @@ -1246,6 +1261,23 @@ impl MemoryPlaylistSource { pub fn into_source(self) -> PlaylistSource { PlaylistSource::memory(self) } + fn retain_items_mut(&mut self, skip_set: Option<&HashSet>, mut predicate: F) + where + F: FnMut(&mut PlaylistItem) -> bool, + { + let playlist = Arc::make_mut(&mut self.playlist); + playlist.retain_mut(|group| { + if skip_set.is_some_and(|clusters| clusters.contains(&group.xtream_cluster)) { + return false; + } + let was_empty = group.channels.is_empty(); + group.channels.retain_mut(|item| { + !skip_set.is_some_and(|clusters| clusters.contains(&item.header.xtream_cluster)) && predicate(item) + }); + was_empty || !group.channels.is_empty() + }); + } + /// Merge a batch of groups into the in-memory playlist in a single pass. /// /// Equivalent to calling [`PlaylistSourceOps::update_playlist`] for every @@ -1555,6 +1587,25 @@ mod tests { assert_eq!(cloned.get_channel_count(), 2); } + #[test] + fn retaining_memory_items_preserves_copy_on_write() { + let group = PlaylistGroup { + id: 1, + title: "Series".intern(), + channels: vec![make_item("keep", "Series", 1), make_item("remove", "Series", 1)], + xtream_cluster: XtreamCluster::Series, + }; + let source = MemoryPlaylistSource::new(vec![group]).into_source(); + let mut original = source.clone_source().expect("memory source clone should succeed"); + let mut retained = source.clone_source().expect("memory source clone should succeed"); + + assert!(retained.retain_memory_items_mut(|item| item.header.title.as_ref() != "remove")); + + assert_eq!(original.get_channel_count(), 2); + assert_eq!(retained.get_channel_count(), 1); + assert_eq!(retained.items_mut().next().map(|item| item.header.title.as_ref()), Some("keep")); + } + #[tokio::test] async fn xtream_disk_clone_source_returns_error_when_query_clone_fails() { let app_config = test_app_config(); diff --git a/config/source.yml b/config/source.yml index b1cd3d88c..4001c5f27 100644 --- a/config/source.yml +++ b/config/source.yml @@ -26,6 +26,7 @@ sources: filter: "!final_channel_lineup!" options: ignore_logo: false + required_epg: false epg_output: lowercase_ids: false lowercase_xmltv_display_names: false @@ -65,6 +66,7 @@ sources: filter: "!final_channel_lineup!" options: ignore_logo: false + required_epg: false epg_output: lowercase_ids: false lowercase_xmltv_display_names: false diff --git a/docs/src/configuration/source.md b/docs/src/configuration/source.md index 0c37b3f7a..67b7df382 100644 --- a/docs/src/configuration/source.md +++ b/docs/src/configuration/source.md @@ -1145,6 +1145,7 @@ sources: sort: { } options: ignore_logo: false + required_epg: false epg_output: lowercase_ids: false lowercase_xmltv_display_names: false @@ -1438,6 +1439,7 @@ targets: use_output: xtream options: ignore_logo: false + required_epg: false epg_output: lowercase_ids: true lowercase_xmltv_display_names: false @@ -1456,6 +1458,7 @@ targets: | Parameter | Type | Required | Default | Technical Impact & Background | |:-------------------------------------------|:-----|:--------:|:--------|:--------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------| | `ignore_logo` | Bool | No | `false` | Ignores `tvg-logo` and `tvg-logo-small` attributes. This reduces downstream device-side logo caching and can keep generated M3U playlists leaner for clients with limited storage or poor cache invalidation behavior. | +| `required_epg` | Bool | No | `false` | Keeps only live playlist entries whose EPG ID resolves to programme data from an available EPG source. Filtering and mappings in the regular processing stage run first, so EPG matching only processes the reduced target playlist. VOD, series, and local-library entries are unaffected. If no EPG source was successfully materialized for the input, this option leaves its live entries unchanged. | | `share_live_streams.hls` | Bool | No | `false` | Enables HLS live sharing for the new HLS cache proxy path. This is a configuration switch for the HLS cache feature and is independent from MPEG-TS stream sharing. | | `share_live_streams.mpeg_ts` | Bool | No | `false` | Allows Tuliprox to share MPEG-TS live stream connections in reverse proxy mode. This can reduce upstream provider connection usage when multiple clients watch the same channel, but it increases memory usage per shared channel. | | `remove_duplicates` | Bool | No | `false` | Legacy pre-transform identity deduplication. It runs independently for each input before the F/R/M pipe and removes repeated source identities before mapping can emit additional items. The field remains supported for backward compatibility. | diff --git a/frontend/public/assets/i18n/ar.json b/frontend/public/assets/i18n/ar.json index d6de14bdc..6d42def5f 100644 --- a/frontend/public/assets/i18n/ar.json +++ b/frontend/public/assets/i18n/ar.json @@ -379,6 +379,7 @@ "WATCH": "قائمة أنماط المجموعات لمراقبة التغييرات. يتم إرسال الإشعارات عبر المراسلة المكوّنة." }, "CONFIG_TARGET_OPTIONS": { + "REQUIRED_EPG": "يحتفظ فقط بالقنوات المباشرة التي تتوفر لها بيانات EPG مطابقة. لا تتأثر عناصر الفيديو حسب الطلب والمسلسلات والمكتبة المحلية.", "SHARE_LIVE_STREAMS": "يُفعّل مشاركة البث المباشر. المفتاح الرئيسي يضبط HLS و MPEG-TS معًا؛ يمكن تغيير كل تنسيق أيضًا بشكل منفصل." }, "CONFIG_TARGET_SHARE_LIVE_STREAMS": { @@ -1369,6 +1370,7 @@ "REGEXP": "تعبير منتظم", "RELEASES": "الإصدارات", "REMOVE_DUPLICATES": "إزالة التكرارات", + "REQUIRED_EPG": "اشتراط EPG للقنوات المباشرة", "RENAME": "إعادة تسمية", "RENAME_SETTINGS": "إعادة تسمية الإعدادات", "REPEAT_PASSWORD": "تكرار كلمة المرور", diff --git a/frontend/public/assets/i18n/en.json b/frontend/public/assets/i18n/en.json index 0a5da8cd3..506f0f34b 100644 --- a/frontend/public/assets/i18n/en.json +++ b/frontend/public/assets/i18n/en.json @@ -379,6 +379,7 @@ "WATCH": "List of group patterns to monitor for changes. Notifications are sent via configured messaging." }, "CONFIG_TARGET_OPTIONS": { + "REQUIRED_EPG": "Keeps only live channels that match programme data from an available EPG source. VOD, series, and local library items are unaffected.", "SHARE_LIVE_STREAMS": "Enables live stream sharing. The master toggle sets HLS and MPEG-TS together; each format can also be changed separately." }, "EPG_OUTPUT_OPTIONS": { @@ -1454,6 +1455,7 @@ "REGEXP": "Regexp", "RELEASES": "Releases", "REMOVE_DUPLICATES": "Remove duplicates", + "REQUIRED_EPG": "Require EPG for live channels", "RENAME": "Rename", "RENAME_SETTINGS": "Rename Settings", "REPEAT_PASSWORD": "Repeat Password", diff --git a/frontend/public/assets/i18n/ru.json b/frontend/public/assets/i18n/ru.json index 7a337fed8..9c5599dde 100644 --- a/frontend/public/assets/i18n/ru.json +++ b/frontend/public/assets/i18n/ru.json @@ -355,6 +355,9 @@ "USE_MEMORY_CACHE": "Если включено, плейлист кэшируется в оперативной памяти для уменьшения нагрузки на диск (увеличивает использование памяти).", "WATCH": "Список шаблонов групп для мониторинга изменений. Уведомления отправляются через настроенный мессенджер." }, + "CONFIG_TARGET_OPTIONS": { + "REQUIRED_EPG": "Оставляет только прямые каналы, для которых найдены данные программ, но только при наличии успешно загруженного источника EPG. Если такой источник отсутствует, плейлист не изменяется. VOD, сериалы и элементы локальной библиотеки не затрагиваются." + }, "EPG_OUTPUT_OPTIONS": { "LOWERCASE_IDS": "Преобразует технические идентификаторы EPG в нижний регистр ASCII во всех выходах M3U, Xtream, XMLTV и API EPG. После изменения этой опции требуется полное обновление целевого плейлиста.", "LOWERCASE_XMLTV_DISPLAY_NAMES": "Преобразует в нижний регистр Unicode только текст display-name XMLTV. Имена плейлиста и названия программ остаются без изменений." @@ -1308,6 +1311,7 @@ "REGEXP": "Регулярное выражение", "RELEASES": "Релизы", "REMOVE_DUPLICATES": "Удалить дубликаты", + "REQUIRED_EPG": "Требовать EPG для прямых каналов", "RENAME": "Переименовать", "RENAME_SETTINGS": "Настройки переименования", "REPEAT_PASSWORD": "Повторите пароль", diff --git a/frontend/src/app/components/source_editor/target_form.rs b/frontend/src/app/components/source_editor/target_form.rs index 7862e2b18..32cc5f600 100644 --- a/frontend/src/app/components/source_editor/target_form.rs +++ b/frontend/src/app/components/source_editor/target_form.rs @@ -29,6 +29,7 @@ const LABEL_ADD_WATCH: &str = "LABEL.ADD_WATCH"; const LABEL_USE_MEMORY_CACHE: &str = "LABEL.USE_MEMORY_CACHE"; const LABEL_PROCESSING_ORDER: &str = "LABEL.PROCESSING_ORDER"; const LABEL_IGNORE_LOGO: &str = "LABEL.IGNORE_LOGO"; +const LABEL_REQUIRED_EPG: &str = "LABEL.REQUIRED_EPG"; const LABEL_SHARE_LIVE_STREAMS: &str = "LABEL.SHARE_LIVE_STREAMS"; const LABEL_HLS: &str = "LABEL.HLS"; const LABEL_MPEG_TS: &str = "LABEL.MPEG_TS"; @@ -104,6 +105,7 @@ impl HasFormData for ConfigTargetOptionsFormState { #[derive(Clone)] pub enum ConfigTargetOptionsFormAction { IgnoreLogo(bool), + RequiredEpg(bool), ShareLiveStreamsHls(bool), ShareLiveStreamsMpegTs(bool), RemoveDuplicates(bool), @@ -125,6 +127,10 @@ impl yew::prelude::Reducible for ConfigTargetOptionsFormState { form.ignore_logo = value; modified = true; } + ConfigTargetOptionsFormAction::RequiredEpg(value) => { + form.required_epg = value; + modified = true; + } ConfigTargetOptionsFormAction::ShareLiveStreamsHls(value) => { form.share_live_streams.hls = value; modified = true; @@ -307,6 +313,7 @@ pub fn ConfigTargetView(props: &ConfigTargetViewProps) -> Html {
{ edit_field_bool!(target_options_state, translate.t(LABEL_IGNORE_LOGO), ignore_logo, ConfigTargetOptionsFormAction::IgnoreLogo) } + { edit_field_bool!(target_options_state, translate.t(LABEL_REQUIRED_EPG), required_epg, ConfigTargetOptionsFormAction::RequiredEpg) }
{ translate.t(LABEL_SHARE_LIVE_STREAMS) } @@ -359,6 +366,7 @@ pub fn ConfigTargetView(props: &ConfigTargetViewProps) -> Html {
{ config_field_bool!(target_options_state.form, translate.t(LABEL_IGNORE_LOGO), ignore_logo) } + { config_field_bool!(target_options_state.form, translate.t(LABEL_REQUIRED_EPG), required_epg) }
{ translate.t(LABEL_SHARE_LIVE_STREAMS) } @@ -553,6 +561,14 @@ mod tests { assert!(!state.form.epg_output.lowercase_xmltv_display_names); } + #[test] + fn required_epg_action_updates_target_option() { + let state = default_options_state().reduce(ConfigTargetOptionsFormAction::RequiredEpg(true)); + + assert!(state.form.required_epg); + assert!(state.modified); + } + #[test] fn lowercase_xmltv_display_names_action_updates_nested_option() { let state = default_options_state().reduce(ConfigTargetOptionsFormAction::LowercaseXmltvDisplayNames(true)); diff --git a/shared/src/model/config/target.rs b/shared/src/model/config/target.rs index 681f0071d..9bb80c0f9 100644 --- a/shared/src/model/config/target.rs +++ b/shared/src/model/config/target.rs @@ -94,6 +94,8 @@ pub struct DeduplicateConfig { pub struct ConfigTargetOptions { #[serde(default, skip_serializing_if = "is_false")] pub ignore_logo: bool, + #[serde(default, skip_serializing_if = "is_false")] + pub required_epg: bool, #[serde( default, deserialize_with = "deserialize_share_live_streams", @@ -113,6 +115,7 @@ pub struct ConfigTargetOptions { impl ConfigTargetOptions { pub fn is_empty(&self) -> bool { !self.ignore_logo + && !self.required_epg && self.share_live_streams.is_empty() && !self.remove_duplicates && self.deduplicate.is_none() @@ -124,6 +127,8 @@ impl ConfigTargetOptions { pub const fn lowercase_xmltv_display_names(&self) -> bool { self.epg_output.lowercase_xmltv_display_names } + pub const fn required_epg(&self) -> bool { self.required_epg } + pub fn share_live_hls_enabled(&self) -> bool { self.share_live_streams.hls } pub fn share_live_mpeg_ts_enabled(&self) -> bool { self.share_live_streams.mpeg_ts } @@ -654,6 +659,7 @@ share_live_streams: false let options = ConfigTargetOptions::default(); assert!(options.is_empty()); + assert!(!options.required_epg()); let serialized = serde_saphyr::to_string(&options).expect("default options should serialize"); assert!( @@ -676,6 +682,18 @@ share_live_streams: false assert_eq!(reparsed.share_live_streams, options.share_live_streams); } + #[test] + fn target_options_required_epg_roundtrips_and_makes_options_nonempty() { + let options = serde_saphyr::from_str::("required_epg: true\n") + .expect("required_epg should deserialize"); + + assert!(options.required_epg()); + assert!(!options.is_empty()); + + let serialized = serde_saphyr::to_string(&options).expect("required_epg should serialize"); + assert!(serialized.contains("required_epg: true")); + } + #[test] fn target_options_default_epg_output_is_disabled_and_omitted() { let options = serde_saphyr::from_str::("{}")