From e652205de01e344da419fdd4753c83d647817648 Mon Sep 17 00:00:00 2001 From: euzu Date: Fri, 13 Dec 2024 17:51:47 +0100 Subject: [PATCH] xtream resolve vod wip --- src/api/xmltv_api.rs | 6 +- src/model/mapping.rs | 4 +- src/processing/playlist_processor.rs | 8 +- src/processing/xtream_processor.rs | 39 ++++++---- src/repository/bplustree.rs | 108 +++++++++++++------------- src/repository/indexed_document.rs | 51 ++++++++---- src/repository/kodi_repository.rs | 63 +++++++++------ src/repository/mod.rs | 2 +- src/repository/playlist_repository.rs | 2 +- src/repository/xtream_repository.rs | 11 ++- src/utils/download.rs | 14 ++-- src/utils/file_utils.rs | 4 +- src/utils/request_utils.rs | 2 +- 13 files changed, 178 insertions(+), 136 deletions(-) diff --git a/src/api/xmltv_api.rs b/src/api/xmltv_api.rs index 5b799b6b5..9b8ae8e2e 100644 --- a/src/api/xmltv_api.rs +++ b/src/api/xmltv_api.rs @@ -2,7 +2,7 @@ use std::fs::File; use std::path::{Path, PathBuf}; use actix_web::{HttpRequest, HttpResponse, web, http::header}; -use log::{info}; +use log::{error, info}; use quick_xml::{Reader, Writer}; use flate2::write::GzEncoder; use flate2::Compression; @@ -151,7 +151,7 @@ fn serve_epg_with_timeshift(epg_file: File, offset_minutes: i32) -> HttpResponse elem.push_attribute(attr); } Err(e) => { - println!("Error parsing attribute: {e}"); + error!("Error parsing attribute: {e}"); } } } @@ -165,7 +165,7 @@ fn serve_epg_with_timeshift(epg_file: File, offset_minutes: i32) -> HttpResponse xml_writer.write_event(event).expect("Failed to write event"); } Err(e) => { - println!("Error: {e}"); + error!("Error: {e}"); break; } } diff --git a/src/model/mapping.rs b/src/model/mapping.rs index d19dc7d39..d0420ebd0 100644 --- a/src/model/mapping.rs +++ b/src/model/mapping.rs @@ -300,7 +300,7 @@ impl MappingValueProcessor<'_> { } } - fn apply_tags(&self, value: &String, captures: &HashMap<&str, &str>) -> Option { + fn apply_tags(&self, value: &str, captures: &HashMap<&str, &str>) -> Option { let mut new_value = String::from(value); let tag_captures = self.mapper.t_tagre.as_ref().unwrap().captures_iter(value) .filter(|caps| caps.len() > 1) @@ -510,7 +510,7 @@ impl Mappings { self.mappings.prepare() } - pub fn get_mapping(&self, mapping_id: &String) -> Option { + pub fn get_mapping(&self, mapping_id: &str) -> Option { for mapping in &self.mappings.mapping { if mapping.id.eq(mapping_id) { return Some(mapping.clone()); diff --git a/src/processing/playlist_processor.rs b/src/processing/playlist_processor.rs index ad743c361..299a3ce76 100644 --- a/src/processing/playlist_processor.rs +++ b/src/processing/playlist_processor.rs @@ -533,15 +533,15 @@ fn process_watch(target: &ConfigTarget, cfg: &Config, new_playlist: &Vec, targets: Arc) { let (stats, errors) = process_sources(cfg.clone(), targets.clone()).await; + // log errors + for err in &errors { + error!("{}", err.message); + } let stats_msg = format!("{{\"stats\": {}}}", stats.iter().map(std::string::ToString::to_string).collect::>().join("\n")); // print stats info!("{}", stats_msg); // send stats send_message(&MsgKind::Stats, cfg.messaging.as_ref(), stats_msg.as_str()); - // log errors - for err in &errors { - error!("{}", err.message); - } // send errors if let Some(message) = get_errors_notify_message!(errors, 255) { let error_msg = format!("{{\"errors\": \"{}\"}}", message.as_str()); diff --git a/src/processing/xtream_processor.rs b/src/processing/xtream_processor.rs index f644fa11f..55b492f88 100644 --- a/src/processing/xtream_processor.rs +++ b/src/processing/xtream_processor.rs @@ -4,7 +4,7 @@ use crate::model::playlist::{FetchedPlaylist, PlaylistEntry, PlaylistItem, Playl use crate::processing::playlist_processor::ProcessingPipe; use crate::repository::storage::get_input_storage_path; use crate::repository::xtream_repository::{xtream_get_info_file_paths, xtream_update_input_vod_info_file, xtream_update_input_vod_tmdb_file}; -use crate::repository::IndexedDocumentQuery; +use crate::repository::IndexedDocumentIndex; use crate::utils::download; use serde_json::{Map, Value}; use std::collections::HashSet; @@ -152,11 +152,13 @@ pub async fn playlist_resolve_vod(cfg: &Config, target: &ConfigTarget, errors: & // we cant write to the indexed-document directly because of the write lock and time-consuming operation. // All readers would be waiting for the lock and the app would be unresponsive. // We collect the content into a temp file and write it once we collected everything. - let Some((temp_file_info, temp_file_tmdb)) = create_resolve_vod_info_temp_files(errors) else { return }; + let Some((mut temp_file_info, mut temp_file_tmdb)) = create_resolve_vod_info_temp_files(errors) else { return }; let mut processed_vod_ids = read_processed_vod_info_ids(cfg, errors, fpl).await; let mut info_writer = BufWriter::new(&temp_file_info); let mut tmdb_writer = BufWriter::new(&temp_file_tmdb); + let mut info_updated = false; + let mut tmdb_updated = false; for pli in fpl.playlistgroups.iter().flat_map(|plg| &plg.channels) { let a = pli.header.borrow_mut().get_provider_id().as_ref().map_or(false, |pid| processed_vod_ids.contains(pid)); if !a { @@ -166,32 +168,37 @@ pub async fn playlist_resolve_vod(cfg: &Config, target: &ConfigTarget, errors: & errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve vod, could not write to temporary file {err}"))); return; } + info_updated = true; processed_vod_ids.insert(provider_id); if tmdb_id > 0 { if let Err(err) = write_vod_info_tmdb_to_temp_file(&mut tmdb_writer, provider_id, tmdb_id) { errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve vod tmdb, could not write to temporary file {err}"))); return; } + tmdb_updated = true; } - // TODO create tmdb_id index for kodi export } } } } - if let Err(err) = info_writer.flush() { - errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve vod, could not write to temporary file {err}"))); + if info_updated { + if let Err(err) = info_writer.flush() { + errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve vod, could not write to temporary file {err}"))); + } + drop(info_writer); + if let Err(err) = xtream_update_input_vod_info_file(cfg, fpl.input, &mut temp_file_info).await { + errors.push(err); + } } - if let Err(err) = tmdb_writer.flush() { - errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve vod tmdb, could not write to temporary file {err}"))); + if tmdb_updated { + if let Err(err) = tmdb_writer.flush() { + errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve vod tmdb, could not write to temporary file {err}"))); + } + drop(tmdb_writer); + if let Err(err) = xtream_update_input_vod_tmdb_file(cfg, fpl.input, &mut temp_file_tmdb).await { + errors.push(err); + } } - - if let Err(err) = xtream_update_input_vod_info_file(cfg, fpl.input, &temp_file_info).await { - errors.push(err); - } - if let Err(err) = xtream_update_input_vod_tmdb_file(cfg, fpl.input, &temp_file_tmdb).await { - errors.push(err); - } - } async fn read_processed_vod_info_ids(cfg: &Config, errors: &mut Vec, fpl: &FetchedPlaylist<'_>) -> HashSet { @@ -201,7 +208,7 @@ async fn read_processed_vod_info_ids(cfg: &Config, errors: &mut Vec { match cfg.file_locks.read_lock(&file_path).await { Ok(file_lock) => { - if let Ok(mut info_id_mapping) = IndexedDocumentQuery::::try_new(&idx_path) { + if let Ok(info_id_mapping) = IndexedDocumentIndex::::load(&idx_path) { info_id_mapping.traverse(|keys, _| { for doc_id in keys { processed_vod_ids.insert(*doc_id); } }); diff --git a/src/repository/bplustree.rs b/src/repository/bplustree.rs index 8c17a7ca7..0cdea080d 100644 --- a/src/repository/bplustree.rs +++ b/src/repository/bplustree.rs @@ -455,32 +455,32 @@ where }; } } - -fn traverse_tree(file: &mut R, offset: u64, callback: &mut F) -where - K: Ord + Serialize + for<'de> Deserialize<'de> + Clone, - V: Serialize + for<'de> Deserialize<'de> + Clone, - F: FnMut(&Vec, &Vec), -{ - let current_offset = offset; - let mut buffer = vec![0u8; BLOCK_SIZE]; - - match BPlusTreeNode::::deserialize_from_block(file, &mut buffer, current_offset, false) { - Ok((node, pointers)) => { - if node.is_leaf { - callback(&node.keys, &node.values); - } else if let Some(child_pointers) = pointers { - for &child_offset in &child_pointers { - traverse_tree(file, child_offset, callback); - } - } - // if it's a leaf we return. - } - Err(err) => { - error!("Failed to read tree node at offset {current_offset}: {err}"); - } - } -} +// +// fn traverse_tree(file: &mut R, offset: u64, callback: &mut F) +// where +// K: Ord + Serialize + for<'de> Deserialize<'de> + Clone, +// V: Serialize + for<'de> Deserialize<'de> + Clone, +// F: FnMut(&Vec, &Vec), +// { +// let current_offset = offset; +// let mut buffer = vec![0u8; BLOCK_SIZE]; +// +// match BPlusTreeNode::::deserialize_from_block(file, &mut buffer, current_offset, false) { +// Ok((node, pointers)) => { +// if node.is_leaf { +// callback(&node.keys, &node.values); +// } else if let Some(child_pointers) = pointers { +// for &child_offset in &child_pointers { +// traverse_tree(file, child_offset, callback); +// } +// } +// // if it's a leaf we return. +// } +// Err(err) => { +// error!("Failed to read tree node at offset {current_offset}: {err}"); +// } +// } +// } /// `BPlusTreeQuery` can be used to query the `BPlusTree` on-disk. /// If you intend to do frequent queries then use `BPlusTree` instead which loads the tree into memory. @@ -513,12 +513,12 @@ where query_tree(&mut self.file, key) } - pub fn traverse(&mut self, mut visit: F) - where - F: FnMut(&Vec, &Vec), - { - traverse_tree(&mut self.file, 0, &mut visit); - } + // pub fn traverse(&mut self, mut visit: F) + // where + // F: FnMut(&Vec, &Vec), + // { + // traverse_tree(&mut self.file, 0, &mut visit); + // } } pub struct BPlusTreeUpdate { @@ -678,28 +678,28 @@ mod tests { Ok(()) } - #[test] - fn traverse_test() -> io::Result<()> { - let mut tree = BPlusTree::::new(); - for i in 0u32..=500 { - tree.insert(i, Record { - id: i, - data: format!("Entry {i}"), - }); - } - let filepath = PathBuf::from("/tmp/tree.bin"); - // Serialize the tree to a file - tree.store(&filepath)?; - - let mut tree_query: BPlusTreeQuery = BPlusTreeQuery::try_new(&filepath)?; - - // Traverse the tree - tree_query.traverse(|keys, values| { - // TODO real test - println!("Node: {:?} {:?}", keys, values); - }); - Ok(()) - } + // #[test] + // fn traverse_test() -> io::Result<()> { + // let mut tree = BPlusTree::::new(); + // for i in 0u32..=500 { + // tree.insert(i, Record { + // id: i, + // data: format!("Entry {i}"), + // }); + // } + // let filepath = PathBuf::from("/tmp/tree.bin"); + // // Serialize the tree to a file + // tree.store(&filepath)?; + // + // let mut tree_query: BPlusTreeQuery = BPlusTreeQuery::try_new(&filepath)?; + // + // // Traverse the tree + // tree_query.traverse(|keys, values| { + // // TODO real test + // println!("Node: {:?} {:?}", keys, values); + // }); + // Ok(()) + // } #[test] fn insert_dulplicate_test() -> io::Result<()> { diff --git a/src/repository/indexed_document.rs b/src/repository/indexed_document.rs index 7de001aa1..41aa8d159 100644 --- a/src/repository/indexed_document.rs +++ b/src/repository/indexed_document.rs @@ -4,10 +4,10 @@ use std::io::{BufReader, Error, ErrorKind, Read, Seek, SeekFrom, Write}; use std::marker::PhantomData; use std::path::{Path, PathBuf}; -use log::error; -use serde::{Deserialize, Serialize}; use crate::repository::bplustree::{BPlusTree, BPlusTreeQuery, BPlusTreeUpdate}; use crate::utils::file_utils; +use log::error; +use serde::{Deserialize, Serialize}; const BLOCK_SIZE: usize = 4096; const LEN_SIZE: usize = 4; @@ -41,7 +41,8 @@ impl IndexedDocument { pub(in crate::repository) fn get_offset(index_path: &Path, doc_id: &K) -> Result - where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug + where + K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug, { match BPlusTreeQuery::::try_new(index_path) { Ok(mut tree) => { @@ -70,7 +71,8 @@ impl IndexedDocument { * index file is a bplustree */ pub(in crate::repository) struct IndexedDocumentWriter -where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug +where + K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug, { main_path: PathBuf, index_path: PathBuf, @@ -81,8 +83,9 @@ where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug fragmented: bool, } -impl IndexedDocumentWriter -where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug +impl IndexedDocumentWriter +where + K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug, { fn new_with_mode(main_path: PathBuf, index_path: PathBuf, append: bool) -> Result { let append_mode = append && main_path.exists(); @@ -147,8 +150,16 @@ where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug pub fn store(&mut self) -> std::io::Result<()> { if self.dirty { self.dirty = false; - self.main_file.flush()?; - self.index_tree.store(&self.index_path).map(|_| ()) + match self.index_tree.store(&self.index_path) { + Ok(written_bytes) => { + if written_bytes > 0 { + self.main_file.flush() + } else { + Ok(()) + } + } + Err(err) => { Err(err) } + } } else { Ok(()) } @@ -166,7 +177,8 @@ where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug if size == encoded_bytes.len() { // check if it is equal let mut record_buffer = vec![0; size]; - record_buffer.resize(size, 0); self.main_file.read_exact(&mut record_buffer)?; + record_buffer.resize(size, 0); + self.main_file.read_exact(&mut record_buffer)?; if record_buffer == encoded_bytes { return Ok(()); } @@ -208,8 +220,9 @@ where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug } } -impl Drop for IndexedDocumentWriter -where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug +impl Drop for IndexedDocumentWriter +where + K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug, { fn drop(&mut self) { let _ = self.store(); @@ -234,7 +247,8 @@ pub(in crate::repository) struct IndexedDocumentReader { } impl IndexedDocumentReader -where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug +where + K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug, { pub fn new(main_path: &Path, index_path: &Path) -> Result { if main_path.exists() && index_path.exists() { @@ -269,7 +283,7 @@ where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug pub fn get_path(&self) -> &Path { - &self.main_path + &self.main_path } pub const fn has_error(&self) -> bool { @@ -323,7 +337,8 @@ where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug } impl Iterator for IndexedDocumentReader -where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug +where + K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug, { type Item = T; @@ -351,8 +366,9 @@ pub(in crate::repository) struct IndexedDocumentGarbageCollector { index_tree: BPlusTree, } -impl IndexedDocumentGarbageCollector -where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone +impl IndexedDocumentGarbageCollector +where + K: Ord + Serialize + for<'de> Deserialize<'de> + Clone, { pub fn new(main_path: PathBuf, index_path: PathBuf) -> Result { if main_path.exists() && index_path.exists() { @@ -527,4 +543,5 @@ mod tests { } pub type IndexedDocumentQuery = BPlusTreeQuery; -pub type IndexedDocumentUpdate = BPlusTreeUpdate; \ No newline at end of file +pub type IndexedDocumentUpdate = BPlusTreeUpdate; +pub type IndexedDocumentIndex = BPlusTree; \ No newline at end of file diff --git a/src/repository/kodi_repository.rs b/src/repository/kodi_repository.rs index c65e59477..dbb83b573 100644 --- a/src/repository/kodi_repository.rs +++ b/src/repository/kodi_repository.rs @@ -1,7 +1,7 @@ use crate::create_m3u_filter_error_result; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::{Config, ConfigTarget}; -use crate::model::playlist::{PlaylistGroup}; +use crate::model::playlist::PlaylistGroup; use crate::repository::bplustree::BPlusTree; use crate::repository::storage::get_input_storage_path; use crate::repository::xtream_repository::xtream_get_vod_tmdb_file_path; @@ -27,7 +27,7 @@ fn sanitize_for_filename(text: &str, underscore_whitespace: bool) -> String { .collect::() } -fn kodi_style_rename_year(name: &String, style: &KodiStyle) -> (String, Option) { +fn kodi_style_rename_year(name: &str, style: &KodiStyle) -> (String, Option) { let current_date = chrono::Utc::now(); let cur_year = current_date.year(); match style.year.find(name) { @@ -44,7 +44,7 @@ fn kodi_style_rename_year(name: &String, style: &KodiStyle) -> (String, Option (String, Option) { +fn kodi_style_rename_season(name: &str, style: &KodiStyle) -> (String, Option) { style.season.find(name).map_or_else(|| (String::from(name), Some(String::from("01"))), |m| { let s_season = &name[m.start()..m.end()]; let season = Some(String::from(&s_season[1..])); @@ -53,7 +53,7 @@ fn kodi_style_rename_season(name: &String, style: &KodiStyle) -> (String, Option }) } -fn kodi_style_rename_episode(name: &String, style: &KodiStyle) -> (String, Option) { +fn kodi_style_rename_episode(name: &str, style: &KodiStyle) -> (String, Option) { style.episode.find(name).map_or_else(|| (String::from(name), None), |m| { let s_episode = &name[m.start()..m.end()]; let episode = Some(String::from(&s_episode[1..])); @@ -62,15 +62,26 @@ fn kodi_style_rename_episode(name: &String, style: &KodiStyle) -> (String, Optio }) } -fn kodi_style_rename(name: &String, style: &KodiStyle) -> String { +fn kodi_style_rename(name: &str, style: &KodiStyle) -> (Vec, String) { let (work_name_1, year) = kodi_style_rename_year(name, style); let (work_name_2, season) = kodi_style_rename_season(&work_name_1, style); let (work_name_3, episode) = kodi_style_rename_episode(&work_name_2, style); - if year.is_some() && season.is_some() && episode.is_some() { - let formatted = format!("{} ({}) S{}E{}", work_name_3, year.unwrap(), season.unwrap(), episode.unwrap()); - return String::from(style.whitespace.replace_all(formatted.as_str(), " ").as_ref()); + let mut filename = work_name_3; + let mut filedir = vec![String::from(style.whitespace.replace_all(filename.as_str(), " "))]; + if year.is_some() || season.is_some() { + if year.is_some() { + filename = format!("{filename} ({})", year.unwrap()); + filedir = vec![String::from(style.whitespace.replace_all(filename.as_str(), " "))]; + } + if season.is_some() && episode.is_some() { + let season_value = season.unwrap(); + filedir.push(format!("Season {season_value}")); + filename = format!("{filename} S{season_value}E{}", episode.unwrap()); + } + return (filedir, String::from(style.whitespace.replace_all(filename.as_str(), " "))); } - String::from(name) + let new_name = String::from(style.whitespace.replace_all(name, " ")); + (vec![new_name.clone()], new_name) } static KODY_STYLE: LazyLock = LazyLock::new(|| KodiStyle { @@ -115,7 +126,7 @@ async fn get_tmdb_id(cfg: &Config, provider_id: Option, input_id: u16, } } -pub async fn kodi_write_strm_playlist(target: &ConfigTarget, cfg: &Config, new_playlist: &[PlaylistGroup], filename: Option<&String>) -> Result<(), M3uFilterError> { +pub async fn kodi_write_strm_playlist(target: &ConfigTarget, cfg: &Config, new_playlist: &[PlaylistGroup], filename: Option<&str>) -> Result<(), M3uFilterError> { if !new_playlist.is_empty() { if filename.is_none() { return Err(M3uFilterError::new(M3uFilterErrorKind::Notify, "write strm playlist failed: ".to_string())); @@ -137,24 +148,28 @@ pub async fn kodi_write_strm_playlist(target: &ConfigTarget, cfg: &Config, new_p for pg in new_playlist { for pli in &pg.channels { - let provider_id = pli.header.borrow_mut().get_provider_id(); - let input_id = pli.header.borrow().input_id; - let tmdb_id = get_tmdb_id(cfg, provider_id, input_id, &mut input_tmdb_indexes).await; - let additional_info = match tmdb_id { - None => {String::new()} - Some(id) => {format!(" {{tmdb={id}}}")} - }; + let header = &mut pli.header.borrow_mut(); + let mut dir_path = path.join(sanitize_for_filename(&header.group, underscore_whitespace)); + let mut kodi_file_name = sanitize_for_filename(&header.title, underscore_whitespace); + let mut additional_info = String::new(); + if kodi_style { + let provider_id = header.get_provider_id(); + let input_id = header.input_id; + let (kodi_file_dir_name, file_name) = kodi_style_rename(&kodi_file_name, &KODY_STYLE); + kodi_file_name = file_name; + kodi_file_dir_name.iter().for_each(|p| dir_path = dir_path.join(p)); - let header = &pli.header.borrow(); - let dir_path = path.join(sanitize_for_filename(&header.group, underscore_whitespace)); + let tmdb_id = get_tmdb_id(cfg, provider_id, input_id, &mut input_tmdb_indexes).await; + additional_info = match tmdb_id { + None => { String::new() } + Some(id) => { format!(" {{tmdb={id}}}") } + }; + + } if let Err(e) = std::fs::create_dir_all(&dir_path) { - error!("cant create directory: {:?}", &path); + error!("cant create directory: {:?}", &dir_path); return create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "failed to write strm playlist: {}", e); }; - let mut kodi_file_name = sanitize_for_filename(&header.title, underscore_whitespace); - if kodi_style { - kodi_file_name = kodi_style_rename(&kodi_file_name, &KODY_STYLE); - } let file_path = dir_path.join(format!("{kodi_file_name}{additional_info}.strm")); match File::create(&file_path) { Ok(mut strm_file) => { diff --git a/src/repository/mod.rs b/src/repository/mod.rs index 27df2e822..e4caf9039 100644 --- a/src/repository/mod.rs +++ b/src/repository/mod.rs @@ -2,7 +2,7 @@ pub mod storage; pub mod target_id_mapping; pub mod bplustree; mod indexed_document; -pub use indexed_document::IndexedDocumentQuery; +pub use indexed_document::{IndexedDocumentIndex}; pub mod playlist_repository; pub mod m3u_repository; pub mod xtream_repository; diff --git a/src/repository/playlist_repository.rs b/src/repository/playlist_repository.rs index 61b42944a..ccf6b7679 100644 --- a/src/repository/playlist_repository.rs +++ b/src/repository/playlist_repository.rs @@ -48,7 +48,7 @@ pub async fn persist_playlist(playlist: &mut [PlaylistGroup], epg: Option<&Epg>, let result = match output.target { TargetType::M3u => m3u_write_playlist(target, cfg, &target_path, playlist).await, TargetType::Xtream => xtream_write_playlist(target, cfg, playlist).await, - TargetType::Strm => kodi_write_strm_playlist(target, cfg, playlist, output.filename.as_ref()).await, + TargetType::Strm => kodi_write_strm_playlist(target, cfg, playlist, output.filename.as_ref().map(|x| x.as_str())).await, }; if let Err(err) = result { diff --git a/src/repository/xtream_repository.rs b/src/repository/xtream_repository.rs index beb3330fd..ce198b5c6 100644 --- a/src/repository/xtream_repository.rs +++ b/src/repository/xtream_repository.rs @@ -1,3 +1,4 @@ +use std::io::{Seek, SeekFrom}; use std::collections::HashMap; use std::fs::File; use std::io::{BufReader, Error, ErrorKind, Read}; @@ -104,7 +105,7 @@ pub fn xtream_get_vod_tmdb_file_path(storage_path: &Path) -> PathBuf { async fn write_playlists_to_file( cfg: &Config, storage_path: &Path, - collections: Vec<(XtreamCluster, &mut [PlaylistItem])>, + collections: Vec<(XtreamCluster, &mut [&PlaylistItem])>, ) -> Result<(), M3uFilterError> { for (cluster, playlist) in collections { let (xtream_path, idx_path) = xtream_get_file_paths(storage_path, cluster); @@ -247,7 +248,7 @@ pub async fn xtream_write_playlist( TAG_PARENT_ID: 0 })); - for pli in plg.channels.drain(..) { + for pli in &plg.channels { let mut header = pli.header.borrow_mut(); let col = match header.item_type { PlaylistItemType::Series => { @@ -874,7 +875,7 @@ pub async fn xtream_write_series_info( pub async fn xtream_update_input_vod_info_file( cfg: &Config, input: &ConfigInput, - temp_file: &File, + temp_file: &mut File, ) -> Result<(), M3uFilterError> { match get_input_storage_path(input, &cfg.working_dir) .map(|storage_path| xtream_get_info_file_paths(&storage_path, XtreamCluster::Video)) @@ -882,6 +883,7 @@ pub async fn xtream_write_series_info( Ok(Some((info_path, idx_path))) => { match cfg.file_locks.write_lock(&info_path).await { Ok(_file_lock) => { + temp_file.seek(SeekFrom::Start(0)).map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not read vod info {err}")))?; let mut reader = BufReader::new(temp_file); match IndexedDocumentWriter::::new_append(info_path, idx_path) { Ok(mut writer) => { @@ -917,7 +919,7 @@ pub async fn xtream_write_series_info( pub async fn xtream_update_input_vod_tmdb_file( cfg: &Config, input: &ConfigInput, - temp_file: &File, + temp_file: &mut File, ) -> Result<(), M3uFilterError> { match get_input_storage_path(input, &cfg.working_dir) .map(|storage_path| xtream_get_vod_tmdb_file_path(&storage_path)) @@ -925,6 +927,7 @@ pub async fn xtream_write_series_info( Ok(tmdb_path) => { match cfg.file_locks.write_lock(&tmdb_path).await { Ok(_file_lock) => { + temp_file.seek(SeekFrom::Start(0)).map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not read vod tmdb info {err}")))?; let mut reader = BufReader::new(temp_file); let mut provider_id_bytes = [0u8; 4]; let mut tmdb_id_bytes = [0u8; 4]; diff --git a/src/utils/download.rs b/src/utils/download.rs index 88e06a04b..0e494b982 100644 --- a/src/utils/download.rs +++ b/src/utils/download.rs @@ -18,9 +18,9 @@ const ACTION_GET_SERIES_INFO: &str = "get_series_info"; const ACTION_GET_VOD_INFO: &str = "get_vod_info"; const ACTION_GET_LIVE_INFO: &str = "get_live_info"; -fn prepare_file_path(persist: Option<&String>, working_dir: &str, action: &str) -> Option { +fn prepare_file_path(persist: Option<&str>, working_dir: &str, action: &str) -> Option { let persist_file: Option = - persist.map(|persist_path| file_utils::prepare_persist_path(persist_path.as_str(), action)); + persist.map(|persist_path| file_utils::prepare_persist_path(persist_path, action)); if persist_file.is_some() { let file_path = file_utils::get_file_path(working_dir, persist_file); debug!("persist to file: {file_path:?}"); @@ -30,9 +30,9 @@ fn prepare_file_path(persist: Option<&String>, working_dir: &str, action: &str) } } -pub async fn get_m3u_playlist(cfg: &Config, input: &ConfigInput, working_dir: &String) -> (Vec, Vec) { +pub async fn get_m3u_playlist(cfg: &Config, input: &ConfigInput, working_dir: &str) -> (Vec, Vec) { let url = input.url.clone(); - let persist_file_path = prepare_file_path(input.persist.as_ref(), working_dir, ""); + let persist_file_path = prepare_file_path(input.persist.as_ref().map(|x| x.as_str()), working_dir, ""); match request_utils::get_input_text_content(input, working_dir, &url, persist_file_path).await { Ok(text) => { (m3u_parser::parse_m3u(cfg, input, text.lines()), vec![]) @@ -188,8 +188,8 @@ pub async fn get_xtream_playlist(input: &ConfigInput, working_dir: &str) -> (Vec if !skip_cluster.contains(xtream_cluster) { let category_url = format!("{base_url}&action={category}"); let stream_url = format!("{base_url}&action={stream}"); - let category_file_path = prepare_file_path(input.persist.as_ref(), working_dir, format!("{category}_").as_str()); - let stream_file_path = prepare_file_path(input.persist.as_ref(), working_dir, format!("{stream}_").as_str()); + let category_file_path = prepare_file_path(input.persist.as_ref().map(|x| x.as_str()), working_dir, format!("{category}_").as_str()); + let stream_file_path = prepare_file_path(input.persist.as_ref().map(|x| x.as_str()), working_dir, format!("{stream}_").as_str()); match futures::join!( request_utils::get_input_json_content(input, category_url.as_str(), category_file_path), @@ -228,7 +228,7 @@ pub async fn get_xmltv(_cfg: &Config, input: &ConfigInput, working_dir: &str) -> None => (None, vec![]), Some(url) => { debug!("Getting epg file path for url: {}", url); - let persist_file_path = prepare_file_path(input.persist.as_ref(), working_dir, "") + let persist_file_path = prepare_file_path(input.persist.as_ref().map(|x| x.as_str()), working_dir, "") .map(|path| file_utils::add_prefix_to_filename(&path, "epg_", Some("xml"))); match request_utils::get_input_text_content_as_file(input, working_dir, url, persist_file_path).await { diff --git a/src/utils/file_utils.rs b/src/utils/file_utils.rs index 8cb61ca75..ce5c79ce3 100644 --- a/src/utils/file_utils.rs +++ b/src/utils/file_utils.rs @@ -85,7 +85,7 @@ pub fn get_default_api_proxy_config_path(config_path: &str) -> String { get_default_file_path(config_path, API_PROXY_FILE) } -pub fn get_working_path(wd: &String) -> String { +pub fn get_working_path(wd: &str) -> String { let current_dir = std::env::current_dir().unwrap(); if wd.is_empty() { String::from(current_dir.to_str().unwrap_or(".")) @@ -111,7 +111,7 @@ pub fn open_file(file_name: &Path) -> Result { File::open(file_name) } -pub fn persist_file(persist_file: Option, text: &String) { +pub fn persist_file(persist_file: Option, text: &str) { if let Some(path_buf) = persist_file { let filename = &path_buf.to_str().unwrap_or("?"); match File::create(&path_buf) { diff --git a/src/utils/request_utils.rs b/src/utils/request_utils.rs index 251c89033..2ad1015f6 100644 --- a/src/utils/request_utils.rs +++ b/src/utils/request_utils.rs @@ -73,7 +73,7 @@ pub async fn get_input_text_content_as_file(input: &ConfigInput, working_dir: &s } -pub async fn get_input_text_content(input: &ConfigInput, working_dir: &String, url_str: &str, persist_filepath: Option) -> Result { +pub async fn get_input_text_content(input: &ConfigInput, working_dir: &str, url_str: &str, persist_filepath: Option) -> Result { debug_if_enabled!("getting input text content working_dir: {}, url: {}", working_dir, mask_sensitive_info(url_str)); if url_str.parse::().is_ok() {