diff --git a/CHANGELOG.md b/CHANGELOG.md index c315333bf..8dbfe1198 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,7 +4,7 @@ - no error log for xtream api when content is skipped with options `xtream_skip_[live|vod|series]` - added live channel connection sharing in reverse proxy mode. To activate set `share_live_streams` in target options. - The `kodi` format for movies can contain the `tmdb-id` (_optional_). -To add the `tmdb-id` you can set now `xtream_resolve_movies` and `xtream_resolve_movies_delay` to target options. +To add the `tmdb-id` you can set now `xtream_resolve_video` and `xtream_resolve_video_delay` to target options. # 2.0.10 (2024-12-03) - added Target Output Option `m3u_include_type_in_url`, default false. This adds `live`, `movie`, `series` to the url of the stream in reverse proxy mode. diff --git a/README.md b/README.md index 31270691b..1938f3bc6 100644 --- a/README.md +++ b/README.md @@ -397,15 +397,15 @@ each series info entry needs to be fetched one by one and the provider can ban y The `kodi` format for movies can contain the `tmdb-id` (_optional_). Because xtream api delivers the data only on request, we need to fetch this info for each movie entry. But be aware the provider can ban you if you are doing request too frequently. -- `xtream_resolve_movies` if is set to `true` and you have xtream input, the movies info are fetched and stored. +- `xtream_resolve_video` if is set to `true` and you have xtream input, the movies info are fetched and stored. This can cause a lot of requests to the provider. Be cautious when using this option. -- `xtream_resolve_movies_delay` to avoid a provider ban you can set the seconds between vod_info_request's. Default is 2 seconds. +- `xtream_resolve_video_delay` to avoid a provider ban you can set the seconds between vod_info_request's. Default is 2 seconds. But be aware that the more series entries there are, the longer the process takes. Unlike `series info` `movie info` is only fetched once for each movie. If the data is stored locally there will be no update. -There is a difference for `xtream_resolve_movies` and `xtream_resolve_series`. +There is a difference for `xtream_resolve_video` and `xtream_resolve_series`. `xtream_resolve_series` works only when input: `xtream` and output: `m3u`. -`xtream_resolve_movies` works only when input: `xtream`. +`xtream_resolve_video` works only when input: `xtream`. ### 2.2.2.5 `filter` The filter is a string with a filter statement. diff --git a/src/api/api_utils.rs b/src/api/api_utils.rs index 825dd27d6..07fd99eb6 100644 --- a/src/api/api_utils.rs +++ b/src/api/api_utils.rs @@ -200,5 +200,5 @@ async fn shared_stream_response(app_state: &AppState, stream_url: &str) -> Optio } pub fn is_stream_share_enabled(item_type: PlaylistItemType, target: &ConfigTarget) -> bool { - item_type == PlaylistItemType::Live && target.options.as_ref().map_or(false, |opt| opt.share_live_streams) + item_type == PlaylistItemType::Live && target.options.as_ref().is_some_and(|opt| opt.share_live_streams) } \ No newline at end of file diff --git a/src/model/config.rs b/src/model/config.rs index 55952a7a2..fab932e04 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -304,9 +304,9 @@ pub struct ConfigTargetOptions { #[serde(default = "default_as_two_u16")] pub xtream_resolve_series_delay: u16, #[serde(default)] - pub xtream_resolve_movies: bool, + pub xtream_resolve_video: bool, #[serde(default = "default_as_two_u16")] - pub xtream_resolve_movies_delay: u16, + pub xtream_resolve_video_delay: u16, #[serde(default)] pub m3u_include_type_in_url: bool, #[serde(default)] diff --git a/src/model/playlist.rs b/src/model/playlist.rs index 061ab5d75..d9bbdfcc3 100644 --- a/src/model/playlist.rs +++ b/src/model/playlist.rs @@ -16,6 +16,7 @@ use crate::repository::storage::hash_string; // https://siptv.eu/howto/playlist.html pub trait PlaylistEntry { + fn get_uuid(&self) -> Rc<[u8;32]>; fn get_virtual_id(&self) -> u32; fn get_provider_id(&self) -> Option; fn get_category_id(&self) -> Option; @@ -232,6 +233,7 @@ generate_field_accessor_impl_for_playlist_item_header!(id, /*virtual_id,*/ name, #[derive(Debug, Clone, Serialize, Deserialize)] pub struct M3uPlaylistItem { + pub uuid: Rc, pub virtual_id: u32, pub provider_id: Rc, pub name: Rc, @@ -274,6 +276,10 @@ impl M3uPlaylistItem { } impl PlaylistEntry for M3uPlaylistItem { + fn get_uuid(&self) -> Rc<[u8; 32]> { + Rc::clone(&self.uuid) + } + #[inline] fn get_virtual_id(&self) -> u32 { self.virtual_id @@ -305,6 +311,7 @@ impl PlaylistEntry for M3uPlaylistItem { #[derive(Debug, Clone, Serialize, Deserialize)] pub struct XtreamPlaylistItem { + pub uuid: Rc, pub virtual_id: u32, pub provider_id: u32, pub name: Rc, @@ -331,6 +338,10 @@ impl XtreamPlaylistItem { } impl PlaylistEntry for XtreamPlaylistItem { + fn get_uuid(&self) -> Rc<[u8; 32]> { + Rc::clone(&self.uuid) + } + #[inline] fn get_virtual_id(&self) -> u32 { self.virtual_id @@ -360,6 +371,7 @@ impl PlaylistItem { pub fn to_m3u(&self) -> M3uPlaylistItem { let header = self.header.borrow(); M3uPlaylistItem { + uuid: Rc::clone(&header.uuid), virtual_id: header.virtual_id, provider_id: Rc::clone(&header.id), name: Rc::clone(&header.name), @@ -383,6 +395,7 @@ impl PlaylistItem { let header = self.header.borrow(); let provider_id = header.id.parse::().unwrap_or_default(); XtreamPlaylistItem { + uuid: Rc::clone(&header.uuid), virtual_id: header.virtual_id, provider_id, name: Rc::clone(&header.name), @@ -405,6 +418,11 @@ impl PlaylistItem { } impl PlaylistEntry for PlaylistItem { + #[inline] + fn get_uuid(&self) -> Rc<[u8; 32]> { + Rc::clone(self.header.borrow().get_uuid()) + } + #[inline] fn get_virtual_id(&self) -> u32 { self.header.borrow().virtual_id diff --git a/src/processing/xtream_processor.rs b/src/processing/xtream_processor.rs index bd185be43..b35249959 100644 --- a/src/processing/xtream_processor.rs +++ b/src/processing/xtream_processor.rs @@ -6,7 +6,6 @@ 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::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}; @@ -53,7 +52,7 @@ async fn playlist_resolve_movies_process_playlist_item(pli: &PlaylistItem, input 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) { - result = match get_xtream_stream_info_content(&info_url, input).await { + result = match download::get_xtream_stream_info_content(&info_url, input).await { Ok(content) => Some(content), Err(err) => { errors.push(M3uFilterError::new(M3uFilterErrorKind::Info, format!("{err}"))); @@ -67,7 +66,7 @@ async fn playlist_resolve_movies_process_playlist_item(pli: &PlaylistItem, input result } -fn write_content_to_file(writer: &mut BufWriter<&File>, uuid: &[u8;32], content: &str) -> std::io::Result<()> { +fn write_vod_info_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)?; @@ -77,11 +76,10 @@ fn write_content_to_file(writer: &mut BufWriter<&File>, uuid: &[u8;32], content: 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) + (options.xtream_resolve_video && fpl.input.input_type == InputType::Xtream, options.xtream_resolve_video_delay) } else { (false, 0) }; @@ -118,7 +116,7 @@ pub async fn playlist_resolve_movies(cfg: &Config, target: &ConfigTarget, errors 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) { + 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; } diff --git a/src/repository/xtream_repository.rs b/src/repository/xtream_repository.rs index 1fb421787..856e5ba1d 100644 --- a/src/repository/xtream_repository.rs +++ b/src/repository/xtream_repository.rs @@ -8,10 +8,19 @@ use serde_json::{json, Map, Value}; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::{Config, ConfigInput, ConfigTarget}; -use crate::model::playlist::{PlaylistEntry, PlaylistGroup, PlaylistItem, PlaylistItemType, UUIDType, XtreamCluster, XtreamPlaylistItem}; +use crate::model::playlist::{ + PlaylistEntry, PlaylistGroup, PlaylistItem, PlaylistItemType, UUIDType, XtreamCluster, + XtreamPlaylistItem, +}; use crate::model::xtream::XtreamMappingOptions; -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::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}; @@ -32,11 +41,15 @@ const TAG_PARENT_ID: &str = "parent_id"; const TAG_MOVIE_DATA: &str = "movie_data"; const TAG_STREAM_ID: &str = "stream_id"; - macro_rules! cant_write_result { ($path:expr, $err:expr) => { - create_m3u_filter_error!(M3uFilterErrorKind::Notify, "failed to write xtream playlist: {} - {}", $path.to_str().unwrap() ,$err) - } + create_m3u_filter_error!( + M3uFilterErrorKind::Notify, + "failed to write xtream playlist: {} - {}", + $path.to_str().unwrap(), + $err + ) + }; } macro_rules! try_option_ok { @@ -55,17 +68,24 @@ fn get_collection_path(path: &Path, collection: &str) -> PathBuf { fn ensure_xtream_storage_path(cfg: &Config, target_name: &str) -> Result { if let Some(path) = xtream_get_storage_path(cfg, target_name) { if std::fs::create_dir_all(&path).is_err() { - let msg = format!("Failed to save xtream data, can't create directory {}", &path.to_str().unwrap()); + let msg = format!( + "Failed to save xtream data, can't create directory {}", + &path.to_str().unwrap() + ); return Err(M3uFilterError::new(M3uFilterErrorKind::Notify, msg)); } Ok(path) } else { - let msg = format!("Failed to save xtream data, can't create directory for target {target_name}"); + let msg = + format!("Failed to save xtream data, can't create directory for target {target_name}"); Err(M3uFilterError::new(M3uFilterErrorKind::Notify, msg)) } } -pub fn xtream_get_info_file_paths(storage_path: &Path, cluster: XtreamCluster) -> 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}")); @@ -78,23 +98,33 @@ pub fn xtream_get_info_file_paths(storage_path: &Path, cluster: XtreamCluster) - None } -async fn write_playlists_to_file(cfg: &Config, storage_path: &Path, collections: Vec<(XtreamCluster, &mut [PlaylistItem])>) -> Result<(), M3uFilterError> { +async fn write_playlists_to_file( + cfg: &Config, + storage_path: &Path, + collections: Vec<(XtreamCluster, &mut [PlaylistItem])>, +) -> Result<(), M3uFilterError> { for (cluster, playlist) in collections { let (xtream_path, idx_path) = xtream_get_file_paths(storage_path, cluster); { - let _file_lock = cfg.file_locks.write_lock(&xtream_path).await.map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Info, format!("{err}")))?; + let _file_lock = cfg + .file_locks + .write_lock(&xtream_path) + .await + .map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Info, format!("{err}")))?; match IndexedDocumentWriter::new(xtream_path.clone(), idx_path) { Ok(mut writer) => { for item in playlist { let xtream = item.to_xtream(); match writer.write_doc(item.header.borrow().virtual_id, &xtream) { Ok(()) => {} - Err(err) => return Err(cant_write_result!(&xtream_path, err)) + Err(err) => return Err(cant_write_result!(&xtream_path, err)), } } - writer.store().map_err(|err| cant_write_result!(&xtream_path, err))?; + writer + .store() + .map_err(|err| cant_write_result!(&xtream_path, err))?; } - Err(err) => return Err(cant_write_result!(&xtream_path, err)) + Err(err) => return Err(cant_write_result!(&xtream_path, err)), } } } @@ -116,14 +146,17 @@ fn load_old_category_ids(path: &Path) -> (u32, HashMap) { for col_path in [ get_collection_path(path, COL_CAT_LIVE), get_collection_path(path, COL_CAT_VOD), - get_collection_path(path, COL_CAT_SERIES)] { + get_collection_path(path, COL_CAT_SERIES), + ] { if col_path.exists() { if let Ok(file) = File::open(col_path) { let reader = BufReader::new(file); for entry in json_iter_array::>(reader).flatten() { if let Some(item) = entry.as_object() { if let Some(category_id) = get_map_item_as_str(item, TAG_CATEGORY_ID) { - if let Some(category_name) = get_map_item_as_str(item, TAG_CATEGORY_NAME) { + if let Some(category_name) = + get_map_item_as_str(item, TAG_CATEGORY_NAME) + { if let Ok(cat_id) = category_id.parse::() { result.insert(category_name, cat_id); max_id = max_id.max(cat_id); @@ -139,7 +172,8 @@ fn load_old_category_ids(path: &Path) -> (u32, HashMap) { } pub fn xtream_get_storage_path(cfg: &Config, target_name: &str) -> Option { - get_target_storage_path(cfg, target_name).map(|target_path| target_path.join(PathBuf::from(PATH_XTREAM))) + get_target_storage_path(cfg, target_name) + .map(|target_path| target_path.join(PathBuf::from(PATH_XTREAM))) } pub fn xtream_get_epg_file_path(path: &Path) -> PathBuf { @@ -163,7 +197,10 @@ pub fn xtream_get_file_paths_for_series(storage_path: &Path) -> (PathBuf, PathBu async fn xtream_garbage_collect(config: &Config, target_name: &str) -> std::io::Result<()> { // Garbage collect series 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::Series)); + let (info_path, idx_path) = try_option_ok!(xtream_get_info_file_paths( + &storage_path, + XtreamCluster::Series + )); { let _file_lock = config.file_locks.write_lock(&info_path).await?; IndexedDocumentGarbageCollector::::new(info_path, idx_path)?.garbage_collect()?; @@ -171,7 +208,11 @@ async fn xtream_garbage_collect(config: &Config, target_name: &str) -> std::io:: Ok(()) } -pub async fn xtream_write_playlist(target: &ConfigTarget, cfg: &Config, playlist: &mut [PlaylistGroup]) -> Result<(), M3uFilterError> { +pub async fn xtream_write_playlist( + target: &ConfigTarget, + cfg: &Config, + playlist: &mut [PlaylistGroup], +) -> Result<(), M3uFilterError> { let path = ensure_xtream_storage_path(cfg, target.name.as_str())?; let mut errors = Vec::new(); let mut cat_live_col = vec![]; @@ -196,7 +237,8 @@ pub async fn xtream_write_playlist(target: &ConfigTarget, cfg: &Config, playlist XtreamCluster::Live => &mut cat_live_col, 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 @@ -239,19 +281,22 @@ pub async fn xtream_write_playlist(target: &ConfigTarget, cfg: &Config, playlist for (col_path, data) in [ (get_collection_path(&path, COL_CAT_LIVE), &cat_live_col), (get_collection_path(&path, COL_CAT_VOD), &cat_vod_col), - (get_collection_path(&path, COL_CAT_SERIES), &cat_series_col)] { + (get_collection_path(&path, COL_CAT_SERIES), &cat_series_col), + ] { 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)]).await { + match write_playlists_to_file(cfg, &path, + 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 { if err.kind() != ErrorKind::NotFound { @@ -271,36 +316,69 @@ pub async fn xtream_write_playlist(target: &ConfigTarget, cfg: &Config, playlist Ok(()) } -pub fn xtream_get_collection_path(cfg: &Config, target_name: &str, collection_name: &str) -> Result<(Option, Option), Error> { +pub fn xtream_get_collection_path( + cfg: &Config, + target_name: &str, + collection_name: &str, +) -> Result<(Option, Option), Error> { if let Some(path) = xtream_get_storage_path(cfg, target_name) { let col_path = get_collection_path(&path, collection_name); if col_path.exists() { return Ok((Some(col_path), None)); } } - Err(Error::new(ErrorKind::Other, format!("Cant find collection: {target_name}/{collection_name}"))) + Err(Error::new( + ErrorKind::Other, + format!("Cant find collection: {target_name}/{collection_name}"), + )) } -async fn xtream_read_item_for_stream_id(cfg: &Config, stream_id: u32, storage_path: &Path, cluster: XtreamCluster) -> Result { +async fn xtream_read_item_for_stream_id( + cfg: &Config, + stream_id: u32, + storage_path: &Path, + cluster: XtreamCluster, +) -> Result { let (xtream_path, idx_path) = xtream_get_file_paths(storage_path, cluster); { let _file_lock = cfg.file_locks.read_lock(&xtream_path).await?; - IndexedDocumentReader::::read_indexed_item(&xtream_path, &idx_path, &stream_id) + IndexedDocumentReader::::read_indexed_item( + &xtream_path, + &idx_path, + &stream_id, + ) } } -async fn xtream_read_series_item_for_stream_id(cfg: &Config, stream_id: u32, storage_path: &Path) -> Result { +async fn xtream_read_series_item_for_stream_id( + cfg: &Config, + stream_id: u32, + storage_path: &Path, +) -> Result { let (xtream_path, idx_path) = xtream_get_file_paths_for_series(storage_path); { let _file_lock = cfg.file_locks.read_lock(&xtream_path).await?; - IndexedDocumentReader::::read_indexed_item(&xtream_path, &idx_path, &stream_id) + IndexedDocumentReader::::read_indexed_item( + &xtream_path, + &idx_path, + &stream_id, + ) } } macro_rules! try_cluster { ($xtream_cluster:expr, $item_type:expr, $virtual_id:expr) => { - $xtream_cluster.or_else(|| XtreamCluster::try_from($item_type).ok()) - .ok_or_else(|| Error::new(ErrorKind::Other, format!("Could not determine cluster for xtream item with stream-id {}", $virtual_id))) + $xtream_cluster + .or_else(|| XtreamCluster::try_from($item_type).ok()) + .ok_or_else(|| { + Error::new( + ErrorKind::Other, + format!( + "Could not determine cluster for xtream item with stream-id {}", + $virtual_id + ), + ) + }) }; } @@ -310,32 +388,82 @@ pub async fn xtream_get_item_for_stream_id( target: &ConfigTarget, xtream_cluster: Option, ) -> Result { - 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 storage_path = xtream_get_storage_path(config, target.name.as_str()) - .ok_or_else(|| Error::new(ErrorKind::Other, format!("Could not find path for target {} xtream output", &target.name)))?; + 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 storage_path = xtream_get_storage_path(config, target.name.as_str()).ok_or_else(|| { + Error::new( + ErrorKind::Other, + format!( + "Could not find path for target {} xtream output", + &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::new(ErrorKind::Other, format!("Could not get lock for id mapping for target {} err:{err}", target.name)))?; + let _file_lock = config + .file_locks + .read_lock(&target_id_mapping_file) + .await + .map_err(|err| { + Error::new( + ErrorKind::Other, + format!( + "Could not get lock for id mapping for target {} err:{err}", + target.name + ), + ) + })?; - let mut target_id_mapping = IndexedDocumentQuery::::try_new(&target_id_mapping_file) - .map_err(|err| Error::new(ErrorKind::Other, format!("Could not load id mapping for target {} err:{err}", target.name)))?; + let mut target_id_mapping = + IndexedDocumentQuery::::try_new(&target_id_mapping_file) + .map_err(|err| { + Error::new( + ErrorKind::Other, + format!( + "Could not load id mapping for target {} err:{err}", + target.name + ), + ) + })?; - let mapping = target_id_mapping - .query(&virtual_id) - .ok_or_else(|| Error::new(ErrorKind::Other, format!("Could not find mapping for target {} and id {}", target.name, virtual_id)))?; + let mapping = target_id_mapping.query(&virtual_id).ok_or_else(|| { + Error::new( + ErrorKind::Other, + format!( + "Could not find mapping for target {} and id {}", + target.name, virtual_id + ), + ) + })?; match mapping.item_type { - PlaylistItemType::SeriesInfo => xtream_read_series_item_for_stream_id(config, virtual_id, &storage_path).await, + PlaylistItemType::SeriesInfo => { + xtream_read_series_item_for_stream_id(config, virtual_id, &storage_path).await + } PlaylistItemType::SeriesEpisode => { - let mut item = xtream_read_series_item_for_stream_id(config, mapping.parent_virtual_id, &storage_path).await?; + let mut item = xtream_read_series_item_for_stream_id( + config, + mapping.parent_virtual_id, + &storage_path, + ) + .await?; item.provider_id = mapping.provider_id; Ok(item) } PlaylistItemType::Catchup => { let cluster = try_cluster!(xtream_cluster, mapping.item_type, virtual_id)?; - let mut item = xtream_read_item_for_stream_id(config, mapping.parent_virtual_id, &storage_path, cluster).await?; + let mut item = xtream_read_item_for_stream_id( + config, + mapping.parent_virtual_id, + &storage_path, + cluster, + ) + .await?; item.provider_id = mapping.provider_id; Ok(item) } @@ -347,31 +475,51 @@ pub async fn xtream_get_item_for_stream_id( } } - -pub async fn xtream_load_rewrite_playlist(cluster: XtreamCluster, config: &Config, target: &ConfigTarget, category_id: u32) -> Result>, M3uFilterError> { - Ok(Box::new(XtreamPlaylistIterator::new(cluster, config, target, category_id).await?)) +pub async fn xtream_load_rewrite_playlist( + cluster: XtreamCluster, + config: &Config, + target: &ConfigTarget, + category_id: u32, +) -> Result>, M3uFilterError> { + Ok(Box::new( + XtreamPlaylistIterator::new(cluster, config, target, category_id).await?, + )) } -pub async fn xtream_write_series_info(config: &Config, target_name: &str, - series_info_id: u32, - content: &str) -> Result<(), Error> { +pub async fn xtream_write_series_info( + config: &Config, + target_name: &str, + series_info_id: u32, + content: &str, +) -> Result<(), Error> { let target_path = try_option_ok!(get_target_storage_path(config, target_name)); 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::Series)); + let (info_path, idx_path) = try_option_ok!(xtream_get_info_file_paths( + &storage_path, + XtreamCluster::Series + )); { let _file_lock = config.file_locks.write_lock(&info_path).await?; let mut writer = IndexedDocumentWriter::new_append(info_path, idx_path)?; - writer - .write_doc(series_info_id, content) - .map_err(|_| Error::new(ErrorKind::Other, format!("failed to write xtream series info for target {target_name}")))?; + writer.write_doc(series_info_id, content).map_err(|_| { + Error::new( + ErrorKind::Other, + format!("failed to write xtream series info for target {target_name}"), + ) + })?; writer.store()?; } { 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) { + 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); @@ -382,15 +530,26 @@ pub async fn xtream_write_series_info(config: &Config, target_name: &str, Ok(()) } -pub async fn xtream_write_vod_info(config: &Config, target_name: &str, virtual_id: u32, content: &str) -> Result<(), Error> { +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 (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 series info for target {target_name}")))?; + 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()?; } @@ -398,21 +557,44 @@ pub async fn xtream_write_vod_info(config: &Config, target_name: &str, virtual_i } // Reads the series info entry if exists -pub async fn xtream_load_series_info(config: &Config, target_name: &str, series_id: u32) -> Option { +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) + let _file_lock = config + .file_locks + .read_lock(&target_id_mapping_file) + .await .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()?; + 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() { @@ -425,14 +607,27 @@ pub async fn xtream_load_series_info(config: &Config, target_name: &str, 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) { + 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); + error!( + "Failed to read series info for id {series_id} for {target_name}: {}", + err + ); None } }; @@ -441,18 +636,41 @@ pub async fn xtream_load_series_info(config: &Config, target_name: &str, series_ None } -pub async fn xtream_get_vod_info_mapping(config: &Config, target_name: &str, vod_id: u32) -> Option { +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) + let _file_lock = config + .file_locks + .read_lock(&target_id_mapping_file) + .await .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()?; + 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(&vod_id) { // if id_record.is_expired() { @@ -462,25 +680,43 @@ pub async fn xtream_get_vod_info_mapping(config: &Config, target_name: &str, vod 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 { +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()?; + 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) { + 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); + error!( + "Failed to read vod info for id {vod_id} for {target_name}: {}", + err + ); None } }; @@ -505,18 +741,29 @@ where 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))); + 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 { - movie_data.insert(TAG_DIRECT_SOURCE.to_string(), Value::String(pli.get_provider_url().to_string())); + 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(); + xtream_write_vod_info(config, target.name.as_str(), pli.get_virtual_id(), &result) + .await + .ok(); Ok(result) } @@ -533,28 +780,56 @@ where 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 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") + 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 _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()) { + 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())); + 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())); @@ -566,13 +841,40 @@ where } 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(); + 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)) { +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); + } + } + } + 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) => { @@ -593,15 +895,28 @@ pub async fn xtream_update_input_vod_info_file(cfg: &Config, input: &ConfigInput 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}"))), } } - Err(err) => Err(M3uFilterError::new(M3uFilterErrorKind::Info, format!("{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}"))), + 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 +} diff --git a/src/utils/download.rs b/src/utils/download.rs index 34b6b1a7d..55757e060 100644 --- a/src/utils/download.rs +++ b/src/utils/download.rs @@ -4,7 +4,7 @@ use crate::model::playlist::{FetchedPlaylist, PlaylistEntry, PlaylistGroup, Play use crate::model::xmltv::TVGuide; use crate::processing::xtream_parser::parse_xtream_series_info; use crate::processing::{m3u_parser, xtream_parser}; -use crate::repository::xtream_repository; +use crate::repository::{xtream_repository}; use crate::utils::{file_utils, request_utils}; use log::{debug, info}; use std::cmp::Ordering; @@ -12,6 +12,7 @@ use std::collections::HashSet; use std::io::{Error, ErrorKind}; use std::path::PathBuf; use std::rc::Rc; +use crate::repository::xtream_repository::xtream_get_input_vod_info; const ACTION_GET_SERIES_INFO: &str = "get_series_info"; const ACTION_GET_VOD_INFO: &str = "get_vod_info"; @@ -127,6 +128,12 @@ where if let Some(content) = xtream_repository::xtream_load_vod_info(config, target.name.as_str(), pli.get_virtual_id()).await { return Ok(content); } + let resolve_movies = target.options.as_ref().is_some_and(|opt| opt.xtream_resolve_video); + if resolve_movies { + if let Some(content) = xtream_get_input_vod_info(config, input, &pli.get_uuid()).await { + return xtream_repository::write_and_get_xtream_vod_info(config, target, pli, &content).await; + } + } } if let Ok(content) = get_xtream_stream_info_content(info_url, input).await {