diff --git a/CHANGELOG.md b/CHANGELOG.md index 1e62c193f..10cac4952 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -28,7 +28,16 @@ web_ui: userfile: user.txt ``` - user has now the attribute `ui_enabled` to disable/enable web_ui for user. -- epg processing optimization, auto guessing/assigning epg id's +- !BREAKING_CHANGE! multi epg processing/optimization, auto guessing/assigning epg id's +```yaml +epg_url: ['https://localhost.com/epg.xml'] +``` +```yaml +epg_url: ['file:///${env:M3U_FILTER_HOME}/epg.xml', 'file:///${env:M3U_FILTER_HOME}/epg2.xml'] +``` +```yaml +epg_url: ['http://localhost:3001/xmltv.php?epg_id=1', 'http://localhost:3001/xmltv.php?epg_id=2'] +``` # 2.2.5 (2025-03-27) - fixed web ui playlist regexp search diff --git a/src/model/config.rs b/src/model/config.rs index f3c3c88af..4f8d34c20 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -768,7 +768,7 @@ pub struct ConfigInput { pub headers: HashMap, pub url: String, #[serde(default, skip_serializing_if = "Option::is_none")] - pub epg_url: Option, + pub epg_url: Option>, #[serde(default, skip_serializing_if = "Option::is_none")] pub username: Option, #[serde(default, skip_serializing_if = "Option::is_none")] @@ -808,6 +808,12 @@ impl ConfigInput { check_input_credentials!(self, self.input_type); self.persist = get_trimmed_string(&self.persist); + self.epg_url = self.epg_url.take().map(|list| { + list.into_iter() + .map(|url| url.trim().to_string()) + .filter(|s| !s.is_empty()) + .collect() + }); if let Some(aliases) = self.aliases.as_mut() { let input_type = &self.input_type; diff --git a/src/model/xmltv.rs b/src/model/xmltv.rs index 1a3289248..ea1ab869b 100644 --- a/src/model/xmltv.rs +++ b/src/model/xmltv.rs @@ -68,5 +68,5 @@ impl Epg { #[derive(Debug, Clone)] pub struct TVGuide { - pub file: PathBuf, + pub file_paths: Vec, } diff --git a/src/processing/parser/xmltv.rs b/src/processing/parser/xmltv.rs index 186c70848..fc3c63ea5 100644 --- a/src/processing/parser/xmltv.rs +++ b/src/processing/parser/xmltv.rs @@ -1,9 +1,11 @@ +use std::borrow::Cow; use crate::model::xmltv::{Epg, TVGuide, XmlTag, EPG_ATTRIB_CHANNEL, EPG_ATTRIB_ID, EPG_TAG_CHANNEL, EPG_TAG_DISPLAY_NAME, EPG_TAG_ICON, EPG_TAG_PROGRAMME, EPG_TAG_TV}; use crate::utils::compression::compressed_file_reader::CompressedFileReader; use quick_xml::events::{BytesStart, Event}; use quick_xml::Reader; use regex::Regex; use std::collections::{HashMap, HashSet}; +use std::path::Path; use std::sync::{Arc, LazyLock}; use unidecode::unidecode; @@ -26,11 +28,20 @@ pub fn normalize_channel_name(name: &str) -> String { } impl TVGuide { - pub fn filter(&self, epg_channel_ids: &mut HashSet, normalized_epg_channel_ids: &mut HashMap>) -> Option { - if epg_channel_ids.is_empty() && normalized_epg_channel_ids.is_empty() { + fn merge(mut epgs: Vec) -> Option { + if epgs.is_empty() { return None; } - match CompressedFileReader::new(&self.file) { + let first_epg_attributes = epgs.get_mut(0).unwrap().attributes.take(); + let merged_children: Vec = epgs.into_iter().flat_map(|epg| epg.children).collect(); + Some(Epg { + attributes: first_epg_attributes, + children: merged_children, + }) + } + + fn process_epg_file(epg_channel_ids: &mut HashSet>, normalized_epg_channel_ids: &mut HashMap, Option>>, processed_epg_channel_ids: &mut HashSet, epg_file: &Path) -> Option { + match CompressedFileReader::new(epg_file) { Ok(mut reader) => { let mut children: Vec = vec![]; let mut tv_attributes: Option>> = None; @@ -38,24 +49,31 @@ impl TVGuide { match tag.name.as_str() { EPG_TAG_CHANNEL => { if let Some(epg_id) = tag.get_attribute_value(EPG_ATTRIB_ID) { - if let Some(normalized_epg_id) = tag.normalized_epg_id.as_ref() { - match normalized_epg_channel_ids.entry(normalized_epg_id.to_string()) { - std::collections::hash_map::Entry::Occupied(mut entry) => { - entry.insert(Some(epg_id.to_string())); - epg_channel_ids.insert(epg_id.to_string()); + if !processed_epg_channel_ids.contains(epg_id) { + let id: Cow = Cow::Owned(epg_id.to_string()); + if let Some(normalized_epg_id) = tag.normalized_epg_id.as_ref() { + let key = Cow::Owned(normalized_epg_id.to_string()); + match normalized_epg_channel_ids.entry(key) { + std::collections::hash_map::Entry::Occupied(mut entry) => { + entry.insert(Some(id.clone())); + epg_channel_ids.insert(id.clone()); + } + std::collections::hash_map::Entry::Vacant(_entry) => {} } - std::collections::hash_map::Entry::Vacant(_entry) => {} } - } - if epg_channel_ids.contains(epg_id.as_str()) { - children.push(tag); + if epg_channel_ids.contains(&id) { + children.push(tag); + } } } } EPG_TAG_PROGRAMME => { if let Some(epg_id) = tag.get_attribute_value(EPG_ATTRIB_CHANNEL) { - if epg_channel_ids.contains(epg_id.as_str()) { - children.push(tag); + if !processed_epg_channel_ids.contains(epg_id) { + let borrowed_epg_id = Cow::Borrowed(epg_id.as_str()); + if epg_channel_ids.contains(&borrowed_epg_id) { + children.push(tag); + } } } } @@ -71,6 +89,13 @@ impl TVGuide { if children.is_empty() { return None; } + + children.iter().filter(|tag| tag.name == EPG_TAG_CHANNEL).for_each(|tag| { + if let Some(epg_id) = tag.get_attribute_value(EPG_ATTRIB_ID) { + processed_epg_channel_ids.insert(epg_id.to_string()); + } + }); + Some(Epg { attributes: tv_attributes, children, @@ -79,6 +104,22 @@ impl TVGuide { Err(_) => None } } + + pub fn filter(&self, epg_channel_ids: &mut HashSet>, normalized_epg_channel_ids: &mut HashMap, Option>>) -> Option { + if epg_channel_ids.is_empty() && normalized_epg_channel_ids.is_empty() { + return None; + } + let mut processed_epg_ids: HashSet = HashSet::new(); + let epgs: Vec = self.file_paths.iter() + .map(|path| Self::process_epg_file(epg_channel_ids, normalized_epg_channel_ids, &mut processed_epg_ids, path)) + .flatten() + .collect(); + if epgs.len() == 1 { + epgs.into_iter().next() + } else { + Self::merge(epgs) + } + } } pub fn parse_tvguide(content: R, callback: &mut F) @@ -244,30 +285,26 @@ pub fn flatten_tvguide(tv_guides: &[Epg]) -> Option { #[cfg(test)] mod tests { - use crate::model::xmltv::TVGuide; use crate::processing::parser::xmltv::normalize_channel_name; - use std::collections::{HashMap, HashSet}; - use std::io; - use std::path::PathBuf; - #[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"); + // + // 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 normalize() { diff --git a/src/processing/processor/playlist.rs b/src/processing/processor/playlist.rs index f0ac34ecd..e807fb427 100644 --- a/src/processing/processor/playlist.rs +++ b/src/processing/processor/playlist.rs @@ -6,6 +6,7 @@ use crate::utils::network::epg; use crate::utils::network::m3u; use crate::utils::network::xtream; use core::cmp::Ordering; +use std::borrow::Cow; use std::collections::{HashMap, HashSet}; use std::path::PathBuf; use std::sync::{Arc}; @@ -537,12 +538,12 @@ async fn process_playlist_for_target(client: Arc, // we need to process each input epg. for mut fp in processed_fetched_playlists { // collect all epg_channel ids - let mut epg_channel_ids: HashSet = HashSet::new(); - let mut normalized_epg_channel_ids: HashMap> = HashMap::new(); + let mut epg_channel_ids: HashSet> = HashSet::new(); + let mut normalized_epg_channel_ids: HashMap, Option>> = HashMap::new(); for channel in fp.playlistgroups.iter().flat_map(|g| &g.channels) { match channel.header.epg_channel_id.as_ref() { - None => {normalized_epg_channel_ids.insert(normalize_channel_name(&channel.header.name), None);}, - Some(epg_id) => {epg_channel_ids.insert(epg_id.to_string());}, + None => {normalized_epg_channel_ids.insert(Cow::Owned(normalize_channel_name(&channel.header.name)), None);}, + Some(epg_id) => {epg_channel_ids.insert(Cow::Owned(epg_id.to_string()));}, } } // let epg_channel_ids: HashSet<_> = fp.playlistgroups.iter().flat_map(|g| &g.channels) @@ -561,7 +562,8 @@ async fn process_playlist_for_target(client: Arc, .filter(|c| c.header.epg_channel_id.is_none() || c.header.logo.is_empty() || c.header.logo_small.is_empty()) .for_each(|c| { if c.header.epg_channel_id.as_ref().is_none() { - if let Some(Some(epg_id)) = normalized_epg_channel_ids.get(&normalize_channel_name(&c.header.name)) { + let normalized = normalize_channel_name(&c.header.name); + if let Some(Some(epg_id)) = normalized_epg_channel_ids.get(&Cow::Borrowed(normalized.as_str())) { c.header.epg_channel_id = Some(epg_id.to_string()); } } diff --git a/src/repository/storage.rs b/src/repository/storage.rs index 0c09c10a0..92a65df51 100644 --- a/src/repository/storage.rs +++ b/src/repository/storage.rs @@ -22,6 +22,11 @@ pub fn hash_string(text: &str) -> UUIDType { hash_bytes(text.as_bytes()) } +pub fn short_hash(text: &str) -> String { + let hash = blake3::hash(text.as_bytes()); + hex_encode(&hash.as_bytes()[..8]) +} + #[inline] pub fn hex_encode(bytes: &[u8]) -> String { diff --git a/src/utils/file/config_reader.rs b/src/utils/file/config_reader.rs index 110e1f3c6..165c7d1da 100644 --- a/src/utils/file/config_reader.rs +++ b/src/utils/file/config_reader.rs @@ -271,26 +271,30 @@ pub fn csv_read_inputs_from_reader(batch_input_type: InputType, reader: impl Buf let mut result = vec![]; let mut default_columns = vec![]; default_columns.extend_from_slice(DEFAULT_COLUMNS); + let mut header_defined = false; for line in reader.lines() { let line = line?; if line.is_empty() { continue } if line.starts_with(HEADER_PREFIX) { - default_columns = line[1..].split(CSV_SEPARATOR).map(|s| { - match s { - FIELD_URL => FIELD_URL, - FIELD_MAX_CON => FIELD_MAX_CON, - FIELD_PRIO => FIELD_PRIO, - FIELD_NAME => FIELD_NAME, - FIELD_USERNAME => FIELD_USERNAME, - FIELD_PASSWORD => FIELD_PASSWORD, - _ => { - error!("Field {s} is unsupported for csv input"); - FIELD_UNKNOWN + if !header_defined { + header_defined = true; + default_columns = line[1..].split(CSV_SEPARATOR).map(|s| { + match s { + FIELD_URL => FIELD_URL, + FIELD_MAX_CON => FIELD_MAX_CON, + FIELD_PRIO => FIELD_PRIO, + FIELD_NAME => FIELD_NAME, + FIELD_USERNAME => FIELD_USERNAME, + FIELD_PASSWORD => FIELD_PASSWORD, + _ => { + error!("Field {s} is unsupported for csv input"); + FIELD_UNKNOWN + } } - } - }).collect(); + }).collect(); + } continue; } diff --git a/src/utils/network/epg.rs b/src/utils/network/epg.rs index 65caf2d6a..bb9d33d71 100644 --- a/src/utils/network/epg.rs +++ b/src/utils/network/epg.rs @@ -1,25 +1,45 @@ +use std::path::PathBuf; use std::sync::Arc; use log::debug; use crate::m3u_filter_error::M3uFilterError; use crate::model::config::{Config, ConfigInput}; use crate::model::xmltv::TVGuide; +use crate::repository::storage::{short_hash}; use crate::utils::file::file_utils::prepare_file_path; use crate::utils::network::request; use crate::utils::file::file_utils; + +async fn download_epg_file(url: &str, client: &Arc, input: &ConfigInput, working_dir: &str) -> Result { + debug!("Getting epg file path for url: {url}"); + let file_prefix = short_hash(url); + let persist_file_path = prepare_file_path(input.persist.as_deref(), working_dir, "") + .map(|path| file_utils::add_prefix_to_filename(&path, format!("{file_prefix}_epg_").as_str(), Some("xml"))); + + request::get_input_text_content_as_file(Arc::clone(client), input, working_dir, url, persist_file_path).await +} + pub async fn get_xmltv(client: Arc, _cfg: &Config, input: &ConfigInput, working_dir: &str) -> (Option, Vec) { match &input.epg_url { None => (None, vec![]), - Some(url) => { - debug!("Getting epg file path for url: {url}"); - let persist_file_path = prepare_file_path(input.persist.as_deref(), working_dir, "") - .map(|path| file_utils::add_prefix_to_filename(&path, "epg_", Some("xml"))); - - match request::get_input_text_content_as_file(client, input, working_dir, url, persist_file_path).await { - Ok(file) => { - (Some(TVGuide { file }), vec![]) + Some(urls) => { + let mut errors = vec![]; + let mut file_paths = vec![]; + for url in urls { + match download_epg_file(url, &client, input, working_dir).await { + Ok(file_path) => { + file_paths.push(file_path); + } + Err(err) => { + errors.push(err); + } } - Err(err) => (None, vec![err]) + } + + if file_paths.is_empty() { + (None, errors) + } else { + (Some(TVGuide { file_paths }), errors) } } } diff --git a/src/utils/network/request.rs b/src/utils/network/request.rs index a54652225..0c73a7735 100644 --- a/src/utils/network/request.rs +++ b/src/utils/network/request.rs @@ -20,7 +20,7 @@ use crate::m3u_filter_error::create_m3u_filter_error_result; use crate::m3u_filter_error::{str_to_io_error, M3uFilterError, M3uFilterErrorKind}; use crate::model::config::ConfigInput; use crate::model::stats::format_elapsed_time; -use crate::repository::storage::get_input_storage_path; +use crate::repository::storage::{get_input_storage_path, short_hash}; use crate::repository::xtream_repository::FILE_EPG; use crate::utils::compression::compression_utils::{is_deflate, is_gzip, ENCODING_DEFLATE, ENCODING_GZIP}; use crate::utils::debug_if_enabled; @@ -292,7 +292,7 @@ pub async fn download_text_content_as_file(client: Arc, input: } else { let file_path = persist_filepath.map_or_else(|| match get_input_storage_path(&input.name, working_dir) { Ok(download_path) => { - Ok(download_path.join(FILE_EPG)) + Ok(download_path.join(format!("{}_{FILE_EPG}", short_hash(url_str)))) } Err(err) => Err(err) }, Ok);