diff --git a/CHANGELOG.md b/CHANGELOG.md index 5370bcbf1..5f6ee3056 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,4 +1,7 @@ # Changelog +# 2.0.7 (2024-11-x) +- epg is now first downloaded to disk and not into memory. It is processed than from disk with a sax parser (slower but previously used up to 2GB ram). + # 2.0.6 (2024-11-02) - breaking change virtual_id handling. You need to clear the data directory. - new content storage implementation with BPlusTree indexing. diff --git a/src/filter.rs b/src/filter.rs index ab73fbae2..f951ddfff 100644 --- a/src/filter.rs +++ b/src/filter.rs @@ -5,7 +5,7 @@ use std::collections::HashMap; use std::rc::Rc; use enum_iterator::all; -use log::{debug, error, Level, log_enabled}; +use log::{trace, debug, error, Level, log_enabled}; use pest::iterators::Pair; use pest::Parser; use petgraph::algo::toposort; @@ -252,8 +252,8 @@ fn get_parser_regexp(expr: &Pair, templates: &Vec) -> Res let regexp = re.unwrap(); let captures = regexp.capture_names() .flatten().map(String::from).filter(|x| !x.is_empty()).collect::>(); - if log_enabled!(Level::Debug) { - debug!("Created regex: {} with captures: [{}]", regstr, captures.join(", ")); + if log_enabled!(Level::Trace) { + trace!("Created regex: {} with captures: [{}]", regstr, captures.join(", ")); } return Ok(RegexWithCaptures { restr: regstr, diff --git a/src/model/mapping.rs b/src/model/mapping.rs index b72b68c84..af3ed7a9a 100644 --- a/src/model/mapping.rs +++ b/src/model/mapping.rs @@ -7,7 +7,7 @@ use std::str::FromStr; use std::sync::{Arc, Mutex}; use enum_iterator::Sequence; -use log::{debug, error}; +use log::{trace, debug, error}; use regex::Regex; use crate::{create_m3u_filter_error_result, handle_m3u_filter_error_result, valid_property}; @@ -283,7 +283,7 @@ impl MappingValueProcessor<'_> { if !self.pli.borrow().header.borrow_mut().set_field(key, value) { error!("Cant set unknown field {} to {}", key, value); } - debug!("Property {} set to {}", key, value); + trace!("Property {} set to {}", key, value); } fn apply_attributes(&mut self, captured_names: &HashMap<&str, &str>) { diff --git a/src/model/stats.rs b/src/model/stats.rs index 627b96e1b..12d99a91d 100644 --- a/src/model/stats.rs +++ b/src/model/stats.rs @@ -1,4 +1,5 @@ use std::fmt::Display; + use crate::model::config::InputType; #[derive(Debug, Clone)] @@ -9,7 +10,17 @@ pub(crate) struct PlaylistStats { impl Display for PlaylistStats { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - write!(f, "{}", format_args!("{{\"groups\": {}, \"channels\": {}}}", self.group_count, self.channel_count)) + write!(f, "{}", format_args!("{{\"groups\": {}, \"channels\": {}}}", self.group_count, self.channel_count)) + } +} + +pub(crate) fn format_elapsed_time(seconds: u64) -> String { + if seconds < 60 { + format!("{seconds} secs") + } else { + let minutes = seconds / 60; + let seconds = seconds % 60; + format!("{minutes}:{seconds} mins") } } @@ -20,11 +31,13 @@ pub(crate) struct InputStats { pub error_count: usize, pub raw_stats: PlaylistStats, pub processed_stats: PlaylistStats, + pub secs_took: u64, } impl Display for InputStats { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - let str = format!("{{\"name\": {}, \"type\": {}, \"errors\": {}, \"raw\": {}, \"processed\": {}}}", + let elapsed = format_elapsed_time(self.secs_took); + let str = format!("{{\"name\": {}, \"type\": {}, \"errors\": {}, \"raw\": {}, \"processed\": {}, \"took\": {elapsed}}}", self.name, self.input_type, self.error_count, self.raw_stats, self.processed_stats); write!(f, "{str}") diff --git a/src/processing/playlist_processor.rs b/src/processing/playlist_processor.rs index db2e72dce..074baedee 100644 --- a/src/processing/playlist_processor.rs +++ b/src/processing/playlist_processor.rs @@ -8,7 +8,9 @@ use std::sync::{Arc, Mutex}; use std::thread; use actix_rt::System; -use log::{debug, error, info, Level, log_enabled}; +use log::{trace, debug, error, info, Level, log_enabled}; +use std::time::Instant; + use unidecode::unidecode; use crate::{Config, get_errors_notify_message, model::config}; @@ -40,7 +42,7 @@ fn filter_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Opt for pg in playlist.iter_mut() { let channels = pg.channels.iter() .filter(|&pli| is_valid(pli, target)).cloned().collect::>(); - debug!("Filtered group {} has now {}/{} items", pg.title, channels.len(), pg.channels.len()); + trace!("Filtered group {} has now {}/{} items", pg.title, channels.len(), pg.channels.len()); if !channels.is_empty() { new_playlist.push(PlaylistGroup { id: pg.id, @@ -181,7 +183,7 @@ fn map_channel(channel: PlaylistItem, mapping: &Mapping) -> PlaylistItem { if !mapping.mapper.is_empty() { let header = channel.header.borrow(); let channel_name = if mapping.match_as_ascii { Rc::new(unidecode(&header.name)) } else { header.name.clone() }; - if mapping.match_as_ascii && log_enabled!(Level::Debug) { debug!("Decoded {} for matching to {}", &header.name, &channel_name); }; + if mapping.match_as_ascii && log_enabled!(Level::Trace) { trace!("Decoded {} for matching to {}", &header.name, &channel_name); }; drop(header); let ref_chan = RefCell::new(&channel); let provider = ValueProvider { pli: ref_chan.clone() }; @@ -298,6 +300,7 @@ async fn process_source(cfg: Arc, source_idx: usize, user_targets: Arc

, source_idx: usize, user_targets: Arc

, source_idx: usize, user_targets: Arc

InputStats { +fn create_input_stat(group_count: usize, channel_count: usize, error_count: usize, input_type: InputType, input_name: &str, secs_took: u64) -> InputStats { InputStats { name: input_name.to_string(), input_type, @@ -371,6 +375,7 @@ fn create_input_stat(group_count: usize, channel_count: usize, error_count: usiz group_count: 0, channel_count: 0, }, + secs_took } } diff --git a/src/utils/json_utils.rs b/src/utils/json_utils.rs index ddd3ae489..daa73e6a0 100644 --- a/src/utils/json_utils.rs +++ b/src/utils/json_utils.rs @@ -1,10 +1,11 @@ use std::collections::HashMap; use std::fs::File; -use serde::de::DeserializeOwned; -use serde_json::{self, Deserializer}; use std::io::{self, BufReader, BufWriter, Error, Read, Write}; use std::path::Path; + +use serde::de::DeserializeOwned; use serde::Serialize; +use serde_json::{self, Deserializer, Value}; fn read_skipping_ws(mut reader: impl Read) -> io::Result { loop { @@ -64,25 +65,27 @@ pub(crate) fn json_iter_array( pub(crate) fn json_filter_file(file_path: &Path, filter: &HashMap<&str, &str>) -> Vec { let mut filtered: Vec = Vec::new(); - if file_path.exists() { - if let Ok(file) = File::open(file_path) { - let reader = BufReader::new(file); - for entry in json_iter_array::>(reader).flatten() { - if let Some(item) = entry.as_object() { - for (&key, &value) in filter { - if let Some(field_value) = item.get(key) { - if field_value.is_string() && field_value.eq(value) { - filtered.push(entry.clone()); - } else if let Some(num_val) = field_value.as_i64() { - if let Ok(filter_num_val) = value.parse::() { - if num_val == filter_num_val { - filtered.push(entry.clone()); - } - } - } - } + if !file_path.exists() { + return filtered; // Return early if the file does not exist + } + + let Ok(file) = File::open(file_path) else { return filtered }; + + let reader = BufReader::new(file); + for entry in json_iter_array::>(reader).flatten() { + if let Some(item) = entry.as_object() { + if filter.iter().all(|(&key, &value)| { + if let Some(field_value) = item.get(key) { + match field_value { + Value::String(s) => s == value, + Value::Number(n) => value.parse::().ok() == n.as_i64(), + _ => false, } + } else { + false } + }) { + filtered.push(entry); } } } @@ -91,8 +94,9 @@ pub(crate) fn json_filter_file(file_path: &Path, filter: &HashMap<&str, &str>) - } pub(crate) fn json_write_documents_to_file(file: &Path, value: &T) -> Result<(), Error> - where - T: ?Sized + Serialize { +where + T: ?Sized + Serialize, +{ match File::create(file) { Ok(file) => { let mut writer = BufWriter::new(file); diff --git a/src/utils/request_utils.rs b/src/utils/request_utils.rs index eb68604df..dbe0d4cea 100644 --- a/src/utils/request_utils.rs +++ b/src/utils/request_utils.rs @@ -3,6 +3,7 @@ use std::fs; use std::fs::File; use std::io::{BufWriter, Error, ErrorKind, Read, Write}; use std::path::{Path, PathBuf}; +use std::time::Instant; use flate2::read::{GzDecoder, ZlibDecoder}; use futures::StreamExt; @@ -14,9 +15,10 @@ use url::Url; use crate::create_m3u_filter_error_result; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::ConfigInput; -use crate::repository::storage::{get_input_storage_path}; +use crate::model::stats::format_elapsed_time; +use crate::repository::storage::get_input_storage_path; use crate::repository::xtream_repository::FILE_EPG; -use crate::utils::compression_utils::{ENCODING_GZIP, ENCODING_DEFLATE, is_gzip}; +use crate::utils::compression_utils::{ENCODING_DEFLATE, ENCODING_GZIP, is_deflate, is_gzip}; use crate::utils::file_utils::{get_file_path, persist_file}; pub(crate) fn bytes_to_megabytes(bytes: u64) -> u64 { @@ -179,6 +181,7 @@ fn get_local_file_content(file_path: &PathBuf) -> Result { async fn get_remote_content_as_file(input: &ConfigInput, url: &Url, file_path: &Path) -> Result { + let start_time = Instant::now(); let request = get_client_request(Some(input), url, None); match request.send().await { Ok(response) => { @@ -199,7 +202,8 @@ async fn get_remote_content_as_file(input: &ConfigInput, url: &Url, file_path: & } file.flush()?; - debug!("File downloaded successfully to {file_path:?}"); + let elapsed = start_time.elapsed().as_secs(); + debug!("File downloaded successfully to {file_path:?}, took:{}", format_elapsed_time(elapsed)); Ok(file_path.to_path_buf()) } else { Err(std::io::Error::new(std::io::ErrorKind::Other, format!("Request failed with status {}", response.status()))) @@ -210,6 +214,7 @@ async fn get_remote_content_as_file(input: &ConfigInput, url: &Url, file_path: & } async fn get_remote_content(input: &ConfigInput, url: &Url) -> Result { + let start_time = Instant::now(); let request = get_client_request(Some(input), url, None); match request.send().await { Ok(response) => { @@ -226,9 +231,14 @@ async fn get_remote_content(input: &ConfigInput, url: &Url) -> Result { - if bytes.len() >= 2 && is_gzip(&bytes[0..2]) { - encoding = Some(ENCODING_GZIP.to_string()); + if bytes.len() >= 2 { + if is_gzip(&bytes[0..2]) { + encoding = Some(ENCODING_GZIP.to_string()); + } else if is_deflate(&bytes[0..2]) { + encoding = Some(ENCODING_DEFLATE.to_string()); + } } + let mut decode_buffer = String::new(); if let Some(encoding_type) = encoding { match encoding_type.as_str() { @@ -238,7 +248,7 @@ async fn get_remote_content(input: &ConfigInput, url: &Url) -> Result {} Err(err) => return Err(std::io::Error::new(ErrorKind::Other, format!("failed to decode gzip content {err}"))) }; - }, + } ENCODING_DEFLATE => { let mut decoder = ZlibDecoder::new(&bytes[..]); match decoder.read_to_string(&mut decode_buffer) { @@ -252,10 +262,14 @@ async fn get_remote_content(input: &ConfigInput, url: &Url) -> Result Ok(decoded_content), + Ok(decoded_content) => { + debug!("Request took:{} {url}", format_elapsed_time(start_time.elapsed().as_secs())); + Ok(decoded_content) + }, Err(err) => Err(std::io::Error::new(ErrorKind::Other, format!("failed to plain text content {err}"))) } } else { + debug!("Request took:{}, {url:?}", format_elapsed_time(start_time.elapsed().as_secs())); Ok(decode_buffer) } } @@ -271,7 +285,6 @@ async fn get_remote_content(input: &ConfigInput, url: &Url) -> Result) -> Result { if let Ok(url) = url_str.parse::() { - if url.scheme() == "file" { if let Ok(file_path) = url.to_file_path() { if file_path.exists() {