diff --git a/Cargo.lock b/Cargo.lock index 591dfb2a4..f1a8f6e4c 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1793,6 +1793,7 @@ dependencies = [ "serde", "serde_json", "serde_yaml", + "tempfile", "time", "tokio", "tokio-stream", diff --git a/Cargo.toml b/Cargo.toml index 1c3f1e802..db68fef6c 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -53,3 +53,4 @@ bytes = "1.8.0" async-std = "1.13" tokio-stream = { version = "0.1", features = ["sync"] } tokio = "1.41" +tempfile = "3.14" diff --git a/src/processing/xtream_processor.rs b/src/processing/xtream_processor.rs index 867b2ea05..bd185be43 100644 --- a/src/processing/xtream_processor.rs +++ b/src/processing/xtream_processor.rs @@ -3,11 +3,13 @@ use crate::model::config::{Config, ConfigInput, ConfigTarget, InputType, TargetT use crate::model::playlist::{FetchedPlaylist, PlaylistEntry, PlaylistItem, PlaylistItemType, UUIDType, XtreamCluster}; use crate::processing::playlist_processor::ProcessingPipe; use crate::repository::storage::get_input_storage_path; -use crate::repository::xtream_repository::xtream_get_info_file_paths; +use crate::repository::xtream_repository::{xtream_get_info_file_paths, xtream_update_input_vod_info_file}; use crate::repository::IndexedDocumentQuery; use crate::utils::download; use crate::utils::download::get_xtream_stream_info_content; use std::collections::HashSet; +use std::fs::File; +use std::io::{BufWriter, Error, ErrorKind, Write}; use std::rc::Rc; pub async fn playlist_resolve_series(target: &ConfigTarget, errors: &mut Vec, @@ -65,35 +67,93 @@ async fn playlist_resolve_movies_process_playlist_item(pli: &PlaylistItem, input result } -pub async fn playlist_resolve_movies(cfg: &Config, target: &ConfigTarget, errors: &mut Vec, fpl: &FetchedPlaylist<'_>) { +fn write_content_to_file(writer: &mut BufWriter<&File>, uuid: &[u8;32], content: &str) -> std::io::Result<()> { + let length = u32::try_from(content.len()).map_err(|err| Error::new(ErrorKind::Other, err))?; + if length > 0 { + writer.write_all(uuid)?; + writer.write_all(&length.to_le_bytes())?; + writer.write_all(content.as_bytes())?; + } + Ok(()) +} + + +fn get_resolve_movies_options(target: &ConfigTarget, fpl: &FetchedPlaylist) -> (bool, u16) { let (resolve_movies, resolve_delay) = if let Some(options) = &target.options { (options.xtream_resolve_movies && fpl.input.input_type == InputType::Xtream, options.xtream_resolve_movies_delay) } else { (false, 0) }; + (resolve_movies, resolve_delay) +} + +fn create_temp_file() -> Result { + match tempfile::tempfile() { + Ok(temp_file) => Ok(temp_file), + Err(err) => Err(M3uFilterError::new(M3uFilterErrorKind::Info, format!("Cant resolve movies, could not create temporary file {err}"))) + } +} + +pub async fn playlist_resolve_movies(cfg: &Config, target: &ConfigTarget, errors: &mut Vec, fpl: &FetchedPlaylist<'_>) { + let (resolve_movies, resolve_delay) = get_resolve_movies_options(target, fpl); if !resolve_movies { return; } - match get_input_storage_path(fpl.input, &cfg.working_dir).map(|storage_path| xtream_get_info_file_paths(&storage_path, XtreamCluster::Video)) { - Ok(Some((file_path, idx_path))) => { - match cfg.file_locks.write_lock(&file_path).await { - Ok(_file_lock) => { - match IndexedDocumentWriter::::try_new(&idx_path) { - Ok(mut info_id_mapping) => { - for pli in fpl.playlistgroups.iter().flat_map(|plg| &plg.channels) { - if info_id_mapping.query(pli.header.borrow().get_uuid()).is_none() { - playlist_resolve_movies_process_playlist_item(pli, fpl.input, errors, resolve_delay).await; - } - } - } - Err(err) => errors.push(M3uFilterError::new(M3uFilterErrorKind::Info, format!("Could not load id mapping for input {err}"))), - } + + // 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 temp_file = match create_temp_file() { + Ok(value) => value, + Err(err) => { + errors.push(err); + return; + } + }; + + let mut processed_vod_ids = read_processed_vod_info_ids(cfg, errors, fpl).await; + let mut writer = BufWriter::new(&temp_file); + for pli in fpl.playlistgroups.iter().flat_map(|plg| &plg.channels) { + if !processed_vod_ids.contains(pli.header.borrow().get_uuid().as_ref()) { + if let Some(content) = playlist_resolve_movies_process_playlist_item(pli, fpl.input, errors, resolve_delay).await { + processed_vod_ids.insert(*pli.header.borrow().uuid); + if let Err(err) = write_content_to_file(&mut writer, &pli.header.borrow().uuid, &content) { + errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve movies, could not write to temporary file {err}"))); + return; } - Err(err) => errors.push(M3uFilterError::new(M3uFilterErrorKind::Info, format!("{err}"))), } } - Ok(None) => errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not create storage path for input {}", &fpl.input.name.as_ref().map_or("?", |v| v)))), - Err(err) => errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not create storage path for input {err}"))), } -} \ No newline at end of file + if let Err(err) = writer.flush() { + errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve movies, could not write to temporary file {err}"))); + } + + if let Err(err) = xtream_update_input_vod_info_file(cfg, fpl.input, &temp_file).await { + errors.push(err); + } +} + +async fn read_processed_vod_info_ids(cfg: &Config, errors: &mut Vec, fpl: &FetchedPlaylist<'_>) -> HashSet { + let mut processed_vod_ids = HashSet::new(); + { + match get_input_storage_path(fpl.input, &cfg.working_dir).map(|storage_path| xtream_get_info_file_paths(&storage_path, XtreamCluster::Video)) { + Ok(Some((file_path, idx_path))) => { + match cfg.file_locks.read_lock(&file_path).await { + Ok(file_lock) => { + if let Ok(mut info_id_mapping) = IndexedDocumentQuery::::try_new(&idx_path) { + info_id_mapping.traverse(|keys, _| { + for uuid in keys { processed_vod_ids.insert(*uuid); } + }); + } + drop(file_lock); + } + Err(err) => errors.push(M3uFilterError::new(M3uFilterErrorKind::Info, format!("{err}"))), + } + } + Ok(None) => errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not create storage path for input {}", &fpl.input.name.as_ref().map_or("?", |v| v)))), + Err(err) => errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not create storage path for input {err}"))), + } + } + processed_vod_ids +} diff --git a/src/repository/bplustree.rs b/src/repository/bplustree.rs index b1158df57..439a137a2 100644 --- a/src/repository/bplustree.rs +++ b/src/repository/bplustree.rs @@ -2,8 +2,8 @@ use std::array::TryFromSliceError; use std::fs::{File, OpenOptions}; use std::io::{self, BufReader, BufWriter, Error, ErrorKind, Read, Seek, SeekFrom, Write}; use std::marker::PhantomData; -use std::path::Path; use std::mem::size_of; +use std::path::Path; use flate2::Compression; use log::error; @@ -228,10 +228,10 @@ where if self.is_leaf { let values_encoded = bincode_serialize(&self.values)?; let use_compression = values_encoded.len() + LEN_SIZE + FLAG_SIZE > BLOCK_SIZE; - let compression_byte = if use_compression {1u8.to_le_bytes()} else {0u8.to_le_bytes()}; + let compression_byte = if use_compression { 1u8.to_le_bytes() } else { 0u8.to_le_bytes() }; buffer_slice[write_pos..=write_pos].copy_from_slice(&compression_byte); write_pos += FLAG_SIZE; - let content_bytes = if use_compression { + let content_bytes = if use_compression { let mut encoder = flate2::write::ZlibEncoder::new(Vec::new(), Compression::fast()); encoder.write_all(&values_encoded)?; encoder.finish()? @@ -372,7 +372,7 @@ where } pub fn insert(&mut self, key: K, value: V) { - self.dirty= true; + self.dirty = true; if self.root.keys.is_empty() { self.root.keys.push(key); self.root.values.push(value); @@ -429,7 +429,7 @@ where } } -fn query_tree(file: &mut R, key: &K) -> Option +fn query_tree(file: &mut R, key: &K) -> Option where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone, V: Serialize + for<'de> Deserialize<'de> + Clone, @@ -456,6 +456,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}"); + } + } +} + pub struct BPlusTreeQuery { file: BufReader, _marker_k: PhantomData, @@ -467,8 +493,8 @@ where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone, V: Serialize + for<'de> Deserialize<'de> + Clone, { - pub fn try_new(filepath: &Path) -> io::Result { - let file = is_file_valid(File::open(filepath)?)?; + pub fn try_from_file(file: File) -> io::Result { + let file = is_file_valid(file)?; Ok(Self { file: BufReader::new(file), _marker_k: PhantomData, @@ -476,9 +502,21 @@ where }) } + + pub fn try_new(filepath: &Path) -> io::Result { + Self::try_from_file(File::open(filepath)?) + } + pub fn query(&mut self, key: &K) -> Option { 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 struct BPlusTreeUpdate { @@ -522,7 +560,7 @@ where pub fn update(&mut self, key: &K, value: V) -> io::Result { let mut offset = 0; let mut buffer = vec![0u8; BLOCK_SIZE]; - let mut reader = BufReader::new(&mut self.file); + let mut reader = BufReader::new(&mut self.file); loop { match BPlusTreeNode::::deserialize_from_block(&mut reader, &mut buffer, offset, false) { Ok((mut node, pointers)) => { @@ -638,6 +676,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| { + println!("Node: {:?} {:?}", keys, values); + }); + Ok(()) + } + #[test] fn insert_dulplicate_test() -> io::Result<()> { let mut tree = BPlusTree::::new(); @@ -660,7 +720,6 @@ mod tests { }); }); - Ok(()) } } \ No newline at end of file diff --git a/src/repository/xtream_repository.rs b/src/repository/xtream_repository.rs index 65db4ccd4..1fb421787 100644 --- a/src/repository/xtream_repository.rs +++ b/src/repository/xtream_repository.rs @@ -1,17 +1,17 @@ use std::collections::HashMap; use std::fs::File; -use std::io::{BufReader, Error, ErrorKind}; +use std::io::{BufReader, Error, ErrorKind, Read}; use std::path::{Path, PathBuf}; use log::error; use serde_json::{json, Map, Value}; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; -use crate::model::config::{Config, ConfigTarget}; -use crate::model::playlist::{PlaylistEntry, PlaylistGroup, PlaylistItem, PlaylistItemType, XtreamCluster, XtreamPlaylistItem}; +use crate::model::config::{Config, ConfigInput, ConfigTarget}; +use crate::model::playlist::{PlaylistEntry, PlaylistGroup, PlaylistItem, PlaylistItemType, UUIDType, XtreamCluster, XtreamPlaylistItem}; use crate::model::xtream::XtreamMappingOptions; -use crate::repository::indexed_document::{IndexedDocumentGarbageCollector, IndexedDocumentReader, IndexedDocumentWriter, IndexedDocumentQuery, IndexedDocumentUpdate}; -use crate::repository::storage::{get_target_id_mapping_file, get_target_storage_path, hash_string, FILE_SUFFIX_DB, FILE_SUFFIX_INDEX}; +use crate::repository::indexed_document::{IndexedDocumentGarbageCollector, IndexedDocumentQuery, IndexedDocumentReader, IndexedDocumentUpdate, IndexedDocumentWriter}; +use crate::repository::storage::{get_input_storage_path, get_target_id_mapping_file, get_target_storage_path, hash_string, FILE_SUFFIX_DB, FILE_SUFFIX_INDEX}; use crate::repository::target_id_mapping::{TargetIdMapping, VirtualIdRecord}; use crate::repository::xtream_playlist_iterator::XtreamPlaylistIterator; use crate::utils::json_utils::{json_iter_array, json_write_documents_to_file}; @@ -569,4 +569,39 @@ where xtream_write_series_info(config, target.name.as_str(), virtual_id, &result).await.ok(); Ok(result) +} + +pub async fn xtream_update_input_vod_info_file(cfg: &Config, input: &ConfigInput, temp_file: &File) -> Result<(), M3uFilterError> { + match get_input_storage_path(input, &cfg.working_dir).map(|storage_path| xtream_get_info_file_paths(&storage_path, XtreamCluster::Video)) { + Ok(Some((info_path, idx_path))) => { + match cfg.file_locks.write_lock(&info_path).await { + Ok(_file_lock) => { + let mut reader = BufReader::new(temp_file); + match IndexedDocumentWriter::::new_append(info_path, idx_path) { + Ok(mut writer) => { + let mut uuid_bytes = [0u8; 32]; + let mut length_bytes = [0u8; 4]; + loop { + if reader.read_exact(&mut uuid_bytes).is_err() { + break; // End of file + } + reader.read_exact(&mut length_bytes).map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not read temporary vod info {err}")))?; + let length = u32::from_le_bytes(length_bytes) as usize; + let mut buffer = vec![0u8; length]; + reader.read_exact(&mut buffer).map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not read temporary vod info {err}")))?; + if let Ok(content) = String::from_utf8(buffer) { + let _ = writer.write_doc(uuid_bytes, &content); + } + } + Ok(()) + } + Err(err) => Err(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not create create indexed document writer for vod info {err}"))), + } + } + Err(err) => Err(M3uFilterError::new(M3uFilterErrorKind::Info, format!("{err}"))), + } + } + Ok(None) => Err(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not create storage path for input {}", &input.name.as_ref().map_or("?", |v| v)))), + Err(err) => Err(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not create storage path for input {err}"))), + } } \ No newline at end of file