diff --git a/src/processing/playlist_processor.rs b/src/processing/playlist_processor.rs index 796b5aeca..ad743c361 100644 --- a/src/processing/playlist_processor.rs +++ b/src/processing/playlist_processor.rs @@ -26,7 +26,7 @@ use crate::model::stats::{InputStats, PlaylistStats}; use crate::processing::affix_processor::apply_affixes; use crate::processing::playlist_watch::process_group_watch; use crate::processing::xmltv_parser::flatten_tvguide; -use crate::processing::xtream_processor::{playlist_resolve_series, playlist_resolve_movies}; +use crate::processing::xtream_processor::{playlist_resolve_series, playlist_resolve_vod}; use crate::repository::playlist_repository::persist_playlist; use crate::utils::default_utils::default_as_default; use crate::utils::download; @@ -469,7 +469,7 @@ async fn process_playlist(playlists: &mut [FetchedPlaylist<'_>], for provider_fpl in playlists.iter_mut() { let mut processed_fpl = execute_pipe(target, &pipe, provider_fpl); playlist_resolve_series(target, errors, &pipe, provider_fpl, &mut processed_fpl).await; - playlist_resolve_movies(cfg, target, errors, &processed_fpl).await; + playlist_resolve_vod(cfg, target, errors, &processed_fpl).await; // stats let input_stats = stats.get_mut(&processed_fpl.input.id); if let Some(stat) = input_stats { diff --git a/src/processing/xtream_processor.rs b/src/processing/xtream_processor.rs index b35249959..f644fa11f 100644 --- a/src/processing/xtream_processor.rs +++ b/src/processing/xtream_processor.rs @@ -3,14 +3,20 @@ 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, xtream_update_input_vod_info_file}; +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::utils::download; +use serde_json::{Map, Value}; use std::collections::HashSet; use std::fs::File; use std::io::{BufWriter, Error, ErrorKind, Write}; use std::rc::Rc; +const TAG_VOD_INFO_INFO: &str = "info"; +const TAG_VOD_INFO_MOVIE_DATA: &str = "movie_data"; +const TAG_VOD_INFO_TMDB_ID: &str = "tmdb_id"; +const TAG_VOD_INFO_STREAM_ID: &str = "stream_id"; + pub async fn playlist_resolve_series(target: &ConfigTarget, errors: &mut Vec, pipe: &ProcessingPipe, provider_fpl: &mut FetchedPlaylist<'_>, @@ -48,7 +54,7 @@ pub async fn playlist_resolve_series(target: &ConfigTarget, errors: &mut Vec, resolve_delay: u16) -> Option { +async fn playlist_resolve_vod_process_playlist_item(pli: &PlaylistItem, input: &ConfigInput, errors: &mut Vec, resolve_delay: u16) -> Option { let mut result = None; let provider_id = pli.get_provider_id().unwrap_or(0); if let Some(info_url) = download::get_xtream_player_api_info_url(input, XtreamCluster::Video, provider_id) { @@ -66,82 +72,138 @@ async fn playlist_resolve_movies_process_playlist_item(pli: &PlaylistItem, input result } -fn write_vod_info_content_to_file(writer: &mut BufWriter<&File>, uuid: &[u8;32], content: &str) -> std::io::Result<()> { +fn write_vod_info_content_to_temp_file(writer: &mut BufWriter<&File>, provider_id: u32, 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(&provider_id.to_le_bytes())?; 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) { +fn write_vod_info_tmdb_to_temp_file(writer: &mut BufWriter<&File>, provider_id: u32, tmdb_id: u32) -> std::io::Result<()> { + writer.write_all(&provider_id.to_le_bytes())?; + writer.write_all(&tmdb_id.to_le_bytes())?; + Ok(()) +} + +fn get_resolve_video_options(target: &ConfigTarget, fpl: &FetchedPlaylist) -> (bool, u16) { let (resolve_movies, resolve_delay) = - if let Some(options) = &target.options { - (options.xtream_resolve_video && fpl.input.input_type == InputType::Xtream, options.xtream_resolve_video_delay) - } else { - (false, 0) - }; + target.options.as_ref().map_or((false, 0), |opt| (opt.xtream_resolve_video && fpl.input.input_type == InputType::Xtream, opt.xtream_resolve_video_delay)); (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}"))) +fn get_u32_from_serde_value(value: &Value) -> Option { + match value { + Value::Number(num_val) => num_val.as_u64().and_then(|val| u32::try_from(val).ok()), + Value::String(str_val) => { + match str_val.parse::() { + Ok(sid) => Some(sid), + Err(_) => None + } + } + _ => None, } } -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; - } - - // 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_vod_info_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; +fn extract_provider_id_and_tmdb_id_from_vod_info(content: &str) -> Option<(u32, u32)> { + if let Ok(mut doc) = serde_json::from_str::>(content) { + if let Some(Value::Object(movie_data)) = doc.get_mut(TAG_VOD_INFO_MOVIE_DATA) { + if let Some(stream_id_value) = movie_data.get(TAG_VOD_INFO_STREAM_ID) { + if let Some(stream_id) = get_u32_from_serde_value(stream_id_value) { + if let Some(Value::Object(info)) = doc.get_mut(TAG_VOD_INFO_INFO) { + if let Some(tmdb_id_value) = info.get(TAG_VOD_INFO_TMDB_ID) { + if let Some(tmdb_id) = get_u32_from_serde_value(tmdb_id_value) { + return Some((stream_id, tmdb_id)); + } + } + } + return Some((stream_id, 0)); } } } } - 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); - } + None } -async fn read_processed_vod_info_ids(cfg: &Config, errors: &mut Vec, fpl: &FetchedPlaylist<'_>) -> HashSet { +fn create_resolve_vod_info_temp_files(errors: &mut Vec) -> Option<(File, File)> { + let temp_file_info = match tempfile::tempfile() { + Ok(value) => value, + Err(err) => { + errors.push(M3uFilterError::new(M3uFilterErrorKind::Info, format!("Cant resolve vod, could not create temporary file {err}"))); + return None; + } + }; + let temp_file_tmdb = match tempfile::tempfile() { + Ok(value) => value, + Err(err) => { + errors.push(M3uFilterError::new(M3uFilterErrorKind::Info, format!("Cant resolve vod tmdb, could not create temporary file {err}"))); + return None; + } + }; + Some((temp_file_info, temp_file_tmdb)) +} + +pub async fn playlist_resolve_vod(cfg: &Config, target: &ConfigTarget, errors: &mut Vec, fpl: &FetchedPlaylist<'_>) { + let (resolve_movies, resolve_delay) = get_resolve_video_options(target, fpl); + if !resolve_movies { return; } + + // 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 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); + 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 { + if let Some(content) = playlist_resolve_vod_process_playlist_item(pli, fpl.input, errors, resolve_delay).await { + if let Some((provider_id, tmdb_id)) = extract_provider_id_and_tmdb_id_from_vod_info(&content) { + if let Err(err) = write_vod_info_content_to_temp_file(&mut info_writer, provider_id, &content) { + errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve vod, could not write to temporary file {err}"))); + return; + } + 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; + } + } + // 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 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 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 { 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) { + 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); } + for doc_id in keys { processed_vod_ids.insert(*doc_id); } }); } drop(file_lock); @@ -154,4 +216,4 @@ async fn read_processed_vod_info_ids(cfg: &Config, errors: &mut Vec Result Option<(PathBuf, PathBuf)> { +pub fn xtream_get_info_file_paths(storage_path: &Path, cluster: XtreamCluster) -> Option<(PathBuf, PathBuf)> { if cluster == XtreamCluster::Series { let xtream_path = storage_path.join(format!("{FILE_SERIES_EPISODES}.{FILE_SUFFIX_DB}")); let index_path = storage_path.join(format!("{FILE_SERIES_EPISODES}.{FILE_SUFFIX_INDEX}")); @@ -98,6 +97,10 @@ pub fn xtream_get_info_file_paths( None } +pub fn xtream_get_vod_tmdb_file_path(storage_path: &Path) -> PathBuf { + storage_path.join(format!("{FILE_VOD_INFO_TMDB}.{FILE_SUFFIX_DB}")) +} + async fn write_playlists_to_file( cfg: &Config, storage_path: &Path, @@ -238,7 +241,7 @@ pub async fn xtream_write_playlist( XtreamCluster::Series => &mut cat_series_col, XtreamCluster::Video => &mut cat_vod_col, } - .push(json!({ + .push(json!({ TAG_CATEGORY_ID: format!("{}", &cat_id), TAG_CATEGORY_NAME: plg.title.clone(), TAG_PARENT_ID: 0 @@ -285,17 +288,18 @@ pub async fn xtream_write_playlist( ] { match json_write_documents_to_file(&col_path, data) { Ok(()) => {} - Err(err) => { errors.push(format!("Persisting collection failed: {}: {}", &col_path.to_str().unwrap(), err)); + Err(err) => { + errors.push(format!("Persisting collection failed: {}: {}", &col_path.to_str().unwrap(), err)); } } } match write_playlists_to_file(cfg, &path, - vec![ - (XtreamCluster::Live, &mut live_col), - (XtreamCluster::Video, &mut vod_col), - (XtreamCluster::Series, &mut series_col), - ], + vec![ + (XtreamCluster::Live, &mut live_col), + (XtreamCluster::Video, &mut vod_col), + (XtreamCluster::Series, &mut series_col), + ], ).await { Ok(()) => { if let Err(err) = xtream_garbage_collect(cfg, &target.name).await { @@ -451,7 +455,7 @@ pub async fn xtream_get_item_for_stream_id( mapping.parent_virtual_id, &storage_path, ) - .await?; + .await?; item.provider_id = mapping.provider_id; Ok(item) } @@ -463,7 +467,7 @@ pub async fn xtream_get_item_for_stream_id( &storage_path, cluster, ) - .await?; + .await?; item.provider_id = mapping.provider_id; Ok(item) } @@ -480,7 +484,7 @@ pub async fn xtream_load_rewrite_playlist( config: &Config, target: &ConfigTarget, category_id: u32, -) -> Result>, M3uFilterError> { +) -> Result>, M3uFilterError> { Ok(Box::new( XtreamPlaylistIterator::new(cluster, config, target, category_id).await?, )) @@ -514,58 +518,134 @@ pub async fn xtream_write_series_info( { let target_id_mapping_file = get_target_id_mapping_file(&target_path); let _file_lock = config - .file_locks - .write_lock(&target_id_mapping_file) - .await?; - if let Ok(mut target_id_mapping) = - IndexedDocumentUpdate::::try_new(&target_id_mapping_file) - { - if let Some(record) = target_id_mapping.query(&series_info_id) { - let new_record = record.copy_update_timestamp(); - let _ = target_id_mapping.update(&series_info_id, new_record); - } - }; + .file_locks + .write_lock(&target_id_mapping_file) + .await?; + if let Ok(mut target_id_mapping) = + IndexedDocumentUpdate::::try_new(&target_id_mapping_file) + { + if let Some(record) = target_id_mapping.query(&series_info_id) { + let new_record = record.copy_update_timestamp(); + let _ = target_id_mapping.update(&series_info_id, new_record); + } + }; + } + + Ok(()) } - Ok(()) -} - -pub async fn xtream_write_vod_info( - config: &Config, - target_name: &str, - virtual_id: u32, - content: &str, -) -> Result<(), Error> { - let storage_path = try_option_ok!(xtream_get_storage_path(config, target_name)); - let (info_path, idx_path) = try_option_ok!(xtream_get_info_file_paths( + pub async fn xtream_write_vod_info( + config: &Config, + target_name: &str, + virtual_id: u32, + content: &str, + ) -> Result<(), Error> { + let storage_path = try_option_ok!(xtream_get_storage_path(config, target_name)); + let (info_path, idx_path) = try_option_ok!(xtream_get_info_file_paths( &storage_path, XtreamCluster::Video )); - { - let _file_lock = config.file_locks.write_lock(&info_path).await?; - let mut writer = IndexedDocumentWriter::new_append(info_path, idx_path)?; - writer.write_doc(virtual_id, content).map_err(|_| { - Error::new( - ErrorKind::Other, - format!("failed to write xtream vod info for target {target_name}"), - ) - })?; + { + let _file_lock = config.file_locks.write_lock(&info_path).await?; + let mut writer = IndexedDocumentWriter::new_append(info_path, idx_path)?; + writer.write_doc(virtual_id, content).map_err(|_| { + Error::new( + ErrorKind::Other, + format!("failed to write xtream vod info for target {target_name}"), + ) + })?; - writer.store()?; + writer.store()?; + } + Ok(()) } - Ok(()) -} -// Reads the series info entry if exists -pub async fn xtream_load_series_info( - config: &Config, - target_name: &str, - series_id: u32, -) -> Option { - let target_path = get_target_storage_path(config, target_name)?; - let storage_path = xtream_get_storage_path(config, target_name)?; + // Reads the series info entry if exists + pub async fn xtream_load_series_info( + config: &Config, + target_name: &str, + series_id: u32, + ) -> Option { + let target_path = get_target_storage_path(config, target_name)?; + let storage_path = xtream_get_storage_path(config, target_name)?; - { + { + let target_id_mapping_file = get_target_id_mapping_file(&target_path); + let _file_lock = config + .file_locks + .read_lock(&target_id_mapping_file) + .await + .map_err(|err| { + error!( + "Could not lock id mapping for target {target_name}: {}", + err + ); + Error::new( + ErrorKind::Other, + format!("ID mapping load error for target {target_name}"), + ) + }) + .ok()?; + let mut target_id_mapping = + IndexedDocumentQuery::::try_new(&target_id_mapping_file) + .map_err(|err| { + error!( + "Could not load id mapping for target {target_name}: {}", + err + ); + Error::new( + ErrorKind::Other, + format!("ID mapping load error for target {target_name}"), + ) + }) + .ok()?; + + if let Some(id_record) = target_id_mapping.query(&series_id) { + if id_record.is_expired() { + return None; + } + } + } + + let (info_path, idx_path) = xtream_get_info_file_paths(&storage_path, XtreamCluster::Series)?; + + if info_path.exists() && idx_path.exists() { + { + let _file_lock = config + .file_locks + .read_lock(&info_path) + .await + .map_err(|err| { + error!("Could not lock document {:?}: {}", info_path, err); + Error::new( + ErrorKind::Other, + format!("Document Reader error for target {target_name}"), + ) + }) + .ok()?; + return match IndexedDocumentReader::::read_indexed_item( + &info_path, &idx_path, &series_id, + ) { + Ok(content) => Some(content), + Err(err) => { + error!( + "Failed to read series info for id {series_id} for {target_name}: {}", + err + ); + None + } + }; + } + } + None + } + + pub async fn xtream_get_vod_info_mapping( + config: &Config, + target_name: &str, + vod_id: u32, + ) -> Option { + let target_path = get_target_storage_path(config, target_name)?; let target_id_mapping_file = get_target_id_mapping_file(&target_path); let _file_lock = config .file_locks @@ -573,9 +653,9 @@ pub async fn xtream_load_series_info( .await .map_err(|err| { error!( - "Could not lock id mapping for target {target_name}: {}", - err - ); + "Could not lock id mapping for target {target_name}: {}", + err + ); Error::new( ErrorKind::Other, format!("ID mapping load error for target {target_name}"), @@ -586,9 +666,9 @@ pub async fn xtream_load_series_info( IndexedDocumentQuery::::try_new(&target_id_mapping_file) .map_err(|err| { error!( - "Could not load id mapping for target {target_name}: {}", - err - ); + "Could not load id mapping for target {target_name}: {}", + err + ); Error::new( ErrorKind::Other, format!("ID mapping load error for target {target_name}"), @@ -596,327 +676,276 @@ pub async fn xtream_load_series_info( }) .ok()?; - if let Some(id_record) = target_id_mapping.query(&series_id) { - if id_record.is_expired() { - return None; - } - } + // if let Some(id_record) = target_id_mapping.query(&vod_id) { + // if id_record.is_expired() { + // return None; + // } + // } + target_id_mapping.query(&vod_id) } - let (info_path, idx_path) = xtream_get_info_file_paths(&storage_path, XtreamCluster::Series)?; + // Reads the series info entry if exists + pub async fn xtream_load_vod_info( + config: &Config, + target_name: &str, + vod_id: u32, + ) -> Option { + let storage_path = xtream_get_storage_path(config, target_name)?; - if info_path.exists() && idx_path.exists() { - { - let _file_lock = config - .file_locks - .read_lock(&info_path) - .await - .map_err(|err| { - error!("Could not lock document {:?}: {}", info_path, err); - Error::new( - ErrorKind::Other, - format!("Document Reader error for target {target_name}"), - ) - }) - .ok()?; - return match IndexedDocumentReader::::read_indexed_item( - &info_path, &idx_path, &series_id, - ) { - Ok(content) => Some(content), - Err(err) => { - error!( - "Failed to read series info for id {series_id} for {target_name}: {}", - err - ); - None - } - }; - } - } - None -} + xtream_get_vod_info_mapping(config, target_name, vod_id) + .await + .as_ref()?; -pub async fn xtream_get_vod_info_mapping( - config: &Config, - target_name: &str, - vod_id: u32, -) -> Option { - let target_path = get_target_storage_path(config, target_name)?; - let target_id_mapping_file = get_target_id_mapping_file(&target_path); - let _file_lock = config - .file_locks - .read_lock(&target_id_mapping_file) - .await - .map_err(|err| { - error!( - "Could not lock id mapping for target {target_name}: {}", - err - ); - Error::new( - ErrorKind::Other, - format!("ID mapping load error for target {target_name}"), - ) - }) - .ok()?; - let mut target_id_mapping = - IndexedDocumentQuery::::try_new(&target_id_mapping_file) - .map_err(|err| { - error!( - "Could not load id mapping for target {target_name}: {}", - err - ); - Error::new( - ErrorKind::Other, - format!("ID mapping load error for target {target_name}"), - ) - }) - .ok()?; + let (info_path, idx_path) = xtream_get_info_file_paths(&storage_path, XtreamCluster::Video)?; - // if let Some(id_record) = target_id_mapping.query(&vod_id) { - // if id_record.is_expired() { - // return None; - // } - // } - target_id_mapping.query(&vod_id) -} - -// Reads the series info entry if exists -pub async fn xtream_load_vod_info( - config: &Config, - target_name: &str, - vod_id: u32, -) -> Option { - let storage_path = xtream_get_storage_path(config, target_name)?; - - xtream_get_vod_info_mapping(config, target_name, vod_id) - .await - .as_ref()?; - - let (info_path, idx_path) = xtream_get_info_file_paths(&storage_path, XtreamCluster::Video)?; - - if info_path.exists() && idx_path.exists() { - { - let _file_lock = config - .file_locks - .read_lock(&info_path) - .await - .map_err(|err| { - error!("Could not lock document {:?}: {}", info_path, err); - Error::new( - ErrorKind::Other, - format!("Document Reader error for target {target_name}"), - ) - }) - .ok()?; - return match IndexedDocumentReader::::read_indexed_item( - &info_path, &idx_path, &vod_id, - ) { - Ok(content) => Some(content), - Err(err) => { - error!( + if info_path.exists() && idx_path.exists() { + { + let _file_lock = config + .file_locks + .read_lock(&info_path) + .await + .map_err(|err| { + error!("Could not lock document {:?}: {}", info_path, err); + Error::new( + ErrorKind::Other, + format!("Document Reader error for target {target_name}"), + ) + }) + .ok()?; + return match IndexedDocumentReader::::read_indexed_item( + &info_path, &idx_path, &vod_id, + ) { + Ok(content) => Some(content), + Err(err) => { + error!( "Failed to read vod info for id {vod_id} for {target_name}: {}", err ); - None - } - }; + None + } + }; + } } + None } - None -} -pub async fn write_and_get_xtream_vod_info

( - config: &Config, - target: &ConfigTarget, - pli: &P, - content: &str, -) -> Result -where - P: PlaylistEntry, -{ - let mut doc = serde_json::from_str::>(content) - .map_err(|_| Error::new(ErrorKind::Other, "Failed to parse JSON content"))?; + pub async fn write_and_get_xtream_vod_info

( + config: &Config, + target: &ConfigTarget, + pli: &P, + content: &str, + ) -> Result + where + P: PlaylistEntry, + { + let mut doc = serde_json::from_str::>(content) + .map_err(|_| Error::new(ErrorKind::Other, "Failed to parse JSON content"))?; - // wen need to update the movie data with virtual ids. - if let Some(Value::Object(movie_data)) = doc.get_mut(TAG_MOVIE_DATA) { - let stream_id = pli.get_virtual_id(); - let category_id = pli.get_category_id().unwrap_or(0); - movie_data.insert( - TAG_STREAM_ID.to_string(), - Value::Number(serde_json::value::Number::from(stream_id)), - ); - movie_data.insert( - TAG_CATEGORY_ID.to_string(), - Value::Number(serde_json::value::Number::from(category_id)), - ); - let options = XtreamMappingOptions::from_target_options(target.options.as_ref()); - if options.skip_video_direct_source { - movie_data.insert(TAG_DIRECT_SOURCE.to_string(), Value::String(String::new())); - } else { + // wen need to update the movie data with virtual ids. + if let Some(Value::Object(movie_data)) = doc.get_mut(TAG_MOVIE_DATA) { + let stream_id = pli.get_virtual_id(); + let category_id = pli.get_category_id().unwrap_or(0); movie_data.insert( - TAG_DIRECT_SOURCE.to_string(), - Value::String(pli.get_provider_url().to_string()), + TAG_STREAM_ID.to_string(), + Value::Number(serde_json::value::Number::from(stream_id)), ); + movie_data.insert( + TAG_CATEGORY_ID.to_string(), + Value::Number(serde_json::value::Number::from(category_id)), + ); + let options = XtreamMappingOptions::from_target_options(target.options.as_ref()); + if options.skip_video_direct_source { + movie_data.insert(TAG_DIRECT_SOURCE.to_string(), Value::String(String::new())); + } else { + movie_data.insert( + TAG_DIRECT_SOURCE.to_string(), + Value::String(pli.get_provider_url().to_string()), + ); + } } - } - let result = serde_json::to_string(&doc) - .map_err(|_| Error::new(ErrorKind::Other, "Failed to serialize vod info"))?; - xtream_write_vod_info(config, target.name.as_str(), pli.get_virtual_id(), &result) - .await - .ok(); - - Ok(result) -} - -pub async fn write_and_get_xtream_series_info

( - config: &Config, - target: &ConfigTarget, - pli_series_info: &P, - content: &str, -) -> Result -where - P: PlaylistEntry, -{ - let mut doc = serde_json::from_str::(content) - .map_err(|_| Error::new(ErrorKind::Other, "Failed to parse JSON content"))?; - - let target_path = get_target_storage_path(config, target.name.as_str()).ok_or_else(|| { - Error::new( - ErrorKind::Other, - format!("Could not find path for target {}", target.name), - ) - })?; - - let episodes = doc - .get_mut("episodes") - .and_then(Value::as_object_mut) - .ok_or_else(|| Error::new(ErrorKind::Other, "No episodes found in content"))?; - - let virtual_id = pli_series_info.get_virtual_id(); - { - let target_id_mapping_file = get_target_id_mapping_file(&target_path); - let _file_lock = config - .file_locks - .write_lock(&target_id_mapping_file) + let result = serde_json::to_string(&doc) + .map_err(|_| Error::new(ErrorKind::Other, "Failed to serialize vod info"))?; + xtream_write_vod_info(config, target.name.as_str(), pli.get_virtual_id(), &result) .await - .map_err(|err| { - Error::new( - ErrorKind::Other, - format!( - "Could not load id mapping for target {} err:{err}", - target.name - ), - ) - })?; - let mut target_id_mapping = TargetIdMapping::new(&target_id_mapping_file); - let options = XtreamMappingOptions::from_target_options(target.options.as_ref()); + .ok(); - let provider_url = pli_series_info.get_provider_url(); - for episode_list in episodes.values_mut().filter_map(Value::as_array_mut) { - for episode in episode_list.iter_mut().filter_map(Value::as_object_mut) { - if let Some(provider_id) = episode - .get("id") - .and_then(Value::as_str) - .and_then(|id| id.parse::().ok()) - { - let uuid = hash_string(&format!("{provider_url}/{provider_id}")); - let episode_virtual_id = target_id_mapping.insert_entry( - uuid, - provider_id, - PlaylistItemType::SeriesEpisode, - virtual_id, - ); - episode.insert( - "id".to_string(), - Value::String(episode_virtual_id.to_string()), - ); - } - if options.skip_series_direct_source { - episode.insert(TAG_DIRECT_SOURCE.to_string(), Value::String(String::new())); - } - } - } - - drop(target_id_mapping); + Ok(result) } - let result = serde_json::to_string(&doc) - .map_err(|_| Error::new(ErrorKind::Other, "Failed to serialize updated series info"))?; - xtream_write_series_info(config, target.name.as_str(), virtual_id, &result) - .await - .ok(); - Ok(result) -} - -pub async fn xtream_get_input_vod_info( - cfg: &Config, - input: &ConfigInput, - uuid: &UUIDType, -) -> Option { - if let Ok(Some((info_path, idx_path))) = get_input_storage_path(input, &cfg.working_dir) - .map(|storage_path| xtream_get_info_file_paths(&storage_path, XtreamCluster::Video)) + pub async fn write_and_get_xtream_series_info

( + config: &Config, + target: &ConfigTarget, + pli_series_info: &P, + content: &str, + ) -> Result + where + P: PlaylistEntry, { - if let Ok(_file_lock) = cfg.file_locks.read_lock(&info_path).await { - if let Ok(content) = IndexedDocumentReader::::read_indexed_item( - &info_path, &idx_path, uuid, - ) { - return Some(content); - } - } - } - None -} + let mut doc = serde_json::from_str::(content) + .map_err(|_| Error::new(ErrorKind::Other, "Failed to parse JSON content"))?; -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); - } - } - writer.store().map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not store vod info {err}")))?; - Ok(()) - } - Err(err) => Err(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not create create indexed document writer for vod info {err}"))), + let target_path = get_target_storage_path(config, target.name.as_str()).ok_or_else(|| { + Error::new( + ErrorKind::Other, + format!("Could not find path for target {}", target.name), + ) + })?; + + let episodes = doc + .get_mut("episodes") + .and_then(Value::as_object_mut) + .ok_or_else(|| Error::new(ErrorKind::Other, "No episodes found in content"))?; + + let virtual_id = pli_series_info.get_virtual_id(); + { + let target_id_mapping_file = get_target_id_mapping_file(&target_path); + let _file_lock = config + .file_locks + .write_lock(&target_id_mapping_file) + .await + .map_err(|err| { + Error::new( + ErrorKind::Other, + format!( + "Could not load id mapping for target {} err:{err}", + target.name + ), + ) + })?; + let mut target_id_mapping = TargetIdMapping::new(&target_id_mapping_file); + let options = XtreamMappingOptions::from_target_options(target.options.as_ref()); + + let provider_url = pli_series_info.get_provider_url(); + for episode_list in episodes.values_mut().filter_map(Value::as_array_mut) { + for episode in episode_list.iter_mut().filter_map(Value::as_object_mut) { + if let Some(provider_id) = episode + .get("id") + .and_then(Value::as_str) + .and_then(|id| id.parse::().ok()) + { + let uuid = hash_string(&format!("{provider_url}/{provider_id}")); + let episode_virtual_id = target_id_mapping.insert_entry( + uuid, + provider_id, + PlaylistItemType::SeriesEpisode, + virtual_id, + ); + episode.insert( + "id".to_string(), + Value::String(episode_virtual_id.to_string()), + ); + } + if options.skip_series_direct_source { + episode.insert(TAG_DIRECT_SOURCE.to_string(), Value::String(String::new())); } } - Err(err) => Err(M3uFilterError::new( - M3uFilterErrorKind::Info, - format!("{err}"), - )), + } + + drop(target_id_mapping); + } + let result = serde_json::to_string(&doc) + .map_err(|_| Error::new(ErrorKind::Other, "Failed to serialize updated series info"))?; + xtream_write_series_info(config, target.name.as_str(), virtual_id, &result) + .await + .ok(); + + Ok(result) + } + + pub async fn xtream_get_input_vod_info( + cfg: &Config, + input: &ConfigInput, + uuid: &UUIDType, + ) -> Option { + if let Ok(Some((info_path, idx_path))) = get_input_storage_path(input, &cfg.working_dir) + .map(|storage_path| xtream_get_info_file_paths(&storage_path, XtreamCluster::Video)) + { + if let Ok(_file_lock) = cfg.file_locks.read_lock(&info_path).await { + if let Ok(content) = IndexedDocumentReader::::read_indexed_item( + &info_path, &idx_path, uuid, + ) { + return Some(content); + } } } - 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}"), - )), + None } -} + + 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 provider_id_bytes = [0u8; 4]; + let mut length_bytes = [0u8; 4]; + loop { + if reader.read_exact(&mut provider_id_bytes).is_err() { + break; // End of file + } + let provider_id = u32::from_le_bytes(provider_id_bytes); + 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(provider_id, &content); + } + } + writer.store().map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not store vod info {err}")))?; + 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}"))), + } + } + + pub async fn xtream_update_input_vod_tmdb_file( + cfg: &Config, + input: &ConfigInput, + temp_file: &File, + ) -> Result<(), M3uFilterError> { + match get_input_storage_path(input, &cfg.working_dir) + .map(|storage_path| xtream_get_vod_tmdb_file_path(&storage_path)) + { + Ok(tmdb_path) => { + match cfg.file_locks.write_lock(&tmdb_path).await { + Ok(_file_lock) => { + let mut reader = BufReader::new(temp_file); + let mut provider_id_bytes = [0u8; 4]; + let mut tmdb_id_bytes = [0u8; 4]; + let mut tree_tmdb_index: BPlusTree:: = BPlusTree::load(&tmdb_path).unwrap_or_else(|_| BPlusTree::new()); + loop { + if reader.read_exact(&mut provider_id_bytes).is_err() { + break; // End of file + } + let provider_id = u32::from_le_bytes(provider_id_bytes); + if reader.read_exact(&mut tmdb_id_bytes).is_err() { + break; // End of file + } + let tmdb_id = u32::from_le_bytes(tmdb_id_bytes); + tree_tmdb_index.insert(provider_id, tmdb_id); + } + tree_tmdb_index.store(&tmdb_path).map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not store vod tmdb info {err}")))?; + Ok(()) + } + Err(err) => Err(M3uFilterError::new(M3uFilterErrorKind::Info, format!("{err}"))), + } + } + Err(err) => Err(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not create storage path for input {err}"))), + } + } \ No newline at end of file