From 38b1cd6d313ca2515ee9f17c8575150a652efb72 Mon Sep 17 00:00:00 2001 From: euzu Date: Wed, 1 Oct 2025 14:40:44 +0200 Subject: [PATCH] Fixed invalid epg text content parsing --- CHANGELOG.md | 1 + backend/Cargo.toml | 2 +- backend/src/api/endpoints/xmltv_api.rs | 143 +++++++++--------- backend/src/model/xmltv.rs | 20 ++- backend/src/processing/parser/xmltv.rs | 95 +++++++----- backend/src/processing/processor/epg.rs | 8 +- backend/src/processing/processor/playlist.rs | 6 +- .../compressed_file_reader_async.rs | 63 ++++++++ backend/src/utils/compression/mod.rs | 1 + 9 files changed, 217 insertions(+), 122 deletions(-) create mode 100644 backend/src/utils/compression/compressed_file_reader_async.rs diff --git a/CHANGELOG.md b/CHANGELOG.md index ec965951a..8fb69d7d0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -11,6 +11,7 @@ to reduce transient HTTP errors (400, 408, 425, 429, 5xx) - WebSocket now reconnects on disconnect; added WebSocket connection status icon in Web UI - Added Playlist EPG view with timeline, channels, `now` line, and program details - EPG data can now be fetched from selected targets and custom URLs +- Invalid EPG text data fix - Added new sidebar entry and icon for quick EPG access - Added CBOR (binary JSON) support for large API data diff --git a/backend/Cargo.toml b/backend/Cargo.toml index fba18c997..47e47e319 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -63,7 +63,7 @@ dashmap = "6.1" hyper = "1.7" hyper-util = "0.1" socket2 = { version = "0.6", features = ["all"] } -async-compression = "0" +async-compression = { version = "0.4", features = ["tokio"] } #[cfg(target_os = "macos")] libc = "0.2" #[cfg(target_os = "windows")] diff --git a/backend/src/api/endpoints/xmltv_api.rs b/backend/src/api/endpoints/xmltv_api.rs index 324524f08..66e5af570 100644 --- a/backend/src/api/endpoints/xmltv_api.rs +++ b/backend/src/api/endpoints/xmltv_api.rs @@ -1,19 +1,16 @@ use axum::response::IntoResponse; use chrono::{Duration, NaiveDateTime, TimeDelta}; -use flate2::write::GzEncoder; -use flate2::Compression; use log::{error, trace}; use quick_xml::events::{BytesStart, Event}; -use quick_xml::{Reader, Writer}; -use std::fs::File; use std::path::{Path, PathBuf}; use std::sync::Arc; - +use tokio::io::AsyncWriteExt; +use tokio_util::io::ReaderStream; use crate::api::api_utils::try_unwrap_body; use crate::api::api_utils::{get_user_target, serve_file}; use crate::api::model::AppState; use crate::api::model::UserApiRequest; -use crate::model::Config; +use crate::model::{Config, EPG_TAG_PROGRAMME}; use crate::model::{ConfigTarget, ProxyUserCredentials, TargetOutput}; use crate::repository::m3u_repository::m3u_get_epg_file_path; use crate::repository::storage::get_target_storage_path; @@ -105,7 +102,7 @@ async fn serve_epg( epg_path: &Path, user: &ProxyUserCredentials, ) -> impl axum::response::IntoResponse + Send { - match File::open(epg_path) { + match tokio::fs::File::open(epg_path).await { Ok(epg_file) => match parse_timeshift(user.epg_timeshift.as_ref()) { None => serve_file(epg_path, mime::TEXT_XML).await.into_response(), Some(duration) => serve_epg_with_timeshift(epg_file, duration).into_response(), @@ -115,84 +112,84 @@ async fn serve_epg( } fn serve_epg_with_timeshift( - epg_file: File, + epg_file: tokio::fs::File, offset_minutes: i32, ) -> impl axum::response::IntoResponse + Send { - let reader = utils::file_reader(epg_file); - let encoder = GzEncoder::new(Vec::with_capacity(4096), Compression::default()); - let mut xml_reader = Reader::from_reader(reader); - let mut xml_writer = Writer::new(encoder); - let mut buf = Vec::with_capacity(1024); - let duration = Duration::minutes(i64::from(offset_minutes)); + let reader = tokio::io::BufReader::new(epg_file); + let (tx, rx) = tokio::io::duplex(8192); + tokio::spawn(async move { + let encoder = async_compression::tokio::write::GzipEncoder::new(tx); + let mut xml_reader = quick_xml::reader::Reader::from_reader(tokio::io::BufReader::new(reader)); + let mut xml_writer = quick_xml::writer::Writer::new(encoder); + let mut buf = Vec::with_capacity(4096); + let duration = Duration::minutes(i64::from(offset_minutes)); - loop { - match xml_reader.read_event_into(&mut buf) { - Ok(Event::Start(ref e)) if e.name().as_ref() == b"programme" => { - // Modify the attributes - let mut elem = BytesStart::new("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(xml_reader.decoder()) { - // Modify the start attribute value as needed - elem.push_attribute(("start", time_correct(&start_value, &duration).as_str())); - } else { - // keep original attribute unchanged ? + loop { + match xml_reader.read_event_into_async(&mut buf).await { + Ok(Event::Start(ref e)) if e.name().as_ref() == b"programme" => { + // Modify the attributes + 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(xml_reader.decoder()) { + // Modify the start attribute value as needed + elem.push_attribute(("start", time_correct(&start_value, &duration).as_str())); + } else { + // keep original attribute unchanged ? + elem.push_attribute(attr); + } + } + Ok(attr) if attr.key.as_ref() == b"stop" => { + if let Ok(stop_value) = attr.decode_and_unescape_value(xml_reader.decoder()) { + // Modify the stop attribute value as needed + elem.push_attribute(("stop", time_correct(&stop_value, &duration).as_str())); + } else { + elem.push_attribute(attr); + } + } + Ok(attr) => { + // Copy any other attributes as they are elem.push_attribute(attr); } - } - Ok(attr) if attr.key.as_ref() == b"stop" => { - if let Ok(stop_value) = attr.decode_and_unescape_value(xml_reader.decoder()) { - // Modify the stop attribute value as needed - elem.push_attribute(("stop", time_correct(&stop_value, &duration).as_str())); - } else { - elem.push_attribute(attr); + Err(e) => { + error!("Error parsing attribute: {e}"); } } - Ok(attr) => { - // Copy any other attributes as they are - elem.push_attribute(attr); - } - Err(e) => { - error!("Error parsing attribute: {e}"); - } } + + // Write the modified start event + xml_writer + .write_event_async(Event::Start(elem)).await + .expect("Failed to write event"); } + Ok(Event::Eof) => break, // End of file + Ok(event) => { + // Write any other event as is + xml_writer + .write_event_async(event).await + .expect("Failed to write event"); + } + Err(e) => { + error!("Error: {e}"); + break; + } + } - // Write the modified start event - xml_writer - .write_event(Event::Start(elem)) - .expect("Failed to write event"); - } - Ok(Event::Eof) => break, // End of file - Ok(event) => { - // Write any other event as is - xml_writer - .write_event(event) - .expect("Failed to write event"); - } - Err(e) => { - error!("Error: {e}"); - break; - } + buf.clear(); } + let _ = xml_writer.into_inner().shutdown().await; + }); - buf.clear(); - } - match xml_writer.into_inner().finish() { - Ok(compressed_data) => try_unwrap_body!(axum::response::Response::builder() - .header( - axum::http::header::CONTENT_TYPE, - mime::TEXT_XML.to_string() - ) - .header(axum::http::header::CONTENT_ENCODING, "gzip") // Set Content-Encoding header - .body(axum::body::Body::from(compressed_data))), - Err(err) => ( - axum::http::StatusCode::INTERNAL_SERVER_ERROR, - err.to_string(), + let body_stream = ReaderStream::new(rx); + try_unwrap_body!(axum::response::Response::builder() + .header( + axum::http::header::CONTENT_TYPE, + mime::TEXT_XML.to_string() ) - .into_response(), - } + .header(axum::http::header::CONTENT_ENCODING, "gzip") // Set Content-Encoding header + .body(axum::body::Body::from_stream(body_stream))) + .into_response() } /// Handles XMLTV EPG API requests, serving the appropriate EPG file with optional time-shifting based on user configuration. @@ -247,8 +244,6 @@ pub fn xmltv_api_register() -> axum::Router> { #[cfg(test)] mod tests { - use std::io::BufReader; - use shared::model::{PlaylistCategoriesResponse}; use super::*; #[test] diff --git a/backend/src/model/xmltv.rs b/backend/src/model/xmltv.rs index 9796410df..a9ef618b5 100644 --- a/backend/src/model/xmltv.rs +++ b/backend/src/model/xmltv.rs @@ -191,6 +191,18 @@ pub async fn parse_xmltv_for_web_ui_from_url(app_state: &Arc, url: &st } } +fn concat_text(t1: &String, t2: &str) -> String { + if t1.is_empty() { + t2.to_string() + } else if t1.ends_with('\\') { + let mut t = t1.to_string(); + t.pop(); + format!("{t}'{t2}") + } else { + format!("{t1}{t2}") + } +} + #[allow(clippy::too_many_lines)] async fn parse_xmltv_for_web_ui(reader: R) -> Result { @@ -286,14 +298,14 @@ async fn parse_xmltv_for_web_ui(reader: R) -> Resul let text = decoded.trim(); if !text.is_empty() { if let Some(channel) = &mut current_channel { - if current_tag == EPG_TAG_DISPLAY_NAME && channel.title.is_empty() { - channel.title = text.to_string(); + if current_tag == EPG_TAG_DISPLAY_NAME { + channel.title = concat_text(&channel.title, text); } } if let Some(program) = &mut current_programme { - if current_tag == "title" && program.title.is_empty() { - program.title = text.to_string(); + if current_tag == "title" { + program.title = concat_text(&program.title, text); } } } diff --git a/backend/src/processing/parser/xmltv.rs b/backend/src/processing/parser/xmltv.rs index 2d6905b01..07336b838 100644 --- a/backend/src/processing/parser/xmltv.rs +++ b/backend/src/processing/parser/xmltv.rs @@ -1,11 +1,9 @@ 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::compressed_file_reader::CompressedFileReader; use dashmap::DashMap; use deunicode::deunicode; use quick_xml::events::{BytesStart, BytesText, Event}; -use quick_xml::Reader; use rayon::iter::{IntoParallelRefIterator, ParallelIterator}; use shared::model::EpgNamePrefix; use shared::utils::CONSTANTS; @@ -14,6 +12,8 @@ use std::cmp::min; use std::collections::HashMap; use std::mem; use std::sync::{Mutex}; +use tokio::io::AsyncRead; +use crate::utils::compressed_file_reader_async::CompressedFileReaderAsync; /// Splits a string at the first delimiter if the prefix matches a known country code. /// @@ -254,8 +254,8 @@ impl TVGuide { /// assert!(!epg.children.is_empty()); /// } /// ``` - fn process_epg_file(id_cache: &mut EpgIdCache, epg_source: &PersistedEpgSource) -> Option { - match CompressedFileReader::new(&epg_source.file_path) { + async fn process_epg_file(id_cache: &mut EpgIdCache<'_>, epg_source: &PersistedEpgSource) -> Option { + match CompressedFileReaderAsync::new(&epg_source.file_path).await { Ok(mut reader) => { let mut children: Vec = vec![]; let mut tv_attributes: Option> = None; @@ -298,7 +298,7 @@ impl TVGuide { } }; - parse_tvguide(&mut reader, &mut filter_tags); + parse_tvguide(&mut reader, &mut filter_tags).await; if children.is_empty() { return None; @@ -315,13 +315,16 @@ impl TVGuide { } } - pub fn filter(&self, id_cache: &mut EpgIdCache) -> Option> { + pub async fn filter(&self, id_cache: &mut EpgIdCache<'_>) -> Option> { if id_cache.channel_epg_id.is_empty() && id_cache.normalized.is_empty() { return None; } - let mut epg_sources: Vec = self.get_epg_sources().iter() - .filter_map(|epg_source| Self::process_epg_file(id_cache, epg_source)) - .collect(); + let mut epg_sources: Vec = vec![]; + for epg_source in self.get_epg_sources() { + if let Some(epg) = Self::process_epg_file(id_cache, epg_source).await { + epg_sources.push(epg); + } + } epg_sources.sort_by(|a, b| a.priority.cmp(&b.priority)); Some(epg_sources) } @@ -383,28 +386,36 @@ where } fn handle_text_tag(stack: &mut [XmlTag], e: &BytesText) { - if !stack.is_empty() { + if let Some(tag) = stack.last_mut() { if let Ok(text) = e.decode() { let t = text.trim(); if !t.is_empty() { - if let Some(tag) = stack.last_mut() { - tag.value = Some(t.to_string()); - } + let t_fixed: Cow = if t.ends_with('\\') { + let mut owned = t.to_string(); + owned.pop(); + owned.push('\''); + Cow::Owned(owned) + } else { + Cow::Borrowed(t) + }; + + let old = tag.value.get_or_insert_with(String::new); + old.push_str(&t_fixed); } } } } -pub fn parse_tvguide(content: R, callback: &mut F) +pub async fn parse_tvguide(content: R, callback: &mut F) where - R: std::io::BufRead, + R: AsyncRead + Unpin, F: FnMut(XmlTag), { let mut stack: Vec = vec![]; - let mut reader = Reader::from_reader(content); + let mut xml_reader = quick_xml::reader::Reader::from_reader(tokio::io::BufReader::new(content)); let mut buf = Vec::::new(); loop { - match reader.read_event_into(&mut buf) { + match xml_reader.read_event_into_async(&mut buf).await { Ok(Event::Eof) => break, Ok(Event::Start(e)) => handle_tag_start(callback, &mut stack, &e), Ok(Event::Empty(e)) => { @@ -518,7 +529,11 @@ pub fn flatten_tvguide(tv_guides: &[Epg]) -> Option { #[cfg(test)] mod tests { - use crate::model::EpgSmartMatchConfig; + use std::borrow::Cow; + use std::collections::{HashSet}; + use std::io; + use std::path::PathBuf; + use crate::model::{EpgSmartMatchConfig, PersistedEpgSource, TVGuide}; use crate::processing::parser::xmltv::normalize_channel_name; #[test] @@ -537,24 +552,30 @@ mod tests { } - // #[test] - // fn parse_test() -> io::Result<()> { - // let file_path = PathBuf::from("/tmp/epg.xml.gz"); - // - // if file_path.exists() { - // let tv_guide = TVGuide { file: file_path }; - // - // let mut channel_ids = HashSet::from(["channel.1".to_string(), "channel.2".to_string(), "channel.3".to_string()]); - // let mut nomalized = HashMap::new(); - // match tv_guide.filter(&mut channel_ids, &mut nomalized) { - // None => assert!(false, "No epg filtered"), - // Some(epg) => { - // assert_eq!(epg.children.len(), channel_ids.len() * 2, "Epg size does not match") - // } - // } - // } - // Ok(()) - // } + #[test] + fn parse_test() -> io::Result<()> { + //let file_path = PathBuf::from("/tmp/epg.xml.gz"); + let file_path = PathBuf::from("/tmp/invalid_epg.xml"); + + if file_path.exists() { + let tv_guide = TVGuide::new(vec![PersistedEpgSource { file_path, priority: 0, logo_override: false }]); + + let mut id_cache = EpgIdCache::new(None); + id_cache.channel_epg_id.insert(Cow::Owned("342".to_string())); + //id_cache.collect_epg_id(fp); + + let channel_ids = HashSet::from(["342".to_string()]); + match tv_guide.filter(&mut id_cache) { + None => assert!(false, "No epg filtered"), + Some(epgs) => { + for epg in epgs { + assert_eq!(epg.children.len(), channel_ids.len() * 2, "Epg size does not match") + } + } + } + } + Ok(()) + } #[test] /// Tests normalization of channel names with various prefixes, suffixes, and special characters using a configured `EpgSmartMatchConfig`. @@ -580,6 +601,7 @@ mod tests { use rphonetic::{Encoder, Metaphone}; use shared::model::{EpgNamePrefix, EpgSmartMatchConfigDto}; + use crate::processing::processor::epg::EpgIdCache; #[test] /// Demonstrates phonetic encoding (Metaphone) of normalized channel names with various prefixes and suffixes. @@ -612,4 +634,5 @@ mod tests { println!("{}", metaphone.encode(&normalize_channel_name("BU | ODISEA ᵁᴴᴰ ³⁸⁴⁰ᴾ", &epg_smart_cfg))); println!("{}", metaphone.encode(&normalize_channel_name("BG | ODISEA ᵁᴴᴰ ³⁸⁴⁰ᴾ", &epg_smart_cfg))); } + } \ No newline at end of file diff --git a/backend/src/processing/processor/epg.rs b/backend/src/processing/processor/epg.rs index 468730d14..0ed45b56d 100644 --- a/backend/src/processing/processor/epg.rs +++ b/backend/src/processing/processor/epg.rs @@ -153,11 +153,11 @@ impl EpgIdCache<'_> { /// let mut id_cache = EpgIdCache::new(None); /// assign_channel_epg(&mut new_epg, &mut playlist, &mut id_cache); /// ``` -fn assign_channel_epg(new_epg: &mut Vec, fp: &mut FetchedPlaylist, id_cache: &mut EpgIdCache) { +async fn assign_channel_epg(new_epg: &mut Vec, fp: &mut FetchedPlaylist<'_>, id_cache: &mut EpgIdCache<'_>) { //id_cache.normalized.retain(|_, v| v.is_some()); if let Some(tv_guide) = &fp.epg { let mut processed_epgs = vec![]; - if let Some(epg_sources) = tv_guide.filter(id_cache) { + if let Some(epg_sources) = tv_guide.filter(id_cache).await { let mut icon_assigned = HashSet::new(); for epg_source in epg_sources { // icon tags @@ -236,7 +236,7 @@ fn assign_channel_epg(new_epg: &mut Vec, fp: &mut FetchedPlaylist, id_cache /// let mut epg_data = Vec::new(); /// process_playlist_epg(&mut playlist, &mut epg_data); /// ``` -pub fn process_playlist_epg(fp: &mut FetchedPlaylist, epg: &mut Vec) { +pub async fn process_playlist_epg(fp: &mut FetchedPlaylist<'_>, epg: &mut Vec) { // collect all epg_channel ids let mut id_cache = EpgIdCache::new(fp.input.epg.as_ref()); id_cache.collect_epg_id(fp); @@ -244,7 +244,7 @@ pub fn process_playlist_epg(fp: &mut FetchedPlaylist, epg: &mut Vec) { if id_cache.is_empty() && !id_cache.smart_match_enabled { debug!("No epg ids found"); } else { - assign_channel_epg(epg, fp, &mut id_cache); + assign_channel_epg(epg, fp, &mut id_cache).await; } } diff --git a/backend/src/processing/processor/playlist.rs b/backend/src/processing/processor/playlist.rs index 962733bd8..8ca8c8337 100644 --- a/backend/src/processing/processor/playlist.rs +++ b/backend/src/processing/processor/playlist.rs @@ -499,7 +499,7 @@ async fn process_playlist_for_target(app_config: &AppConfig, processed_fetched_playlists.push(processed_fpl); } step.tick("filter rename map"); - let (new_epg, mut new_playlist) = process_epg(&mut processed_fetched_playlists); + let (new_epg, mut new_playlist) = process_epg(&mut processed_fetched_playlists).await; step.tick("epg"); if new_playlist.is_empty() { @@ -552,14 +552,14 @@ async fn trakt_playlist(client: &Arc, target: &ConfigTarget, errors: &mu true } -fn process_epg(processed_fetched_playlists: &mut Vec) -> (Vec, Vec) { +async fn process_epg(processed_fetched_playlists: &mut Vec>) -> (Vec, Vec) { let mut new_playlist = vec![]; let mut new_epg = vec![]; // each fetched playlist can have its own epgl url. // we need to process each input epg. for fp in processed_fetched_playlists { - process_playlist_epg(fp, &mut new_epg); + process_playlist_epg(fp, &mut new_epg).await; new_playlist.append(&mut fp.playlistgroups); } (new_epg, new_playlist) diff --git a/backend/src/utils/compression/compressed_file_reader_async.rs b/backend/src/utils/compression/compressed_file_reader_async.rs new file mode 100644 index 000000000..cb695f4af --- /dev/null +++ b/backend/src/utils/compression/compressed_file_reader_async.rs @@ -0,0 +1,63 @@ +use std::path::Path; +use std::pin::Pin; +use std::task::{Context, Poll}; +use tokio::fs::File; +use tokio::io::{ + self, AsyncRead, BufReader, AsyncSeekExt, AsyncReadExt, ReadBuf, +}; +use async_compression::tokio::bufread::{GzipDecoder, ZlibDecoder}; + +use crate::utils::compression::compression_utils::{is_deflate, is_gzip}; + +pub struct CompressedFileReaderAsync { + reader: BufReader>, +} + +impl CompressedFileReaderAsync { + pub async fn new(path: &Path) -> std::io::Result { + let file: File = tokio::fs::File::open(path).await?; + + let mut buffered_file = BufReader::new(file); + let mut header = [0u8; 2]; + buffered_file.read_exact(&mut header).await?; + buffered_file.seek(io::SeekFrom::Start(0)).await?; + + let reader: Box = if is_gzip(&header) { + Box::new(GzipDecoder::new(buffered_file)) + } else if is_deflate(&header) { + Box::new(ZlibDecoder::new(buffered_file)) + } else { + Box::new(buffered_file) + }; + + Ok(Self { + reader: BufReader::new(reader), + }) + } +} + +impl AsyncRead for CompressedFileReaderAsync { + fn poll_read( + mut self: Pin<&mut Self>, + cx: &mut Context<'_>, + buf: &mut ReadBuf<'_>, + ) -> Poll> { + Pin::new(&mut self.reader).poll_read(cx, buf) + } +} +// +// impl AsyncBufRead for CompressedFileReaderAsync { +// fn poll_fill_buf( +// self: Pin<&mut Self>, +// cx: &mut Context<'_>, +// ) -> Poll> { +// unsafe { +// let this = self.get_unchecked_mut(); +// Pin::new_unchecked(&mut this.reader).poll_fill_buf(cx) +// } +// } +// +// fn consume(mut self: Pin<&mut Self>, amt: usize) { +// Pin::new(&mut self.reader).consume(amt) +// } +// } diff --git a/backend/src/utils/compression/mod.rs b/backend/src/utils/compression/mod.rs index d28df9813..0014d2cf0 100644 --- a/backend/src/utils/compression/mod.rs +++ b/backend/src/utils/compression/mod.rs @@ -1,2 +1,3 @@ pub mod compressed_file_reader; pub mod compression_utils; +pub mod compressed_file_reader_async;