epg timeshift refactoring (#550)

epg timeshift refactoring
This commit is contained in:
euzu
2026-01-29 18:08:22 +01:00
committed by GitHub
parent 2a63fdc155
commit 0be9258fc0
8 changed files with 316 additions and 558 deletions
+28 -17
View File
@@ -9,9 +9,9 @@ use crate::repository::XML_PREAMBLE;
use crate::repository::{get_target_storage_path, BPlusTreeQuery};
use crate::repository::{xtream_get_epg_file_path_for_target, xtream_get_storage_path};
use crate::utils;
use crate::utils::{deobscure_text, file_exists_async, format_xmltv_time_utc, get_epg_processing_options, obscure_text, EpgProcessingOptions};
use crate::utils::{deobscure_text, file_exists_async, format_xmltv_time_utc, get_epg_processing_options, obscure_text, EpgProcessingOptions, EpgTimeShift};
use axum::response::IntoResponse;
use chrono::DateTime;
use chrono::{DateTime, TimeZone};
use log::{error, trace};
use quick_xml::events::{BytesEnd, BytesStart, BytesText, Event};
use shared::concat_string;
@@ -135,7 +135,7 @@ async fn serve_epg_with_rewrites(
let epg_processing_options = get_epg_processing_options(app_state, user, target);
let base_url = if epg_processing_options.offset_minutes != 0 || epg_processing_options.rewrite_urls {
let base_url = if !matches!(epg_processing_options.time_shift, EpgTimeShift::None) || epg_processing_options.rewrite_urls {
let server_info = app_state.app_config.get_user_server_info(user);
Some(concat_string!(&server_info.get_base_url(), "/", storage_const::EPG_RESOURCE_PATH, "/", &user.username, "/", &user.password))
} else {
@@ -195,8 +195,8 @@ async fn serve_epg_with_rewrites(
for programme in programmes {
let mut elem = BytesStart::new("programme");
let (user_start, user_stop) = (programme.start, programme.stop);
elem.push_attribute(("start", format_xmltv_time_utc(user_start, epg_processing_options.offset_minutes).as_str()));
elem.push_attribute(("stop", format_xmltv_time_utc(user_stop, epg_processing_options.offset_minutes).as_str()));
elem.push_attribute(("start", format_xmltv_time_utc(user_start, &epg_processing_options.time_shift).as_str()));
elem.push_attribute(("stop", format_xmltv_time_utc(user_stop, &epg_processing_options.time_shift).as_str()));
elem.push_attribute(("channel", channel.id.as_ref()));
continue_on_err!(writer.write_event_async(Event::Start(elem)).await);
@@ -261,29 +261,40 @@ fn format_xmltv_time(ts: i64) -> String {
}
}
fn apply_user_offset(start: i64, stop: i64, offset_minutes: i32) -> (i64, i64) {
let offset = i64::from(offset_minutes) * 60;
let user_start = start + offset;
let user_end = stop + offset;
(user_start, user_end)
fn get_applied_epg_timeshift(programme: &EpgProgramme, epg_processing_options: &EpgProcessingOptions) -> (String, String, i64, i64) {
match &epg_processing_options.time_shift {
EpgTimeShift::None => (format_xmltv_time(programme.start), format_xmltv_time(programme.stop), programme.start, programme.stop),
EpgTimeShift::Fixed(m) => {
let off = i64::from(*m) * 60;
let s = programme.start + off;
let e = programme.stop + off;
(format_xmltv_time(s), format_xmltv_time(e), s, e)
}
EpgTimeShift::TimeZone(tz) => {
let s_dt = chrono::Utc.timestamp_opt(programme.start, 0).unwrap().with_timezone(tz);
let e_dt = chrono::Utc.timestamp_opt(programme.stop, 0).unwrap().with_timezone(tz);
// We use the original timestamps (programme.start/stop) here because TimeZone adjustment
// is only for the visual string representation. The absolute event time (UTC) remains unchanged.
// Unlike 'Fixed' offset which artificially shifts the event time.
(s_dt.format("%Y-%m-%d %H:%M:%S").to_string(), e_dt.format("%Y-%m-%d %H:%M:%S").to_string(), programme.start, programme.stop)
}
}
}
fn from_programme(stream_id: &Arc<str>, epg_id: &Arc<str>, programme: &EpgProgramme, epg_processing_options: &EpgProcessingOptions) -> ShortEpgDto {
let (user_start, user_end) = apply_user_offset(programme.start, programme.stop, epg_processing_options.offset_minutes);
let (start_str, end_str, start_ts, stop_ts) = get_applied_epg_timeshift(programme, epg_processing_options);
ShortEpgDto {
id: Arc::clone(stream_id),
epg_id: Arc::clone(epg_id),
title: programme.title.as_ref().map_or_else(String::new, ToString::to_string),
lang: String::new(),
start: format_xmltv_time(user_start),
end: format_xmltv_time(user_end),
start: start_str,
end: end_str,
description: programme.desc.as_ref().map_or_else(String::new, ToString::to_string),
channel_id: Arc::clone(epg_id),
start_timestamp: user_start.to_string(),
stop_timestamp: user_end.to_string(),
start_timestamp: start_ts.to_string(),
stop_timestamp: stop_ts.to_string(),
stream_id: Arc::clone(stream_id),
now_playing: None,
has_archive: None,
+3 -93
View File
@@ -5,7 +5,7 @@ use crate::utils::request::get_remote_content_as_stream;
use crate::utils::{async_file_reader, parse_xmltv_time};
use chrono::{Datelike, TimeZone, Utc};
use futures::TryFutureExt;
use quick_xml::events::{BytesEnd, BytesStart, BytesText, Event};
use quick_xml::events::{Event};
use shared::concat_string;
use shared::error::{TuliproxError, TuliproxErrorKind};
use shared::model::{EpgChannel, EpgProgramme, InputFetchMethod};
@@ -13,7 +13,7 @@ use shared::utils::{sanitize_sensitive_info, Internable};
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use tokio::io::{AsyncRead, AsyncWrite, AsyncWriteExt};
use tokio::io::{AsyncRead};
use url::Url;
pub const EPG_TAG_TV: &str = "tv";
@@ -68,97 +68,7 @@ pub struct Epg {
pub priority: i16,
pub logo_override: bool,
pub attributes: Option<HashMap<Arc<str>, Arc<str>>>,
pub children: Vec<Arc<XmlTag>>,
}
impl Epg {
pub async fn write_to_async<W: AsyncWrite + Unpin>(
&self,
writer: &mut quick_xml::writer::Writer<W>,
rename_map: Option<&HashMap<Arc<str>, Arc<str>>>,
) -> Result<(), quick_xml::Error> {
// Start tv-element
let mut elem = BytesStart::new("tv");
if let Some(attrs) = &self.attributes {
for (k, v) in attrs {
elem.push_attribute((k.as_ref(), v.as_ref()));
}
}
writer.write_event_async(Event::Start(elem)).await?;
// Stack for iterative writing
// bool = End-Event written?
let mut stack: Vec<(&XmlTag, bool)> = self
.children
.iter()
.rev()
.map(|c| (c.as_ref(), false))
.collect();
let mut write_counter = 0usize;
let mut current_channel_id: Option<Arc<str>> = None;
let epg_id_key = EPG_ATTRIB_ID.intern();
while let Some((tag, ended)) = stack.pop() {
if ended {
if tag.name.as_ref() == EPG_TAG_CHANNEL {
current_channel_id = None;
}
// End-Event
writer
.write_event_async(Event::End(BytesEnd::new(tag.name.as_ref())))
.await?;
} else {
// Start-Event for the tag
let mut elem = BytesStart::new(tag.name.as_ref());
if let Some(attrs) = &tag.attributes {
for (k, v) in attrs {
elem.push_attribute((k.as_ref(), v.as_ref()));
}
}
if tag.name.as_ref() == EPG_TAG_CHANNEL {
current_channel_id = tag.get_attribute_value(&epg_id_key).cloned();
}
writer.write_event_async(Event::Start(elem)).await?;
// write text
let value_to_write = if tag.name.as_ref() == EPG_TAG_DISPLAY_NAME {
current_channel_id.as_ref()
.and_then(|cid| rename_map.and_then(|m| m.get(cid))
.or(tag.value.as_ref()))
} else {
tag.value.as_ref()
};
if let Some(text) = value_to_write {
writer.write_event_async(Event::Text(BytesText::new(text))).await?;
}
// End-Marker push + children push
stack.push((tag, true));
if let Some(children) = &tag.children {
for child in children.iter().rev() {
stack.push((child.as_ref(), false));
}
}
}
write_counter += 1;
if write_counter >= 50 {
writer.get_mut().flush().await?; // flush underlying writer
write_counter = 0;
}
}
// write tv-end
writer.write_event_async(Event::End(BytesEnd::new("tv"))).await?;
let inner = writer.get_mut();
inner.flush().await?;
Ok(())
}
pub children: Vec<Arc<EpgChannel>>,
}
#[derive(Debug, Clone)]
+180 -96
View File
@@ -1,19 +1,18 @@
use crate::model::{Epg, TVGuide, XmlTag, XmlTagIcon, EPG_ATTRIB_CHANNEL, EPG_ATTRIB_ID, EPG_TAG_CHANNEL, EPG_TAG_DISPLAY_NAME, EPG_TAG_ICON, EPG_TAG_PROGRAMME, EPG_TAG_TV};
use crate::model::{EpgSmartMatchConfig, PersistedEpgSource};
use crate::processing::processor::epg::EpgIdCache;
use crate::utils::async_file_reader;
use crate::utils::compressed_file_reader_async::CompressedFileReaderAsync;
use dashmap::DashMap;
use crate::utils::{async_file_reader, parse_xmltv_time};
use log::error;
use quick_xml::events::{BytesStart, BytesText, Event};
use rayon::iter::{IntoParallelRefIterator, ParallelIterator};
use shared::concat_string;
use shared::model::EpgNamePrefix;
use shared::model::{EpgChannel, EpgNamePrefix, EpgProgramme};
use shared::utils::{deunicode_string, Internable, CONSTANTS};
use std::borrow::Cow;
use std::cmp::min;
use std::collections::HashMap;
use std::mem;
use std::sync::{Arc, Mutex};
use std::collections::{HashMap, HashSet};
use std::sync::Arc;
use tokio::io::AsyncRead;
/// Splits a string at the first delimiter if the prefix matches a known country code.
@@ -96,7 +95,7 @@ impl TVGuide {
pub fn merge(epgs: Vec<Epg>) -> Option<Epg> {
if let Some(first_epg) = epgs.first() {
let first_epg_attributes = first_epg.attributes.clone();
let merged_children: Vec<Arc<XmlTag>> = epgs.into_iter().flat_map(|epg| epg.children).collect();
let merged_children: Vec<Arc<EpgChannel>> = epgs.into_iter().flat_map(|epg| epg.children).collect();
Some(Epg {
logo_override: false,
priority: 0,
@@ -259,9 +258,14 @@ impl TVGuide {
async fn process_epg_file(id_cache: &mut EpgIdCache, epg_source: &PersistedEpgSource) -> Option<Epg> {
let epg_attrib_id = EPG_ATTRIB_ID.intern();
let epg_attrib_channel = EPG_ATTRIB_CHANNEL.intern();
let start_attrib = "start".intern();
let stop_attrib = "stop".intern();
let tag_title = "title".intern();
let tag_desc = "desc".intern();
match CompressedFileReaderAsync::new(&epg_source.file_path).await {
Ok(mut reader) => {
let mut children: Vec<Arc<XmlTag>> = vec![];
let mut children: HashMap<Arc<str>, EpgChannel> = HashMap::with_capacity(5000);
let mut tv_attributes: Option<HashMap<Arc<str>, Arc<str>>> = None;
let smart_match = id_cache.smart_match_config.enabled;
let fuzzy_matching = smart_match && id_cache.smart_match_config.fuzzy_matching;
@@ -271,21 +275,58 @@ impl TVGuide {
let tag_epg_id = tag.get_attribute_value(&epg_attrib_id).map_or_else(|| "".intern(), Internable::intern);
if !tag_epg_id.is_empty() && !id_cache.processed.contains(&tag_epg_id) {
Self::prepare_tag(id_cache, &mut tag, smart_match);
let mut add_channel = false;
if smart_match {
if Self::try_fuzzy_matching(id_cache, &tag_epg_id, &tag, fuzzy_matching) {
children.push(Arc::new(tag));
id_cache.processed.insert(tag_epg_id);
add_channel = true;
}
} else if id_cache.channel_epg_id.contains(&tag_epg_id) {
children.push(Arc::new(tag));
add_channel = true;
}
if add_channel && !children.contains_key(&tag_epg_id) {
let display_name = tag.children.as_ref().and_then(|children| {
children.iter()
.find(|c| c.name.as_ref() == EPG_TAG_DISPLAY_NAME)
.and_then(|c| c.value.clone())
});
children.insert(Arc::clone(&tag_epg_id), EpgChannel {
id: Arc::clone(&tag_epg_id),
title: display_name,
icon: if let XmlTagIcon::Src(src) = &tag.icon { Some(Arc::clone(src)) } else { None },
programmes: vec![],
});
id_cache.processed.insert(tag_epg_id);
}
}
}
EPG_TAG_PROGRAMME => {
if let Some(epg_id) = tag.get_attribute_value(&epg_attrib_channel) {
if id_cache.processed.contains(epg_id) && id_cache.channel_epg_id.contains(epg_id) {
children.push(Arc::new(tag));
if id_cache.processed.contains(epg_id) /*&& id_cache.channel_epg_id.contains(epg_id) */{
if let Some(channel) = children.get_mut(epg_id) {
if let Some((Some(start), Some(stop))) = tag.attributes.as_ref().map(|a| (a.get(&start_attrib), a.get(&stop_attrib))) {
if let (Some(start_time), Some(stop_time)) = (parse_xmltv_time(start), parse_xmltv_time(stop)) {
let mut title = None;
let mut desc = None;
if let Some(children) = tag.children.as_ref() {
for child in children {
if child.name == tag_title {
title.clone_from(&child.value);
} else if child.name == tag_desc {
desc.clone_from(&child.value);
}
}
channel.programmes.push(EpgProgramme::new_all(start_time, stop_time, Arc::clone(epg_id), title, desc));
}
} else {
error!("Failed to parse epg programme time {start} - {stop}");
}
} else {
error!("Missing start or stop attribute in programme tag, skipping");
}
} else {
error!("Channel {epg_id} not found in EPG, dangling programme");
}
}
}
}
@@ -306,7 +347,7 @@ impl TVGuide {
logo_override: epg_source.logo_override,
priority: epg_source.priority,
attributes: tv_attributes,
children,
children: children.into_values().map(Arc::new).collect(),
})
}
Err(e) => {
@@ -338,12 +379,12 @@ where
let binding = e.name();
let name_raw = String::from_utf8_lossy(binding.as_ref());
let name = name_raw.intern();
let (is_tv_tag, is_channel, is_program) = get_tag_types(&name);
let attributes = collect_tag_attributes(e, is_channel, is_program);
let tag_type = get_tag_type(&name);
let attributes = collect_tag_attributes(e, tag_type);
let attribs = if attributes.is_empty() { None } else { Some(attributes) };
let tag = XmlTag::new(name, attribs);
if is_tv_tag {
if tag_type.is_tv() {
callback(tag);
} else {
stack.push(tag);
@@ -425,17 +466,41 @@ where
}
}
fn get_tag_types(name: &str) -> (bool, bool, bool) {
let (is_tv_tag, is_channel, is_program) = match name {
EPG_TAG_TV => (true, false, false),
EPG_TAG_CHANNEL => (false, true, false),
EPG_TAG_PROGRAMME => (false, false, true),
_ => (false, false, false)
};
(is_tv_tag, is_channel, is_program)
#[derive(Copy, Clone, PartialEq, Eq, Hash, Debug)]
enum XmlTagType {
Ignored,
Tv,
Channel,
Programme,
}
fn collect_tag_attributes(e: &BytesStart, is_channel: bool, is_program: bool) -> HashMap<Arc<str>, Arc<str>> {
impl XmlTagType {
#[inline]
pub(crate) fn is_tv(self) -> bool {
self == XmlTagType::Tv
}
#[inline]
pub(crate) fn is_channel(self) -> bool {
self == XmlTagType::Channel
}
#[inline]
pub(crate) fn is_program(self) -> bool {
self == XmlTagType::Programme
}
}
fn get_tag_type(name: &str) -> XmlTagType {
match name {
EPG_TAG_TV => XmlTagType::Tv,
EPG_TAG_CHANNEL => XmlTagType::Channel,
EPG_TAG_PROGRAMME => XmlTagType::Programme,
_ => XmlTagType::Ignored
}
}
fn collect_tag_attributes(e: &BytesStart, tag_type: XmlTagType) -> HashMap<Arc<str>, Arc<str>> {
let attributes = e.attributes().filter_map(Result::ok)
.filter_map(|a| {
let key_binding = a.key;
@@ -444,7 +509,7 @@ fn collect_tag_attributes(e: &BytesStart, is_channel: bool, is_program: bool) ->
if let Ok(value) = a.unescape_value().as_ref() {
if value.is_empty() {
None
} else if (is_channel && key.as_ref() == EPG_ATTRIB_ID) || (is_program && key.as_ref() == EPG_ATTRIB_CHANNEL) {
} else if (tag_type.is_channel() && key.as_ref() == EPG_ATTRIB_ID) || (tag_type.is_program() && key.as_ref() == EPG_ATTRIB_CHANNEL) {
Some((key, value.to_lowercase().intern()))
} else {
Some((key, value.intern()))
@@ -456,79 +521,97 @@ fn collect_tag_attributes(e: &BytesStart, is_channel: bool, is_program: bool) ->
attributes
}
pub fn flatten_tvguide(tv_guides: &[Epg]) -> Option<Epg> {
if tv_guides.is_empty() {
None
} else {
let epg_children: Mutex<Vec<Arc<XmlTag>>> = Mutex::new(Vec::new());
let epg_attributes: Option<HashMap<Arc<str>, Arc<str>>> = tv_guides.first().and_then(|t| t.attributes.clone());
let count = tv_guides.iter().map(|tvg| tvg.children.len()).sum();
let channel_mapping: DashMap<Arc<str>, i16> = DashMap::with_capacity(count);
#[derive(Hash, Eq, PartialEq)]
struct ProgrammeKey {
start: i64,
stop: i64,
}
let mut sorted_guides = tv_guides.to_vec();
let epg_attrib_id = EPG_ATTRIB_ID.intern();
let epg_attrib_channel = EPG_ATTRIB_CHANNEL.intern();
// sort by priority
sorted_guides.sort_by_key(|a| a.priority);
// if executed parallel it does not matter how we sort.
sorted_guides.par_iter().for_each(|guide| {
let mut children = vec![];
guide.children.iter().for_each(|c| {
if c.name.as_ref() == EPG_TAG_CHANNEL {
if let Some(chan_id) = c.get_attribute_value(&epg_attrib_id) {
let chan_id = chan_id.intern();
let should_add = {
// if not stored
!channel_mapping.contains_key(&chan_id) ||
// or if priority is higher (less means higher priority)
channel_mapping.get(&chan_id).as_deref().is_none_or(|&priority| guide.priority < priority)
};
if should_add {
if let Some(mut existing) = channel_mapping.get_mut(&chan_id) {
if guide.priority < *existing {
*existing = guide.priority;
children.push(c.clone());
}
} else {
channel_mapping.insert(chan_id.clone(), guide.priority);
children.push(c.clone());
}
}
}
}
});
guide.children.iter().for_each(|c| {
if c.name.as_ref() == EPG_TAG_PROGRAMME {
if let Some(chan_id) = c.get_attribute_value(&epg_attrib_channel) {
let chan_id = chan_id.intern();
if let Some(stored_priority) = channel_mapping.get(&chan_id) {
if *stored_priority == guide.priority {
children.push(c.clone());
}
}
}
}
});
if let Ok(mut guard) = epg_children.lock() {
guard.extend(children);
}
});
let children = if let Ok(mut children) = epg_children.lock() {
mem::take(&mut *children)
} else {
vec![]
};
let epg = Epg {
logo_override: false,
priority: 0,
attributes: epg_attributes,
children,
};
Some(epg)
impl From<&EpgProgramme> for ProgrammeKey {
fn from(p: &EpgProgramme) -> Self {
Self {
start: p.start,
stop: p.stop,
}
}
}
struct ChannelAcc {
priority: i16,
channel: EpgChannel,
programmes: HashSet<ProgrammeKey>,
}
pub fn flatten_tvguide(mut tv_guides: Vec<Epg>) -> Option<Epg> {
if tv_guides.is_empty() {
return None;
}
let epg_attributes = tv_guides
.first()
.and_then(|t| t.attributes.clone());
let mut channels: HashMap<Arc<str>, ChannelAcc> = HashMap::new();
for guide in tv_guides.drain(..) {
for channel_arc in guide.children {
let Ok(mut channel) = Arc::try_unwrap(channel_arc) else {
error!("Failed to unwrap epg channel");
continue;
};
match channels.entry(Arc::clone(&channel.id)) {
std::collections::hash_map::Entry::Occupied(mut entry) => {
let acc = entry.get_mut();
if guide.priority < acc.priority {
// high priority
acc.priority = guide.priority;
acc.channel = channel;
acc.programmes.clear();
for p in &acc.channel.programmes {
acc.programmes.insert(ProgrammeKey::from(p));
}
} else if guide.priority == acc.priority {
// same priority → merge
for p in channel.programmes.drain(..) {
let key = ProgrammeKey::from(&p);
if acc.programmes.insert(key) {
acc.channel.programmes.push(p);
}
}
}
}
std::collections::hash_map::Entry::Vacant(entry) => {
let mut set = HashSet::new();
for p in &channel.programmes {
set.insert(ProgrammeKey::from(p));
}
entry.insert(ChannelAcc {
priority: guide.priority,
channel,
programmes: set,
});
}
}
}
}
let children = channels
.into_values()
.map(|acc| Arc::new(acc.channel))
.collect();
Some(Epg {
logo_override: false,
priority: 0,
attributes: epg_attributes,
children,
})
}
#[cfg(test)]
mod tests {
use crate::model::{EpgSmartMatchConfig, PersistedEpgSource, TVGuide};
@@ -553,6 +636,7 @@ mod tests {
}
#[ignore]
#[test]
fn parse_test() -> io::Result<()> {
let run_test = async move || {
+8 -9
View File
@@ -1,11 +1,11 @@
use crate::model::{Epg, TVGuide, XmlTag, XmlTagIcon, EPG_ATTRIB_ID};
use crate::model::{Epg, TVGuide};
use crate::model::{EpgConfig, EpgSmartMatchConfig};
use crate::model::FetchedPlaylist;
use crate::processing::parser::xmltv::normalize_channel_name;
use log::{debug, trace, warn};
use rphonetic::{DoubleMetaphone, Encoder};
use std::collections::{HashMap, HashSet};
use shared::model::{EpgSmartMatchConfigDto, PlaylistItem, XtreamCluster};
use shared::model::{EpgChannel, EpgSmartMatchConfigDto, PlaylistItem, XtreamCluster};
use std::sync::Arc;
use shared::utils::Internable;
@@ -167,12 +167,11 @@ async fn assign_channel_epg(new_epg: &mut Vec<Epg>, fp: &mut FetchedPlaylist<'_>
let mut processed_epgs = vec![];
if let Some(epg_sources) = tv_guide.filter(id_cache).await {
let mut icon_assigned = HashSet::new();
let epg_attrib_id = EPG_ATTRIB_ID.intern();
for epg_source in epg_sources {
// icon tags
let icon_tags: HashMap<&Arc<str>, &Arc<XmlTag>> = epg_source.children.iter()
.filter(|tag| tag.icon != XmlTagIcon::Undefined)
.filter_map(|tag| tag.get_attribute_value(&epg_attrib_id).map(|id| (id, tag)))
let icon_tags: HashMap<&Arc<str>, &Arc<EpgChannel>> = epg_source.children.iter()
.filter(|tag| tag.icon.as_ref().is_some_and(|i| !i.is_empty()))
.map(|tag| (&tag.id, tag))
.collect();
let assign_values = |chan: &mut PlaylistItem| {
@@ -204,14 +203,14 @@ async fn assign_channel_epg(new_epg: &mut Vec<Epg>, fp: &mut FetchedPlaylist<'_>
if !icon_assigned.contains(epg_channel_id) &&
(epg_source.logo_override || chan.header.logo.is_empty() || chan.header.logo_small.is_empty()) {
if let Some(icon_tag) = icon_tags.get(epg_channel_id) {
if let XmlTagIcon::Src(icon) = &icon_tag.icon {
if let Some(icon) = icon_tag.icon.as_ref() {
icon_assigned.insert(epg_channel_id.clone());
if epg_source.logo_override || chan.header.logo.is_empty() {
trace!("Matched channel {} to epg icon {icon}", chan.header.name);
chan.header.logo = icon.clone();
chan.header.logo = Arc::clone(icon);
}
if epg_source.logo_override || chan.header.logo_small.is_empty() {
chan.header.logo_small = icon.clone();
chan.header.logo_small = Arc::clone(icon);
}
}
}
+1 -1
View File
@@ -769,7 +769,7 @@ async fn process_playlist_for_target(ctx: &PlaylistProcessingContext,
if process_watch(&ctx.config, &ctx.client, target, &flat_new_playlist).await {
step.tick("group watches");
}
let result = persist_playlist(&ctx.config, &mut flat_new_playlist, flatten_tvguide(&new_epg).as_ref(), target, ctx.playlist_state.as_ref()).await;
let result = persist_playlist(&ctx.config, &mut flat_new_playlist, flatten_tvguide(new_epg).as_ref(), target, ctx.playlist_state.as_ref()).await;
step.stop("Persisting playlists");
result
}
+15 -79
View File
@@ -1,11 +1,10 @@
use crate::model::{Config, ConfigTarget, TargetOutput, XmlTagIcon};
use crate::model::{Epg, EPG_ATTRIB_CHANNEL, EPG_ATTRIB_ID, EPG_TAG_CHANNEL, EPG_TAG_DISPLAY_NAME, EPG_TAG_ICON, EPG_TAG_PROGRAMME};
use crate::model::{Config, ConfigTarget, TargetOutput};
use crate::model::{Epg};
use crate::repository::{m3u_get_epg_file_path_for_target, BPlusTree};
use crate::repository::{xtream_get_epg_file_path_for_target, xtream_get_storage_path};
use crate::utils::{debug_if_enabled, parse_xmltv_time};
use crate::utils::{debug_if_enabled};
use shared::error::{notify_err, TuliproxError};
use shared::model::{EpgChannel, EpgProgramme, PlaylistGroup};
use shared::utils::Internable;
use shared::model::{EpgChannel, PlaylistGroup};
use std::collections::HashMap;
use std::path::Path;
use std::sync::Arc;
@@ -25,17 +24,6 @@ pub const XML_PREAMBLE: &str = r#"<?xml version="1.0" encoding="utf-8"?>
// writer.write_event_async(quick_xml::events::Event::DocType(quick_xml::events::BytesText::new(r#"tv SYSTEM "xmltv.dtd""#)))
// .await.map_err(|e| notify_err!("failed to write doctype: {}", e))?;
pub fn epg_write_file(target: &ConfigTarget, epg: &Epg, path: &Path, playlist: Option<&[PlaylistGroup]>) -> Result<(), TuliproxError> {
let tag_channel = EPG_TAG_CHANNEL.intern();
let tag_programme = EPG_TAG_PROGRAMME.intern();
let tag_display_name = EPG_TAG_DISPLAY_NAME.intern();
let tag_icon = EPG_TAG_ICON.intern();
let tag_title = "title".intern();
let tag_desc = "desc".intern();
let epg_id_attrib = EPG_ATTRIB_ID.intern();
let channel_id_attrib = EPG_ATTRIB_CHANNEL.intern();
let start_attrib = "start".intern();
let stop_attrib = "stop".intern();
if epg.children.is_empty() {
return Ok(());
}
@@ -55,71 +43,19 @@ pub fn epg_write_file(target: &ConfigTarget, epg: &Epg, path: &Path, playlist: O
}
}
let mut channels: HashMap<Arc<str>, EpgChannel> =
epg.children
.iter()
.filter(|tag| tag.name == tag_channel)
.filter_map(|tag| {
let channel_id = tag.get_attribute_value(&epg_id_attrib)?;
let mut title = rename_map.get(channel_id).map(|v| Arc::clone(v));
let mut icon = match tag.icon {
XmlTagIcon::Src(ref url) => Some(Arc::clone(url)),
XmlTagIcon::Undefined | XmlTagIcon::Exists => None,
};
if let Some(children) = tag.children.as_ref() {
for child in children {
if child.name == tag_display_name {
if title.is_none() {
title.clone_from(&child.value);
}
} else if icon.is_none() && child.name == tag_icon {
icon.clone_from(&child.value);
}
}
}
let channel = EpgChannel {
id: Arc::clone(channel_id),
title,
icon,
programmes: vec![],
};
Some((Arc::clone(channel_id), channel))
})
.collect();
let mut tree = BPlusTree::<Arc<str>, EpgChannel>::new();
for channel in &epg.children {
if !channel.programmes.is_empty() {
let mut chan = (**channel).clone();
if let Some(&title) = rename_map.get(&chan.id) {
chan.title = Some(Arc::clone(title));
}
chan.programmes.sort_by_key(|p| p.start);
tree.insert(Arc::clone(&channel.id), chan);
}
}
drop(rename_map);
epg.children.iter().filter(|tag| tag.name == tag_programme).for_each(|tag| {
if let Some(attribs) = tag.attributes.as_ref() {
let opt_channel_id = attribs.get(&channel_id_attrib);
let opt_start = attribs.get(&start_attrib);
let opt_stop = attribs.get(&stop_attrib);
if let (Some(channel_id), Some(start), Some(stop)) = (opt_channel_id, opt_start, opt_stop) {
if let (Some(start_time), Some(stop_time)) = (parse_xmltv_time(start), parse_xmltv_time(stop)) {
if let Some(channel) = channels.get_mut(channel_id) {
let mut title = None;
let mut desc = None;
if let Some(children) = tag.children.as_ref() {
for child in children {
if child.name == tag_title {
title.clone_from(&child.value);
} else if child.name == tag_desc {
desc.clone_from(&child.value);
}
}
channel.programmes.push(EpgProgramme::new_all(start_time, stop_time, Arc::clone(channel_id), title, desc));
}
}
}
}
}
});
let mut tree = BPlusTree::<Arc<str>, EpgChannel>::new();
for (key, mut channel) in channels {
channel.programmes.sort_by_key(|p| p.start);
tree.insert(key, channel);
}
tree.store(path).map_err(|err| notify_err!("Failed to write epg for target {}: {} - {err}", target.name, path.display()))?;
debug_if_enabled!("Epg for target {} written to {}", target.name, path.display());
+80 -262
View File
@@ -1,5 +1,5 @@
use std::sync::Arc;
use chrono::{DateTime, Offset, TimeZone, Utc};
use chrono::{DateTime, TimeZone, Utc};
use chrono_tz::Tz;
use crate::model::{ConfigTarget, ProxyUserCredentials};
use shared::model::PlaylistItemType;
@@ -9,17 +9,15 @@ use crate::api::model::AppState;
/// Parses user-defined EPG timeshift configuration.
/// Supports either a numeric offset (e.g. "+2:30", "-1:15")
/// or a timezone name (e.g. "`Europe/Berlin`", "`UTC`", "`America/New_York`").
///
/// Returns the total offset in minutes (i32).
fn parse_timeshift(time_shift: Option<&str>) -> Option<i32> {
time_shift.and_then(|offset| {
fn parse_timeshift(time_shift: Option<&str>) -> EpgTimeShift {
if let Some(offset) = time_shift {
if offset.is_empty() {
return EpgTimeShift::None;
}
// Try to parse as timezone name first
if let Ok(tz) = offset.parse::<Tz>() {
// Determine the current UTC offset of that timezone (including DST)
let now = Utc::now();
let local_time = tz.from_utc_datetime(&now.naive_utc());
let offset_minutes = local_time.offset().fix().local_minus_utc() / 60;
return Some(offset_minutes);
return EpgTimeShift::TimeZone(tz);
}
// If not a timezone, try to parse as numeric offset
@@ -31,15 +29,28 @@ fn parse_timeshift(time_shift: Option<&str>) -> Option<i32> {
let minutes: i32 = parts.get(1).and_then(|m| m.parse().ok()).unwrap_or(0);
let total_minutes = hours * 60 + minutes;
(total_minutes > 0).then_some(sign_factor * total_minutes)
})
if total_minutes > 0 {
EpgTimeShift::Fixed(sign_factor * total_minutes)
} else {
EpgTimeShift::None
}
} else {
EpgTimeShift::None
}
}
#[derive(Debug, Clone)]
pub enum EpgTimeShift {
None,
Fixed(i32),
TimeZone(Tz),
}
#[derive(Debug, Clone)]
pub struct EpgProcessingOptions {
pub rewrite_urls: bool,
pub offset_minutes: i32,
pub time_shift: EpgTimeShift,
pub encrypt_secret: [u8; 16],
}
@@ -53,227 +64,14 @@ pub fn get_epg_processing_options(app_state: &Arc<AppState>, user: &ProxyUserCre
let redirect = user.proxy.is_redirect(PlaylistItemType::Live) || target.is_force_redirect(PlaylistItemType::Live);
let rewrite_urls = !redirect && rewrite_resources;
// Use 0 for timeshift if None
let timeshift = parse_timeshift(user.epg_timeshift.as_deref()).unwrap_or(0);
let timeshift = parse_timeshift(user.epg_timeshift.as_deref());
EpgProcessingOptions {
rewrite_urls,
offset_minutes: timeshift,
time_shift: timeshift,
encrypt_secret,
}
}
//
// pub trait EpgConsumer: Send {
// fn handle_event(&mut self, event: &Event<'_>, decoder: quick_xml::Decoder) -> impl std::future::Future<Output = Result<(), TuliproxError>> + Send;
// }
//
// pub struct EpgProcessor<R: AsyncBufRead + Send + Unpin> {
// reader: Reader<R>,
// epg_processing_options: EpgProcessingOptions,
// rewrite_base_url: String,
// filter_channel_id: Option<std::sync::Arc<str>>,
// limit: u32
// }
//
// impl<R: AsyncRead + Send + Unpin> EpgProcessor<BufReader<R>> {
// pub fn new(
// reader: R,
// epg_processing_options: EpgProcessingOptions,
// rewrite_base_url: String,
// filter_channel_id: Option<std::sync::Arc<str>>,
// limit: u32
// ) -> Self {
// Self {
// reader: Reader::from_reader(BufReader::new(reader)),
// epg_processing_options,
// rewrite_base_url,
// filter_channel_id,
// limit
// }
// }
// }
//
// #[allow(clippy::too_many_lines)]
// impl<R: AsyncBufRead + Send + Unpin> EpgProcessor<R> {
// pub async fn process<C: EpgConsumer>(&mut self, consumer: &mut C) -> Result<(), TuliproxError> {
// let mut buf = Vec::with_capacity(4096);
// let duration = Duration::minutes(i64::from(self.epg_processing_options.offset_minutes));
// let mut skip_depth = None;
//
// loop {
// buf.clear();
// let event = match self.reader.read_event_into_async(&mut buf).await {
// Ok(e) => e,
// Err(e) => {
// error!("Error reading epg XML event: {e}");
// return Err(info_err!("Error reading epg XML event: {}", e));
// }
// };
//
// if let Some(flt) = &self.filter_channel_id {
// match &event {
// Event::Start(e) => {
// if skip_depth.is_none() {
// let should_skip = match e.name().as_ref() {
// b"channel" => {
// e.attributes()
// .filter_map(Result::ok)
// .find(|a| a.key.as_ref() == b"id")
// .and_then(|a| a.unescape_value().ok())
// .is_some_and(|v| flt.as_ref() != v.as_ref())
// }
// b"programme" => {
// e.attributes()
// .filter_map(Result::ok)
// .find(|a| a.key.as_ref() == b"channel")
// .and_then(|a| a.unescape_value().ok())
// .is_some_and(|v| flt.as_ref() != v.as_ref())
// }
// _ => false,
// };
//
// if should_skip {
// skip_depth = Some(1);
// continue;
// }
// } else {
// skip_depth = skip_depth.map(|d| d + 1);
// continue;
// }
// }
// Event::End(_) => {
// if let Some(depth) = skip_depth {
// if depth == 1 {
// skip_depth = None;
// } else {
// skip_depth = Some(depth - 1);
// }
// continue;
// }
// }
// Event::Empty(_) => {
// if skip_depth.is_some() {
// continue;
// }
// }
// _ => {}
// }
//
// if skip_depth.is_some() {
// continue;
// }
// }
//
// match event {
// Event::Start(ref e) if self.epg_processing_options.offset_minutes != 0 && e.name().as_ref() == b"programme" => {
// let mut elem = BytesStart::new(EPG_TAG_PROGRAMME);
// for attr in e.attributes() {
// match attr {
// Ok(attr) if attr.key.as_ref() == b"start" => {
// if let Ok(start_value) = attr.decode_and_unescape_value(self.reader.decoder()) {
// elem.push_attribute(("start", time_correct(&start_value, &duration).as_str()));
// } else {
// elem.push_attribute(attr);
// }
// }
// Ok(attr) if attr.key.as_ref() == b"stop" => {
// if let Ok(stop_value) = attr.decode_and_unescape_value(self.reader.decoder()) {
// elem.push_attribute(("stop", time_correct(&stop_value, &duration).as_str()));
// } else {
// elem.push_attribute(attr);
// }
// }
// Ok(attr) => {
// elem.push_attribute(attr);
// }
// Err(e) => {
// error!("Error parsing epg attribute: {e}");
// }
// }
// }
// consumer.handle_event(&Event::Start(elem), self.reader.decoder()).await?;
// }
// ref event @ (Event::Empty(ref e) | Event::Start(ref e)) if self.epg_processing_options.rewrite_urls && e.name().as_ref() == b"icon" => {
// let mut elem = BytesStart::new(EPG_TAG_ICON);
// for attr in e.attributes() {
// match attr {
// Ok(attr) if attr.key.as_ref() == b"src" => {
// if let Some(icon) = get_attr_value_unescaped(&attr, self.reader.decoder()) {
// if icon.is_empty() {
// elem.push_attribute(attr);
// } else {
// let rewritten_url = if let Ok(encrypted) = obscure_text(&self.epg_processing_options.encrypt_secret, &icon) {
// format!("{}{}", self.rewrite_base_url, encrypted)
// } else {
// icon
// };
// elem.push_attribute(("src", rewritten_url.as_str()));
// }
// } else {
// elem.push_attribute(attr);
// }
// }
// Ok(attr) => {
// elem.push_attribute(attr);
// }
// Err(e) => {
// error!("Error parsing epg attribute: {e}");
// }
// }
// }
//
// let out_event = match event {
// Event::Empty(_) => Event::Empty(elem),
// Event::Start(_) => Event::Start(elem),
// _ => unreachable!(),
// };
// consumer.handle_event(&out_event, self.reader.decoder()).await?;
// }
// Event::Decl(_) | Event::DocType(_) => {},
// Event::Eof => break,
// ref event => {
// consumer.handle_event(event, self.reader.decoder()).await?;
// }
// }
// }
// Ok(())
// }
// }
//
// /// # Panics
// /// unwrap for `FixedOffset` should not panic!
// pub fn time_correct(original: &str, shift: &Duration) -> String {
// let (datetime_part, tz_part) = if let Some((dt, tz)) = original.trim().rsplit_once(' ') {
// (dt, tz)
// } else {
// (original.trim(), "+0000")
// };
//
// let Ok(naive_dt) = NaiveDateTime::parse_from_str(datetime_part, "%Y%m%d%H%M%S") else { return original.to_string() };
//
// let tz_offset_minutes = if tz_part.len() == 5 {
// let bytes = tz_part.as_bytes();
// let sign = if bytes.first() == Some(&b'-') { -1 } else { 1 };
// let hours: i32 = tz_part.get(1..3).and_then(|s| s.parse().ok()).unwrap_or(0);
// let mins: i32 = tz_part.get(3..5).and_then(|s| s.parse().ok()).unwrap_or(0);
// sign * (hours * 60 + mins)
// } else {
// 0
// };
//
// let tz = FixedOffset::east_opt(tz_offset_minutes * 60).unwrap_or(FixedOffset::east_opt(0).unwrap()); // should not panic
//
// let dt: DateTime<FixedOffset> = tz
// .from_local_datetime(&naive_dt)
// .single()
// .unwrap_or_else(|| tz.from_utc_datetime(&naive_dt));
//
// let shifted_dt = dt + *shift;
//
// format!("{} {}", shifted_dt.format("%Y%m%d%H%M%S"), format_offset(tz_offset_minutes))
// }
/// # Panics
/// unwrap for `FixedOffset` should not panic!
pub fn apply_offset(ts_utc: i64, offset_minutes: i32) -> i64 {
ts_utc + i64::from(offset_minutes) * 60
}
@@ -292,19 +90,23 @@ pub fn parse_xmltv_time(t: &str) -> Option<i64> {
.map(|dt| dt.with_timezone(&Utc).timestamp())
}
pub fn format_xmltv_time_utc(ts: i64, offset_minutes: i32) -> String {
pub fn format_xmltv_time_utc(ts: i64, time_shift: &EpgTimeShift) -> String {
let dt = Utc.timestamp_opt(ts, 0).unwrap();
if offset_minutes == 0 {
dt.format("%Y%m%d%H%M%S %z").to_string()
} else {
match chrono::FixedOffset::east_opt(offset_minutes * 60) {
Some(offset) => {
dt.with_timezone(&offset).format("%Y%m%d%H%M%S %z").to_string()
}
None => {
dt.format("%Y%m%d%H%M%S %z").to_string()
match time_shift {
EpgTimeShift::None => dt.format("%Y%m%d%H%M%S %z").to_string(),
EpgTimeShift::Fixed(minutes) => {
match chrono::FixedOffset::east_opt(minutes * 60) {
Some(offset) => {
dt.with_timezone(&offset).format("%Y%m%d%H%M%S %z").to_string()
}
None => {
dt.format("%Y%m%d%H%M%S %z").to_string()
}
}
}
EpgTimeShift::TimeZone(tz) => {
dt.with_timezone(tz).format("%Y%m%d%H%M%S %z").to_string()
}
}
}
@@ -314,36 +116,52 @@ mod tests {
#[test]
fn test_parse_timeshift() {
assert_eq!(parse_timeshift(Some(&String::from("2"))), Some(120));
assert_eq!(parse_timeshift(Some(&String::from("-1:30"))), Some(-90));
assert_eq!(parse_timeshift(Some(&String::from("+0:15"))), Some(15));
assert_eq!(parse_timeshift(Some(&String::from("1:45"))), Some(105));
assert_eq!(parse_timeshift(Some(&String::from(":45"))), Some(45));
assert_eq!(parse_timeshift(Some(&String::from("-:45"))), Some(-45));
assert_eq!(parse_timeshift(Some(&String::from("0:30"))), Some(30));
assert_eq!(parse_timeshift(Some(&String::from(":3"))), Some(3));
assert_eq!(parse_timeshift(Some(&String::from("2:"))), Some(120));
assert_eq!(parse_timeshift(Some(&String::from("+2:00"))), Some(120));
assert_eq!(parse_timeshift(Some(&String::from("-0:10"))), Some(-10));
assert_eq!(parse_timeshift(Some(&String::from("invalid"))), None);
assert_eq!(parse_timeshift(Some(&String::from("+abc"))), None);
assert_eq!(parse_timeshift(Some(&String::new())), None);
assert_eq!(parse_timeshift(None), None);
assert!(matches!(parse_timeshift(Some(&String::from("2"))), EpgTimeShift::Fixed(120)));
assert!(matches!(parse_timeshift(Some(&String::from("-1:30"))), EpgTimeShift::Fixed(-90)));
assert!(matches!(parse_timeshift(Some(&String::from("+0:15"))), EpgTimeShift::Fixed(15)));
assert!(matches!(parse_timeshift(Some(&String::from("1:45"))), EpgTimeShift::Fixed(105)));
assert!(matches!(parse_timeshift(Some(&String::from(":45"))), EpgTimeShift::Fixed(45)));
assert!(matches!(parse_timeshift(Some(&String::from("-:45"))), EpgTimeShift::Fixed(-45)));
assert!(matches!(parse_timeshift(Some(&String::from("0:30"))), EpgTimeShift::Fixed(30)));
assert!(matches!(parse_timeshift(Some(&String::from(":3"))), EpgTimeShift::Fixed(3)));
assert!(matches!(parse_timeshift(Some(&String::from("2:"))), EpgTimeShift::Fixed(120)));
assert!(matches!(parse_timeshift(Some(&String::from("+2:00"))), EpgTimeShift::Fixed(120)));
assert!(matches!(parse_timeshift(Some(&String::from("-0:10"))), EpgTimeShift::Fixed(-10)));
assert!(matches!(parse_timeshift(Some(&String::from("invalid"))), EpgTimeShift::None));
assert!(matches!(parse_timeshift(Some(&String::from("+abc"))), EpgTimeShift::None));
assert!(matches!(parse_timeshift(Some(&String::new())), EpgTimeShift::None));
assert!(matches!(parse_timeshift(None), EpgTimeShift::None));
}
#[test]
fn test_parse_timezone() {
// This will depend on current DST; we just check it’s within a valid range
let berlin = parse_timeshift(Some(&"Europe/Berlin".to_string())).unwrap();
assert!(berlin == 60 || berlin == 120, "Berlin offset should be 60 or 120, got {berlin}");
// Check timezone parsing creates the correct variant
let amterdam = parse_timeshift(Some(&"Europe/Amsterdam".to_string()));
if let EpgTimeShift::TimeZone(tz) = amterdam {
assert_eq!(tz.name(), "Europe/Amsterdam");
} else {
panic!("Expected TimeZone for Europe/Amsterdam");
}
let new_york = parse_timeshift(Some(&"America/New_York".to_string())).unwrap();
assert!(new_york == -300 || new_york == -240, "New York offset should be -300 or -240, got {new_york}");
let new_york = parse_timeshift(Some(&"America/New_York".to_string()));
if let EpgTimeShift::TimeZone(tz) = new_york {
assert_eq!(tz.name(), "America/New_York");
} else {
panic!("Expected TimeZone for America/New_York");
}
let tokyo = parse_timeshift(Some(&"Asia/Tokyo".to_string())).unwrap();
assert_eq!(tokyo, 540); // always UTC+9
let tokyo = parse_timeshift(Some(&"Asia/Tokyo".to_string()));
if let EpgTimeShift::TimeZone(tz) = tokyo {
assert_eq!(tz.name(), "Asia/Tokyo");
} else {
panic!("Expected TimeZone for Asia/Tokyo");
}
let utc = parse_timeshift(Some(&"UTC".to_string())).unwrap();
assert_eq!(utc, 0);
let utc = parse_timeshift(Some(&"UTC".to_string()));
if let EpgTimeShift::TimeZone(tz) = utc {
assert_eq!(tz.name(), "UTC");
} else {
panic!("Expected TimeZone for UTC");
}
}
}
+1 -1
View File
@@ -452,7 +452,7 @@ pub struct M3uPlaylistItem {
#[serde(with = "arc_str_serde")]
pub input_name: Arc<str>,
pub item_type: PlaylistItemType,
#[serde(with = "arc_str_serde")]
#[serde(skip_serializing, default)]
pub t_stream_url: Arc<str>,
#[serde(skip)]
pub t_resource_url: Option<String>,