From 925f2b85fcd54ee60793afbae719b2b72a3329ad Mon Sep 17 00:00:00 2001 From: euzu Date: Tue, 10 Dec 2024 23:16:04 +0100 Subject: [PATCH] wip resolve vod info --- CHANGELOG.md | 3 + README.md | 13 +- src/api/xtream_api.rs | 87 ++--------- src/model/config.rs | 4 + src/model/playlist.rs | 98 ++++++++++++- src/processing/playlist_processor.rs | 23 +-- src/processing/playlist_watch.rs | 3 +- src/processing/xtream_parser.rs | 14 +- src/processing/xtream_processor.rs | 81 +++++++++-- src/repository/indexed_document.rs | 91 +++++++----- src/repository/kodi_repository.rs | 4 + src/repository/m3u_playlist_iterator.rs | 4 +- src/repository/m3u_repository.rs | 2 +- src/repository/mod.rs | 10 +- src/repository/storage.rs | 9 +- src/repository/target_id_mapping.rs | 10 +- src/repository/xtream_playlist_iterator.rs | 6 +- src/repository/xtream_repository.rs | 159 +++++++++++++++++---- src/utils/download.rs | 83 +++++++++-- 19 files changed, 507 insertions(+), 197 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 969b28bc5..c315333bf 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,7 +1,10 @@ # Changelog # 2.0.11 (2024-12-xx) +- watch files are now moved inside the `target` folder. Move them manually from `watch__.bin` to `/watch_.bin` - 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. # 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 c79541caf..31270691b 100644 --- a/README.md +++ b/README.md @@ -389,12 +389,23 @@ Target options are: - `xtream_skip_series_direct_source` if true the direct_source property from provider for series is ignored Because xtream api delivers only the metadata to series, we need to fetch the series and resolve them. But be aware, -each series info entry needs to be fetched one by one. +each series info entry needs to be fetched one by one and the provider can ban you if you are doing request too frequently. - `xtream_resolve_series` if is set to `true` and you have xtream input and m3u output, the series are fetched and resolved. This can cause a lot of requests to the provider. Be cautious when using this option. - `xtream_resolve_series_delay` to avoid a provider ban you can set the seconds between series_info_request's. Default is 2 seconds. But be aware that the more series entries there are, the longer the process takes. +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. + 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. + 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`. +`xtream_resolve_series` works only when input: `xtream` and output: `m3u`. +`xtream_resolve_movies` works only when input: `xtream`. ### 2.2.2.5 `filter` The filter is a string with a filter statement. diff --git a/src/api/xtream_api.rs b/src/api/xtream_api.rs index 12cf6c7b8..dc7a03643 100644 --- a/src/api/xtream_api.rs +++ b/src/api/xtream_api.rs @@ -22,17 +22,17 @@ use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::api_proxy::{ProxyType, ProxyUserCredentials}; use crate::model::config::TargetType; use crate::model::config::{Config, ConfigInput, ConfigTarget}; -use crate::model::playlist::{PlaylistItemType, XtreamCluster, XtreamPlaylistItem}; -use crate::model::xtream::XtreamMappingOptions; +use crate::model::playlist::{PlaylistItemType, XtreamCluster}; use crate::repository::storage::{get_target_storage_path, hash_string}; use crate::repository::target_id_mapping::TargetIdMapping; use crate::repository::xtream_repository; use crate::utils::request_utils::mask_sensitive_info; -use crate::utils::{json_utils, request_utils}; +use crate::utils::{download, json_utils, request_utils}; const ACTION_GET_SERIES_INFO: &str = "get_series_info"; const ACTION_GET_VOD_INFO: &str = "get_vod_info"; -const ACTION_GET_LIVE_INFO: &str = "get_live_info"; +// const ACTION_GET_LIVE_INFO: &str = "get_live_info"; + const ACTION_GET_EPG: &str = "get_epg"; const ACTION_GET_SHORT_EPG: &str = "get_short_epg"; const ACTION_GET_CATCHUP_TABLE: &str = "get_simple_data_table"; @@ -45,9 +45,7 @@ const ACTION_GET_SERIES: &str = "get_series"; const TAG_ID: &str = "id"; const TAG_CATEGORY_ID: &str = "category_id"; -const TAG_DIRECT_SOURCE: &str = "direct_source"; const TAG_STREAM_ID: &str = "stream_id"; -const TAG_MOVIE_DATA: &str = "movie_data"; const TAG_EPG_LISTINGS: &str = "epg_listings"; macro_rules! try_option_bad_request { @@ -142,28 +140,6 @@ pub fn serve_query(file_path: &Path, filter: &HashMap<&str, &str>) -> HttpRespon HttpResponse::Ok().json(filtered) } -fn get_xtream_player_api_action_url(input: &ConfigInput, action: &str) -> Option { - if let Some(user_info) = input.get_user_info() { - Some(format!("{}/player_api.php?username={}&password={}&action={}", - &user_info.base_url, - &user_info.username, - &user_info.password, - action - )) - } else { - None - } -} - -fn get_xtream_player_api_info_url(input: &ConfigInput, cluster: XtreamCluster, stream_id: u32) -> Option { - let (action, stream_id_field) = match cluster { - XtreamCluster::Live => (ACTION_GET_LIVE_INFO, "live_id"), - XtreamCluster::Video => (ACTION_GET_VOD_INFO, "vod_id"), - XtreamCluster::Series => (ACTION_GET_SERIES_INFO, "series_id"), - }; - get_xtream_player_api_action_url(input, action).map(|action_url| format!("{action_url}&{stream_id_field}={stream_id}")) -} - fn get_xtream_player_api_stream_url(input: &ConfigInput, context: &str, action_path: &str, fallback_url: &str) -> Option { let ctx_path = if context.is_empty() { String::new() } else { format!("{context}/") }; if let Some(user_info) = input.get_user_info() { @@ -280,51 +256,6 @@ async fn xtream_player_api_timeshift_stream( xtream_player_api_stream(&req, &api_query_req, &app_state, XtreamApiStreamRequest::from(XtreamApiStreamContext::Timeshift, username, password, stream_id, &action_path)).await } -fn get_xtream_vod_info(target: &ConfigTarget, pli: &XtreamPlaylistItem, content: &str) -> Result { - if let Ok(mut doc) = serde_json::from_str::>(content) { - if let Some(Value::Object(movie_data)) = doc.get_mut(TAG_MOVIE_DATA) { - let stream_id = pli.virtual_id; - let category_id = pli.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.url.to_string())); - } - if let Ok(result) = serde_json::to_string(&doc) { - return Ok(result); - } - } - } - Err(Error::new(ErrorKind::Other, format!("Failed to get vod info for id {}", pli.virtual_id))) -} - -async fn xtream_get_stream_info_content(info_url: &str, input: &ConfigInput) -> Result { - request_utils::download_text_content(input, info_url, None).await -} - -async fn xtream_get_stream_info(config: &Config, input: &ConfigInput, target: &ConfigTarget, - pli: &XtreamPlaylistItem, info_url: &str, cluster: XtreamCluster) -> Result { - if cluster == XtreamCluster::Series { - if let Some(content) = xtream_repository::xtream_load_series_info(config, target.name.as_str(), pli.virtual_id).await { - return Ok(content); - } - } - - if let Ok(content) = xtream_get_stream_info_content(info_url, input).await { - return match cluster { - XtreamCluster::Live => Ok(content), - XtreamCluster::Video => get_xtream_vod_info(target, pli, &content), - XtreamCluster::Series => xtream_repository::write_and_get_xtream_series_info(config, target, pli, &content).await, - }; - } - - Err(Error::new(ErrorKind::Other, format!("Cant find stream with id: {}/{}/{}", - target.name.replace(' ', "_").as_str(), &cluster, pli.virtual_id))) -} - async fn xtream_get_stream_info_response(app_state: &AppState, user: &ProxyUserCredentials, target: &ConfigTarget, stream_id: &str, cluster: XtreamCluster) -> HttpResponse { @@ -336,11 +267,11 @@ async fn xtream_get_stream_info_response(app_state: &AppState, user: &ProxyUserC if let Ok(pli) = xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, Some(cluster)).await { let input_id = pli.input_id; if let Some(input) = app_state.config.get_input_by_id(input_id) { - if let Some(info_url) = get_xtream_player_api_info_url(input, cluster, pli.provider_id) { + if let Some(info_url) = download::get_xtream_player_api_info_url(input, cluster, pli.provider_id) { // Redirect is only possible for live streams, vod and series info needs to be modified if user.proxy == ProxyType::Redirect && cluster == XtreamCluster::Live { return HttpResponse::Found().insert_header(("Location", info_url)).finish(); - } else if let Ok(content) = xtream_get_stream_info(&app_state.config, input, target, &pli, info_url.as_str(), cluster).await { + } else if let Ok(content) = download::get_xtream_stream_info(&app_state.config, input, target, &pli, info_url.as_str(), cluster).await { return HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body(content); } } @@ -364,7 +295,7 @@ async fn xtream_get_short_epg(app_state: &AppState, user: &ProxyUserCredentials, if let Ok(pli) = xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, None).await { let input_id: u16 = pli.input_id; if let Some(input) = app_state.config.get_input_by_id(input_id) { - if let Some(action_url) = get_xtream_player_api_action_url(input, ACTION_GET_SHORT_EPG) { + if let Some(action_url) = download::get_xtream_player_api_action_url(input, ACTION_GET_SHORT_EPG) { let mut info_url = format!("{action_url}&{TAG_STREAM_ID}={}", pli.provider_id); if !(limit.is_empty() || limit.eq("0")) { info_url = format!("{info_url}&limit={limit}"); @@ -413,8 +344,8 @@ async fn xtream_get_catchup_response(app_state: &AppState, target: &ConfigTarget let virtual_id: u32 = try_result_bad_request!(FromStr::from_str(stream_id)); let pli = try_result_bad_request!(xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, Some(XtreamCluster::Live)).await); let input = try_option_bad_request!(app_state.config.get_input_by_id(pli.input_id)); - let info_url = try_option_bad_request!(get_xtream_player_api_action_url(input, ACTION_GET_CATCHUP_TABLE).map(|action_url| format!("{action_url}&{TAG_STREAM_ID}={}&start={start}&end={end}", pli.provider_id))); - let content = try_result_bad_request!(xtream_get_stream_info_content(info_url.as_str(), input).await); + let info_url = try_option_bad_request!(download::get_xtream_player_api_action_url(input, ACTION_GET_CATCHUP_TABLE).map(|action_url| format!("{action_url}&{TAG_STREAM_ID}={}&start={start}&end={end}", pli.provider_id))); + let content = try_result_bad_request!(download::get_xtream_stream_info_content(info_url.as_str(), input).await); let mut doc: Map = try_result_bad_request!(serde_json::from_str(&content)); let epg_listings = try_option_bad_request!(doc.get_mut(TAG_EPG_LISTINGS).and_then(Value::as_array_mut)); let target_path = try_option_bad_request!(get_target_storage_path(&app_state.config, target.name.as_str())); diff --git a/src/model/config.rs b/src/model/config.rs index 32ef27aca..55952a7a2 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -304,6 +304,10 @@ pub struct ConfigTargetOptions { #[serde(default = "default_as_two_u16")] pub xtream_resolve_series_delay: u16, #[serde(default)] + pub xtream_resolve_movies: bool, + #[serde(default = "default_as_two_u16")] + pub xtream_resolve_movies_delay: u16, + #[serde(default)] pub m3u_include_type_in_url: bool, #[serde(default)] pub m3u_mask_redirect_url: bool, diff --git a/src/model/playlist.rs b/src/model/playlist.rs index e54fbb55e..061ab5d75 100644 --- a/src/model/playlist.rs +++ b/src/model/playlist.rs @@ -15,6 +15,13 @@ use crate::repository::storage::hash_string; // https://de.wikipedia.org/wiki/M3U // https://siptv.eu/howto/playlist.html +pub trait PlaylistEntry { + fn get_virtual_id(&self) -> u32; + fn get_provider_id(&self) -> Option; + fn get_category_id(&self) -> Option; + fn get_provider_url(&self) -> Rc; +} + #[derive(Debug, Clone)] pub struct FetchedPlaylist<'a> { // Contains playlist for one input pub input: &'a ConfigInput, @@ -121,9 +128,11 @@ pub trait FieldAccessor { fn set_field(&mut self, field: &str, value: &str) -> bool; } +pub type UUIDType = [u8;32]; + #[derive(Debug, Clone, Serialize, Deserialize, Default)] pub struct PlaylistItemHeader { - pub uuid: Rc<[u8; 32]>, // calculated + pub uuid: Rc, // calculated pub id: Rc, // provider id pub virtual_id: u32, // virtual id pub name: Rc, @@ -154,7 +163,7 @@ impl PlaylistItemHeader { pub fn gen_uuid(&mut self) { self.uuid = Rc::new(hash_string(&self.url)); } - pub const fn get_uuid(&self) -> &Rc<[u8; 32]> { + pub const fn get_uuid(&self) -> &Rc { &self.uuid } @@ -264,6 +273,36 @@ impl M3uPlaylistItem { } } +impl PlaylistEntry for M3uPlaylistItem { + #[inline] + fn get_virtual_id(&self) -> u32 { + self.virtual_id + } + + fn get_provider_id(&self) -> Option { + match self.provider_id.parse::() { + Ok(id) => Some(id), + Err(_) => match extract_id_from_url(&self.url) { + Some(id) => match id.parse::() { + Ok(newid) => { + Some(newid) + } + Err(_) => None, + }, + None => None, + } + } + } + #[inline] + fn get_category_id(&self) -> Option { + None + } + #[inline] + fn get_provider_url(&self) -> Rc { + Rc::clone(&self.url) + } +} + #[derive(Debug, Clone, Serialize, Deserialize)] pub struct XtreamPlaylistItem { pub virtual_id: u32, @@ -291,6 +330,27 @@ impl XtreamPlaylistItem { } } +impl PlaylistEntry for XtreamPlaylistItem { + #[inline] + fn get_virtual_id(&self) -> u32 { + self.virtual_id + } + #[inline] + fn get_provider_id(&self) -> Option { + Some(self.provider_id) + } + #[inline] + fn get_category_id(&self) -> Option { + None + } + #[inline] + fn get_provider_url(&self) -> Rc { + Rc::clone(&self.url) + } + + +} + #[derive(Debug, Clone, Serialize, Deserialize)] pub struct PlaylistItem { pub header: RefCell, @@ -344,6 +404,39 @@ impl PlaylistItem { } } +impl PlaylistEntry for PlaylistItem { + #[inline] + fn get_virtual_id(&self) -> u32 { + self.header.borrow().virtual_id + } + + fn get_provider_id(&self) -> Option { + let header = self.header.borrow(); + match header.id.parse::() { + Ok(id) => Some(id), + Err(_) => match extract_id_from_url(&header.url) { + Some(id) => match id.parse::() { + Ok(newid) => { + Some(newid) + } + Err(_) => None, + }, + None => None, + } + } + } + + #[inline] + fn get_category_id(&self) -> Option { + None + } + #[inline] + fn get_provider_url(&self) -> Rc { + Rc::clone(&self.header.borrow().url) + } + +} + #[derive(Debug, Clone, Serialize, Deserialize)] pub struct PlaylistGroup { pub id: u32, @@ -354,6 +447,7 @@ pub struct PlaylistGroup { } impl PlaylistGroup { + #[inline] pub fn on_load(&mut self) { self.channels.iter().for_each(|pl| pl.header.borrow_mut().gen_uuid()); } diff --git a/src/processing/playlist_processor.rs b/src/processing/playlist_processor.rs index 9c0200b7e..796b5aeca 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; +use crate::processing::xtream_processor::{playlist_resolve_series, playlist_resolve_movies}; use crate::repository::playlist_repository::persist_playlist; use crate::utils::default_utils::default_as_default; use crate::utils::download; @@ -465,29 +465,30 @@ async fn process_playlist(playlists: &mut [FetchedPlaylist<'_>], let pipe = get_processing_pipe(target); debug_if_enabled!("Processing order is {}", &target.processing_order); - let mut new_fetched_playlists: Vec = vec![]; - for fpl in playlists.iter_mut() { - let mut new_fpl = execute_pipe(target, &pipe, fpl); - playlist_resolve_series(target, errors, &pipe, fpl, &mut new_fpl).await; + let mut processed_fetched_playlists: Vec = vec![]; + 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; // stats - let input_stats = stats.get_mut(&new_fpl.input.id); + let input_stats = stats.get_mut(&processed_fpl.input.id); if let Some(stat) = input_stats { - stat.processed_stats.group_count = new_fpl.playlistgroups.len(); - stat.processed_stats.channel_count = new_fpl.playlistgroups.iter() + stat.processed_stats.group_count = processed_fpl.playlistgroups.len(); + stat.processed_stats.channel_count = processed_fpl.playlistgroups.iter() .map(|group| group.channels.len()) .sum(); } - new_fetched_playlists.push(new_fpl); + processed_fetched_playlists.push(processed_fpl); } - apply_affixes(&mut new_fetched_playlists); + apply_affixes(&mut processed_fetched_playlists); let mut new_playlist = vec![]; let mut new_epg = vec![]; // each fetched playlist can have its own epgl url. // we need to process each input epg. - for mut fp in new_fetched_playlists { + for mut fp in processed_fetched_playlists { // collect all epg_channel ids let epg_channel_ids: HashSet<_> = fp.playlistgroups.iter().flat_map(|g| &g.channels) .filter_map(|c| c.header.borrow().epg_channel_id.clone()).collect(); diff --git a/src/processing/playlist_watch.rs b/src/processing/playlist_watch.rs index 3937c05d5..ced59b548 100644 --- a/src/processing/playlist_watch.rs +++ b/src/processing/playlist_watch.rs @@ -15,8 +15,7 @@ pub fn process_group_watch(cfg: &Config, target_name: &str, pl: &PlaylistGroup) new_tree.insert(title); }); - let file_name = format!("watch_{target_name}_{}", &pl.title); - let watch_filename = sanitize_filename(&format!("{file_name}.bin")); + let watch_filename = format!("{}/{}.bin", sanitize_filename(target_name), sanitize_filename(&pl.title)); match file_utils::get_file_path(&cfg.working_dir, Some(std::path::PathBuf::from(&watch_filename))) { Some(path) => { let save_path = path.as_path(); diff --git a/src/processing/xtream_parser.rs b/src/processing/xtream_parser.rs index 26f797089..d3bb45e5e 100644 --- a/src/processing/xtream_parser.rs +++ b/src/processing/xtream_parser.rs @@ -9,6 +9,7 @@ use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::ConfigInput; use crate::model::playlist::{PlaylistGroup, PlaylistItem, PlaylistItemHeader, PlaylistItemType, XtreamCluster}; use crate::model::xtream::{XtreamCategory, XtreamSeriesInfo, XtreamSeriesInfoEpisode, XtreamStream}; +use crate::repository::storage::hash_string; fn map_to_xtream_category(categories: &Value) -> Result, M3uFilterError> { match serde_json::from_value::>(categories.to_owned()) { @@ -45,15 +46,18 @@ pub fn parse_xtream_series_info(info: &Value, group_title: &str, input: &ConfigI match serde_json::from_value::(info.to_owned()) { Ok(series_info) => { - let result: Vec = series_info.episodes.values().flatten().map(|episode| + let result: Vec = series_info.episodes.values().flatten().map(|episode| { + + let episode_url = create_xtream_series_info_url(url, username, password, episode); PlaylistItem { header: RefCell::new(PlaylistItemHeader { id: Rc::new(episode.id.to_string()), + uuid: Rc::new(hash_string(&episode_url)), name: Rc::new(episode.title.clone()), logo: Rc::new(episode.info.movie_image.clone()), group: Rc::new(group_title.to_string()), title: Rc::new(episode.title.clone()), - url: create_xtream_series_info_url(url, username, password, episode), + url: episode_url, item_type: PlaylistItemType::Series, xtream_cluster: XtreamCluster::Series, additional_properties: episode.get_additional_properties(&series_info), @@ -61,7 +65,7 @@ pub fn parse_xtream_series_info(info: &Value, group_title: &str, input: &ConfigI input_id: input.id, ..Default::default() }) - }).collect(); + }}).collect(); if result.is_empty() { Ok(None) } else { Ok(Some(result)) } } Err(err) => { @@ -109,14 +113,16 @@ pub fn parse_xtream(input: &ConfigInput, if let Some(group) = group_map.get(&stream.category_id) { let mut grp = group.borrow_mut(); let category_name = &grp.category_name; + let stream_url = create_xtream_url(xtream_cluster, url, username, password, &stream); let item = PlaylistItem { header: RefCell::new(PlaylistItemHeader { id: Rc::new(stream.get_stream_id().to_string()), + uuid: Rc::new(hash_string(&stream_url)), name: Rc::clone(&stream.name), logo: Rc::clone(&stream.stream_icon), group: Rc::clone(category_name), title: Rc::clone(&stream.name), - url: create_xtream_url(xtream_cluster, url, username, password, &stream), + url: stream_url, epg_channel_id: stream.epg_channel_id.clone(), item_type: PlaylistItemType::from(xtream_cluster), xtream_cluster, diff --git a/src/processing/xtream_processor.rs b/src/processing/xtream_processor.rs index aed662082..867b2ea05 100644 --- a/src/processing/xtream_processor.rs +++ b/src/processing/xtream_processor.rs @@ -1,25 +1,37 @@ -use crate::m3u_filter_error::M3uFilterError; -use crate::model::config::{ConfigTarget, InputType, TargetType}; -use crate::model::playlist::FetchedPlaylist; +use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; +use crate::model::config::{Config, ConfigInput, ConfigTarget, InputType, TargetType}; +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::IndexedDocumentQuery; use crate::utils::download; +use crate::utils::download::get_xtream_stream_info_content; +use std::collections::HashSet; +use std::rc::Rc; pub async fn playlist_resolve_series(target: &ConfigTarget, errors: &mut Vec, pipe: &ProcessingPipe, - fpl: &mut FetchedPlaylist<'_>, - new_fpl: &mut FetchedPlaylist<'_>) { + provider_fpl: &mut FetchedPlaylist<'_>, + processed_fpl: &mut FetchedPlaylist<'_>) { let (resolve_series, resolve_series_delay) = if let Some(options) = &target.options { - (options.xtream_resolve_series && fpl.input.input_type == InputType::Xtream && target.has_output(&TargetType::M3u), + (options.xtream_resolve_series && provider_fpl.input.input_type == InputType::Xtream && target.has_output(&TargetType::M3u), options.xtream_resolve_series_delay) } else { (false, 0) }; if resolve_series { - let mut series_playlist = download::get_xtream_playlist_series(fpl, errors, resolve_series_delay).await; + // collect all series in the processed lists + let to_process_uuids: HashSet> = processed_fpl.playlistgroups.iter() + .filter(|plg| plg.xtream_cluster == XtreamCluster::Series) + .flat_map(|plg| &plg.channels) + .filter(|pli| pli.header.borrow().item_type == PlaylistItemType::SeriesInfo) + .map(|pli| Rc::clone(&pli.header.borrow().uuid)).collect(); + let mut series_playlist = download::get_xtream_playlist_series(provider_fpl, to_process_uuids, errors, resolve_series_delay).await; // original content saved into original list for plg in &series_playlist { - fpl.update_playlist(plg); + provider_fpl.update_playlist(plg); } // run processing pipe over new items for f in pipe { @@ -30,7 +42,58 @@ pub async fn playlist_resolve_series(target: &ConfigTarget, 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) { + result = match get_xtream_stream_info_content(&info_url, input).await { + Ok(content) => Some(content), + Err(err) => { + errors.push(M3uFilterError::new(M3uFilterErrorKind::Info, format!("{err}"))); + None + } + }; + } + if resolve_delay > 0 { + actix_web::rt::time::sleep(std::time::Duration::new(u64::from(resolve_delay), 0)).await; + } + result +} + +pub async fn playlist_resolve_movies(cfg: &Config, target: &ConfigTarget, errors: &mut Vec, fpl: &FetchedPlaylist<'_>) { + 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) + }; + 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}"))), + } + } + 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 diff --git a/src/repository/indexed_document.rs b/src/repository/indexed_document.rs index 71471b735..7de001aa1 100644 --- a/src/repository/indexed_document.rs +++ b/src/repository/indexed_document.rs @@ -1,17 +1,19 @@ +use std::fmt::Debug; use std::fs::{File, OpenOptions}; use std::io::{BufReader, Error, ErrorKind, Read, Seek, SeekFrom, Write}; use std::marker::PhantomData; use std::path::{Path, PathBuf}; use log::error; - -use crate::repository::bplustree::{BPlusTree, BPlusTreeQuery}; +use serde::{Deserialize, Serialize}; +use crate::repository::bplustree::{BPlusTree, BPlusTreeQuery, BPlusTreeUpdate}; use crate::utils::file_utils; const BLOCK_SIZE: usize = 4096; const LEN_SIZE: usize = 4; pub(in crate::repository) type OffsetPointer = u32; +type SizeType = u32; pub(in crate::repository) struct IndexedDocument {} @@ -33,15 +35,17 @@ impl IndexedDocument { { let mut size_bytes = [0u8; LEN_SIZE]; reader.read_exact(&mut size_bytes)?; - let buf_size = u32::from_le_bytes(size_bytes) as usize; + let buf_size = SizeType::from_le_bytes(size_bytes) as usize; Ok(buf_size) } - pub(in crate::repository) fn get_offset(index_path: &Path, doc_id: u32) -> Result { - match BPlusTreeQuery::::try_new(index_path) { + pub(in crate::repository) fn get_offset(index_path: &Path, doc_id: &K) -> Result + where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug + { + match BPlusTreeQuery::::try_new(index_path) { Ok(mut tree) => { - tree.query(&doc_id).map_or_else(|| Err(Error::new(ErrorKind::NotFound, format!("doc_id not found {doc_id}"))), |offset| Ok(u64::from(offset))) + tree.query(doc_id).map_or_else(|| Err(Error::new(ErrorKind::NotFound, format!("doc_id not found {doc_id:?}"))), |offset| Ok(u64::from(offset))) } Err(err) => Err(err) } @@ -65,17 +69,21 @@ impl IndexedDocument { * * index file is a bplustree */ -pub(in crate::repository) struct IndexedDocumentWriter { +pub(in crate::repository) struct IndexedDocumentWriter +where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug +{ main_path: PathBuf, index_path: PathBuf, main_file: File, main_offset: OffsetPointer, - index_tree: BPlusTree, + index_tree: BPlusTree, dirty: bool, fragmented: bool, } -impl IndexedDocumentWriter { +impl IndexedDocumentWriter +where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug +{ fn new_with_mode(main_path: PathBuf, index_path: PathBuf, append: bool) -> Result { let append_mode = append && main_path.exists(); let mut main_file = if append_mode { @@ -96,7 +104,7 @@ impl IndexedDocumentWriter { // Retrieve file size and convert to `u32` for `main_offset`, if possible let mut main_offset = main_file .metadata() - .and_then(|meta| u32::try_from(meta.len()).map_err(|err| Error::new(ErrorKind::Other, err))) + .and_then(|meta| SizeType::try_from(meta.len()).map_err(|err| Error::new(ErrorKind::Other, err))) .unwrap_or(0); let mut fragmented = false; @@ -109,12 +117,12 @@ impl IndexedDocumentWriter { // Initialize the index tree (BPlusTree) - either by deserializing an existing one or creating a new one let index_tree = if append_mode && index_path.exists() { - BPlusTree::::load(&index_path).unwrap_or_else(|err| { + BPlusTree::::load(&index_path).unwrap_or_else(|err| { error!("Failed to load index {:?}: {}", index_path, err); - BPlusTree::::new() + BPlusTree::::new() }) } else { - BPlusTree::::new() + BPlusTree::::new() }; Ok(Self { @@ -146,7 +154,7 @@ impl IndexedDocumentWriter { } } - pub fn write_doc(&mut self, doc_id: u32, doc: &T) -> Result<(), Error> + pub fn write_doc(&mut self, doc_id: K, doc: &T) -> Result<(), Error> where T: ?Sized + serde::Serialize, { @@ -182,13 +190,13 @@ impl IndexedDocumentWriter { self.dirty = true; - let encoded_bytes_len = u32::try_from(encoded_bytes.len()).map_err(|err| Error::new(ErrorKind::Other, err))?; + let encoded_bytes_len = SizeType::try_from(encoded_bytes.len()).map_err(|err| Error::new(ErrorKind::Other, err))?; self.main_file.write_all(&encoded_bytes_len.to_le_bytes())?; match file_utils::check_write(&self.main_file.write_all(&encoded_bytes)) { Ok(()) => { if new_record_appended { self.index_tree.insert(doc_id, self.main_offset); - let written_bytes = u32::try_from(encoded_bytes.len() + LEN_SIZE).map_err(|err| Error::new(ErrorKind::Other, err))?; + let written_bytes = SizeType::try_from(encoded_bytes.len() + LEN_SIZE).map_err(|err| Error::new(ErrorKind::Other, err))?; self.main_offset += written_bytes; } } @@ -200,7 +208,9 @@ impl IndexedDocumentWriter { } } -impl Drop for IndexedDocumentWriter { +impl Drop for IndexedDocumentWriter +where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug +{ fn drop(&mut self) { let _ = self.store(); } @@ -212,7 +222,7 @@ impl Drop for IndexedDocumentWriter { // //////////////////////////////////////////////////////// -pub(in crate::repository) struct IndexedDocumentReader { +pub(in crate::repository) struct IndexedDocumentReader { main_path: PathBuf, main_file: BufReader, offsets: Vec, @@ -220,14 +230,17 @@ pub(in crate::repository) struct IndexedDocumentReader { failed: bool, t_buffer: Vec, t_type: PhantomData, + k_type: PhantomData, } -impl IndexedDocumentReader { +impl IndexedDocumentReader +where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug +{ pub fn new(main_path: &Path, index_path: &Path) -> Result { if main_path.exists() && index_path.exists() { let mut offsets = Vec::::new(); { - let index_tree = BPlusTree::::load(index_path)?; + let index_tree = BPlusTree::::load(index_path)?; index_tree.traverse(|_, values| { offsets.extend(values); }); @@ -243,6 +256,7 @@ impl IndexedDocumentReader { failed: false, t_buffer: Vec::new(), t_type: PhantomData, + k_type: PhantomData, }) } Err(e) => Err(e) @@ -290,7 +304,7 @@ impl IndexedDocumentReader { } } - pub(in crate::repository) fn read_indexed_item(main_path: &Path, index_path: &Path, doc_id: u32) -> Result + pub(in crate::repository) fn read_indexed_item(main_path: &Path, index_path: &Path, doc_id: &K) -> Result { if main_path.exists() && index_path.exists() { // get the offset from index @@ -304,11 +318,13 @@ impl IndexedDocumentReader { return Ok(item); } } - Err(Error::new(ErrorKind::Other, format!("Failed to read item for id {} - {}", doc_id, main_path.to_str().unwrap()))) + Err(Error::new(ErrorKind::Other, format!("Failed to read item for id {:?} - {}", doc_id, main_path.to_str().unwrap()))) } } -impl Iterator for IndexedDocumentReader { +impl Iterator for IndexedDocumentReader +where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug +{ type Item = T; // Implement the next() method @@ -328,14 +344,16 @@ impl Iterator for IndexedDocumentReader { // //////////////////////////////////////////////////////// -pub(in crate::repository) struct IndexedDocumentGarbageCollector { +pub(in crate::repository) struct IndexedDocumentGarbageCollector { main_path: PathBuf, index_path: PathBuf, main_file: File, - index_tree: BPlusTree, + index_tree: BPlusTree, } -impl IndexedDocumentGarbageCollector { +impl IndexedDocumentGarbageCollector +where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone +{ pub fn new(main_path: PathBuf, index_path: PathBuf) -> Result { if main_path.exists() && index_path.exists() { // Attempt to open the main file in the specified mode (append or not) @@ -348,14 +366,14 @@ impl IndexedDocumentGarbageCollector { // Retrieve file size and convert to `u32` for `main_file`, if possible let size = main_file .metadata() - .and_then(|meta| u32::try_from(meta.len()).map_err(|err| Error::new(ErrorKind::Other, err))) + .and_then(|meta| SizeType::try_from(meta.len()).map_err(|err| Error::new(ErrorKind::Other, err))) .unwrap_or(0); if size < 1 { return Err(Error::new(ErrorKind::UnexpectedEof, format!("File empty main:{main_path:?}"))); } // Initialize the index tree (BPlusTree) - by deserializing an existing one - let index_tree = BPlusTree::::load(&index_path)?; + let index_tree = BPlusTree::::load(&index_path)?; Ok(Self { main_path, @@ -383,9 +401,9 @@ impl IndexedDocumentGarbageCollector { .truncate(true) .open(&gc_main_path)?; - let mut key_offset = Vec::<(u32, OffsetPointer)>::new(); + let mut key_offset = Vec::<(K, OffsetPointer)>::new(); self.index_tree.traverse(|keys, values| { - keys.iter().zip(values.iter()).for_each(|(&key, &offset)| key_offset.push((key, offset))); + keys.iter().zip(values.iter()).for_each(|(key, &offset)| key_offset.push((key.clone(), offset))); }); let fragmented_byte = 0u8.to_le_bytes(); @@ -398,7 +416,7 @@ impl IndexedDocumentGarbageCollector { // read old content self.main_file.seek(SeekFrom::Start(u64::from(offset)))?; self.main_file.read_exact(&mut size_bytes)?; - let buf_size = u32::from_le_bytes(size_bytes) as usize; + let buf_size = SizeType::from_le_bytes(size_bytes) as usize; // ensure buffer capacity if buffer.capacity() < buf_size { buffer.reserve(buf_size - buffer.capacity()); @@ -409,7 +427,7 @@ impl IndexedDocumentGarbageCollector { gc_file.write_all(&size_bytes)?; gc_file.write_all(&buffer[0..buf_size])?; - let pointer = u32::try_from(gc_offset).map_err(|err| Error::new(ErrorKind::Other, err))?; + let pointer = OffsetPointer::try_from(gc_offset).map_err(|err| Error::new(ErrorKind::Other, err))?; self.index_tree.insert(key, pointer); gc_offset += size_bytes.len() + buf_size; // gc_file.stream_position(); } @@ -487,14 +505,14 @@ mod tests { { let size_main_file_4 = std::fs::metadata(&main_path)?.len(); - let mut gc = IndexedDocumentGarbageCollector::new(main_path.clone(), index_path.clone())?; + let mut gc = IndexedDocumentGarbageCollector::::new(main_path.clone(), index_path.clone())?; gc.garbage_collect()?; let size_main_file_5 = std::fs::metadata(&main_path)?.len(); assert!(size_main_file_5 < size_main_file_4, "Failed, the filesize should be less"); } { - let reader = IndexedDocumentReader::::new(&main_path, &index_path)?; + let reader = IndexedDocumentReader::::new(&main_path, &index_path)?; let mut i = 0; for doc in reader { assert_eq!(doc.id, i, "Wrong id"); @@ -506,4 +524,7 @@ mod tests { Ok(()) } -} \ No newline at end of file +} + +pub type IndexedDocumentQuery = BPlusTreeQuery; +pub type IndexedDocumentUpdate = BPlusTreeUpdate; \ No newline at end of file diff --git a/src/repository/kodi_repository.rs b/src/repository/kodi_repository.rs index 04f403464..219e45906 100644 --- a/src/repository/kodi_repository.rs +++ b/src/repository/kodi_repository.rs @@ -75,6 +75,10 @@ static KODY_STYLE: LazyLock = LazyLock::new(|| KodiStyle { whitespace: regex::Regex::new(r"\s+").unwrap(), }); +fn _get_tmdb_id(_virtual_id: u32) -> Option { + // TODO + None +} pub fn kodi_write_strm_playlist(target: &ConfigTarget, cfg: &Config, new_playlist: &[PlaylistGroup], filename: Option<&String>) -> Result<(), M3uFilterError> { if !new_playlist.is_empty() { diff --git a/src/repository/m3u_playlist_iterator.rs b/src/repository/m3u_playlist_iterator.rs index 4d7452b2f..3b7958635 100644 --- a/src/repository/m3u_playlist_iterator.rs +++ b/src/repository/m3u_playlist_iterator.rs @@ -11,7 +11,7 @@ use crate::utils::file_lock_manager::FileReadGuard; pub const M3U_STREAM_PATH: &str = "m3u-stream"; pub struct M3uPlaylistIterator { - reader: IndexedDocumentReader, + reader: IndexedDocumentReader, base_url: String, username: String, password: String, @@ -40,7 +40,7 @@ impl M3uPlaylistIterator { })?; let reader = - IndexedDocumentReader::::new(&m3u_path, &idx_path).map_err(|err| { + IndexedDocumentReader::::new(&m3u_path, &idx_path).map_err(|err| { M3uFilterError::new( M3uFilterErrorKind::Info, format!("Could not deserialize file {m3u_path:?} - {err}"), diff --git a/src/repository/m3u_repository.rs b/src/repository/m3u_repository.rs index d6faf4390..7e0af03df 100644 --- a/src/repository/m3u_repository.rs +++ b/src/repository/m3u_repository.rs @@ -94,6 +94,6 @@ pub async fn m3u_get_item_for_stream_id(cfg: &Config, stream_id: u32, m3u_path: } { let _file_lock = cfg.file_locks.read_lock(m3u_path).await?; - IndexedDocumentReader::::read_indexed_item(m3u_path, idx_path, stream_id) + IndexedDocumentReader::::read_indexed_item(m3u_path, idx_path, &stream_id) } } \ No newline at end of file diff --git a/src/repository/mod.rs b/src/repository/mod.rs index e2d3e6496..27df2e822 100644 --- a/src/repository/mod.rs +++ b/src/repository/mod.rs @@ -1,12 +1,12 @@ +pub mod storage; +pub mod target_id_mapping; +pub mod bplustree; +mod indexed_document; +pub use indexed_document::IndexedDocumentQuery; pub mod playlist_repository; pub mod m3u_repository; pub mod xtream_repository; pub mod epg_repository; pub mod kodi_repository; -pub mod storage; - -mod indexed_document; -pub mod target_id_mapping; -pub mod bplustree; pub mod m3u_playlist_iterator; pub mod xtream_playlist_iterator; diff --git a/src/repository/storage.rs b/src/repository/storage.rs index 676ec7f62..a96224254 100644 --- a/src/repository/storage.rs +++ b/src/repository/storage.rs @@ -2,6 +2,7 @@ use std::path::{Path, PathBuf}; use std::fmt::Write; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::{Config, ConfigInput}; +use crate::model::playlist::UUIDType; use crate::utils::file_utils; pub(in crate::repository) const FILE_SUFFIX_DB: &str = "db"; @@ -9,11 +10,13 @@ pub(in crate::repository) const FILE_SUFFIX_INDEX: &str = "idx"; const FILE_ID_MAPPING: &str = "id_mapping.db"; -pub fn hash_string(url: &str) -> [u8; 32] { - let hash = blake3::hash(url.as_bytes()); - hash.into() // convert to hash array +/// generates a hash from a string +#[inline] +pub fn hash_string(url: &str) -> UUIDType { + blake3::hash(url.as_bytes()).into() } +#[inline] fn hex_encode(bytes: &[u8]) -> String { bytes.iter().fold(String::new(), |mut output, b| { let _ = write!(output, "{b:02X}"); diff --git a/src/repository/target_id_mapping.rs b/src/repository/target_id_mapping.rs index f72da1b5c..82c389e20 100644 --- a/src/repository/target_id_mapping.rs +++ b/src/repository/target_id_mapping.rs @@ -7,7 +7,7 @@ use chrono::Local; use log::error; use serde::{Deserialize, Serialize}; -use crate::model::playlist::PlaylistItemType; +use crate::model::playlist::{PlaylistItemType, UUIDType}; use crate::repository::bplustree::BPlusTree; // TODO make configurable @@ -17,14 +17,14 @@ const EXPIRATION_DURATION: i64 = 86400; pub struct VirtualIdRecord { pub virtual_id: u32, pub provider_id: u32, - pub uuid: [u8; 32], + pub uuid: UUIDType, pub item_type: PlaylistItemType, pub parent_virtual_id: u32, // only for series to hold series info id. pub last_updated: i64, } impl VirtualIdRecord { - fn new(provider_id: u32, virtual_id: u32, item_type: PlaylistItemType, parent_virtual_id: u32, uuid: [u8; 32]) -> Self { + fn new(provider_id: u32, virtual_id: u32, item_type: PlaylistItemType, parent_virtual_id: u32, uuid: UUIDType) -> Self { let last_updated = Local::now().timestamp(); Self { virtual_id, provider_id, uuid, item_type, parent_virtual_id, last_updated } } @@ -42,7 +42,7 @@ pub struct TargetIdMapping { dirty: bool, virtual_id_counter: u32, by_virtual_id: BPlusTree, - by_uuid: BTreeMap<[u8; 32], u32>, + by_uuid: BTreeMap, path: PathBuf, } @@ -71,7 +71,7 @@ impl TargetIdMapping { } } - pub fn insert_entry(&mut self, uuid: [u8; 32], provider_id: u32, item_type: PlaylistItemType, parent_virtual_id: u32) -> u32 { + pub fn insert_entry(&mut self, uuid: UUIDType, provider_id: u32, item_type: PlaylistItemType, parent_virtual_id: u32) -> u32 { match self.by_uuid.get(&uuid) { None => { self.dirty = true; diff --git a/src/repository/xtream_playlist_iterator.rs b/src/repository/xtream_playlist_iterator.rs index 393539e8d..fa5d13b67 100644 --- a/src/repository/xtream_playlist_iterator.rs +++ b/src/repository/xtream_playlist_iterator.rs @@ -3,12 +3,12 @@ use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::{Config, ConfigTarget}; use crate::model::playlist::{XtreamCluster, XtreamPlaylistItem}; use crate::model::xtream::XtreamMappingOptions; -use crate::repository::indexed_document::IndexedDocumentReader; +use crate::repository::indexed_document::{IndexedDocumentReader}; use crate::repository::xtream_repository::{xtream_get_file_paths, xtream_get_storage_path}; use crate::utils::file_lock_manager::FileReadGuard; pub struct XtreamPlaylistIterator { - reader: IndexedDocumentReader, + reader: IndexedDocumentReader, options: XtreamMappingOptions, category_id: u32, _file_lock: FileReadGuard, @@ -27,7 +27,7 @@ impl XtreamPlaylistIterator { M3uFilterError::new(M3uFilterErrorKind::Info, format!("Could not lock document {xtream_path:?}: {err}")) )?; - let reader = IndexedDocumentReader::::new(&xtream_path, &idx_path) + let reader = IndexedDocumentReader::::new(&xtream_path, &idx_path) .map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Info, format!("Could not deserialize file {} - {}", &xtream_path.to_str().unwrap(), err)))?; let options = XtreamMappingOptions::from_target_options(target.options.as_ref()); diff --git a/src/repository/xtream_repository.rs b/src/repository/xtream_repository.rs index b8f0f1bc2..65db4ccd4 100644 --- a/src/repository/xtream_repository.rs +++ b/src/repository/xtream_repository.rs @@ -4,24 +4,24 @@ use std::io::{BufReader, Error, ErrorKind}; use std::path::{Path, PathBuf}; use log::error; -use serde_json::{json, Value}; +use serde_json::{json, Map, Value}; -use crate::{create_m3u_filter_error, create_m3u_filter_error_result}; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::{Config, ConfigTarget}; -use crate::model::playlist::{PlaylistGroup, PlaylistItem, PlaylistItemType, XtreamCluster, XtreamPlaylistItem}; +use crate::model::playlist::{PlaylistEntry, PlaylistGroup, PlaylistItem, PlaylistItemType, XtreamCluster, XtreamPlaylistItem}; use crate::model::xtream::XtreamMappingOptions; -use crate::repository::bplustree::{BPlusTreeQuery, BPlusTreeUpdate}; -use crate::repository::indexed_document::{IndexedDocumentGarbageCollector, IndexedDocumentReader, IndexedDocumentWriter}; -use crate::repository::storage::{FILE_SUFFIX_DB, FILE_SUFFIX_INDEX, get_target_id_mapping_file, get_target_storage_path, hash_string}; +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::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}; +use crate::{create_m3u_filter_error, create_m3u_filter_error_result}; pub static COL_CAT_LIVE: &str = "cat_live"; pub static COL_CAT_SERIES: &str = "cat_series"; pub static COL_CAT_VOD: &str = "cat_vod"; const FILE_SERIES_EPISODES: &str = "series_episodes"; +const FILE_VOD_INFO: &str = "vod_info"; const FILE_SERIES: &str = "series"; pub const FILE_EPG: &str = "epg.xml"; const PATH_XTREAM: &str = "xtream"; @@ -29,6 +29,9 @@ const TAG_CATEGORY_ID: &str = "category_id"; const TAG_CATEGORY_NAME: &str = "category_name"; const TAG_DIRECT_SOURCE: &str = "direct_source"; 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) => { @@ -62,11 +65,15 @@ fn ensure_xtream_storage_path(cfg: &Config, target_name: &str) -> 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}")); return Some((xtream_path, index_path)); + } else if cluster == XtreamCluster::Video { + let xtream_path = storage_path.join(format!("{FILE_VOD_INFO}.{FILE_SUFFIX_DB}")); + let index_path = storage_path.join(format!("{FILE_VOD_INFO}.{FILE_SUFFIX_INDEX}")); + return Some((xtream_path, index_path)); } None } @@ -159,7 +166,7 @@ async fn xtream_garbage_collect(config: &Config, target_name: &str) -> std::io:: 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()?; + IndexedDocumentGarbageCollector::::new(info_path, idx_path)?.garbage_collect()?; } Ok(()) } @@ -201,11 +208,11 @@ pub async fn xtream_write_playlist(target: &ConfigTarget, cfg: &Config, playlist PlaylistItemType::Series => { // we skip resolved series, because this is only necessary when writing m3u files None - }, + } PlaylistItemType::LiveUnknown | PlaylistItemType::LiveHls => { header.category_id = *cat_id; Some(&mut live_col) - }, + } _ => { if header.get_provider_id().is_some() { header.category_id = *cat_id; @@ -278,7 +285,7 @@ async fn xtream_read_item_for_stream_id(cfg: &Config, stream_id: u32, storage_pa 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) } } @@ -286,7 +293,7 @@ async fn xtream_read_series_item_for_stream_id(cfg: &Config, stream_id: u32, sto 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) } } @@ -312,7 +319,7 @@ pub async fn xtream_get_item_for_stream_id( 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 = BPlusTreeQuery::::try_new(&target_id_mapping_file) + 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 @@ -346,8 +353,8 @@ pub async fn xtream_load_rewrite_playlist(cluster: XtreamCluster, config: &Confi } pub async fn xtream_write_series_info(config: &Config, target_name: &str, - series_info_id: u32, - content: &str) -> Result<(), Error> { + 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)); @@ -364,7 +371,7 @@ pub async fn xtream_write_series_info(config: &Config, target_name: &str, { 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) = BPlusTreeUpdate::::try_new(&target_id_mapping_file) { + 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); @@ -375,7 +382,22 @@ pub async fn xtream_write_series_info(config: &Config, target_name: &str, Ok(()) } -// Reads the series info entry if exists, otherwise 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 _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.store()?; + } + 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)?; @@ -386,7 +408,7 @@ pub async fn xtream_load_series_info(config: &Config, target_name: &str, series_ 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 = BPlusTreeQuery::::try_new(&target_id_mapping_file) + 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}")) @@ -407,7 +429,7 @@ pub async fn xtream_load_series_info(config: &Config, target_name: &str, series_ 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) { + 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); @@ -419,12 +441,95 @@ pub async fn xtream_load_series_info(config: &Config, target_name: &str, series_ None } -pub async fn write_and_get_xtream_series_info( +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()?; + + // 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!("Failed to read vod info for id {vod_id} for {target_name}: {}", err); + None + } + }; + } + } + None +} + +pub async fn write_and_get_xtream_vod_info

( config: &Config, target: &ConfigTarget, - pli_series_info: &XtreamPlaylistItem, + pli: &P, content: &str, -) -> Result { +) -> 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 { + 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"))?; @@ -435,6 +540,7 @@ pub async fn write_and_get_xtream_series_info( .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 @@ -442,12 +548,13 @@ pub async fn write_and_get_xtream_series_info( 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!("{}/{}", pli_series_info.url, provider_id)); - let virtual_id = target_id_mapping.insert_entry(uuid, provider_id, PlaylistItemType::SeriesEpisode, pli_series_info.virtual_id); - episode.insert("id".to_string(), Value::String(virtual_id.to_string())); + 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())); @@ -459,7 +566,7 @@ pub async fn write_and_get_xtream_series_info( } 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(), pli_series_info.virtual_id, &result).await.ok(); + xtream_write_series_info(config, target.name.as_str(), virtual_id, &result).await.ok(); Ok(result) } \ No newline at end of file diff --git a/src/utils/download.rs b/src/utils/download.rs index 0a267abbb..34b6b1a7d 100644 --- a/src/utils/download.rs +++ b/src/utils/download.rs @@ -1,13 +1,21 @@ -use std::cmp::Ordering; -use std::path::PathBuf; -use log::{debug, info}; use crate::m3u_filter_error::M3uFilterError; -use crate::model::config::{Config, ConfigInput}; -use crate::model::playlist::{FetchedPlaylist, PlaylistGroup, PlaylistItem, PlaylistItemType, XtreamCluster}; +use crate::model::config::{Config, ConfigInput, ConfigTarget}; +use crate::model::playlist::{FetchedPlaylist, PlaylistEntry, PlaylistGroup, PlaylistItem, PlaylistItemType, UUIDType, XtreamCluster}; use crate::model::xmltv::TVGuide; -use crate::processing::{m3u_parser, xtream_parser}; use crate::processing::xtream_parser::parse_xtream_series_info; +use crate::processing::{m3u_parser, xtream_parser}; +use crate::repository::xtream_repository; use crate::utils::{file_utils, request_utils}; +use log::{debug, info}; +use std::cmp::Ordering; +use std::collections::HashSet; +use std::io::{Error, ErrorKind}; +use std::path::PathBuf; +use std::rc::Rc; + +const ACTION_GET_SERIES_INFO: &str = "get_series_info"; +const ACTION_GET_VOD_INFO: &str = "get_vod_info"; +const ACTION_GET_LIVE_INFO: &str = "get_live_info"; fn prepare_file_path(persist: Option<&String>, working_dir: &str, action: &str) -> Option { let persist_file: Option = @@ -32,7 +40,7 @@ pub async fn get_m3u_playlist(cfg: &Config, input: &ConfigInput, working_dir: &S } } -pub async fn get_xtream_playlist_series(fpl: &mut FetchedPlaylist<'_>, errors: &mut Vec, resolve_delay: u16) -> Vec { +pub async fn get_xtream_playlist_series(fpl: &mut FetchedPlaylist<'_>, process_uuids: HashSet>, errors: &mut Vec, resolve_delay: u16) -> Vec { let input = fpl.input; let mut result: Vec = vec![]; for plg in &mut fpl.playlistgroups { @@ -40,7 +48,7 @@ pub async fn get_xtream_playlist_series(fpl: &mut FetchedPlaylist<'_>, errors: & for pli in &plg.channels { let (fetch_series, series_info_url) = { let mut header = pli.header.borrow_mut(); - let fetch_series = !header.series_fetched && header.item_type == PlaylistItemType::SeriesInfo; + let fetch_series = !header.series_fetched && header.item_type == PlaylistItemType::SeriesInfo && process_uuids.contains(header.get_uuid()); if fetch_series { header.series_fetched = true; } @@ -78,6 +86,61 @@ pub async fn get_xtream_playlist_series(fpl: &mut FetchedPlaylist<'_>, errors: & result } + +pub fn get_xtream_player_api_action_url(input: &ConfigInput, action: &str) -> Option { + if let Some(user_info) = input.get_user_info() { + Some(format!("{}/player_api.php?username={}&password={}&action={}", + &user_info.base_url, + &user_info.username, + &user_info.password, + action + )) + } else { + None + } +} + +pub fn get_xtream_player_api_info_url(input: &ConfigInput, cluster: XtreamCluster, stream_id: u32) -> Option { + let (action, stream_id_field) = match cluster { + XtreamCluster::Live => (ACTION_GET_LIVE_INFO, "live_id"), + XtreamCluster::Video => (ACTION_GET_VOD_INFO, "vod_id"), + XtreamCluster::Series => (ACTION_GET_SERIES_INFO, "series_id"), + }; + get_xtream_player_api_action_url(input, action).map(|action_url| format!("{action_url}&{stream_id_field}={stream_id}")) +} + + +pub async fn get_xtream_stream_info_content(info_url: &str, input: &ConfigInput) -> Result { + request_utils::download_text_content(input, info_url, None).await +} + +pub async fn get_xtream_stream_info

(config: &Config, input: &ConfigInput, target: &ConfigTarget, + pli: &P, info_url: &str, cluster: XtreamCluster) -> Result +where + P: PlaylistEntry, +{ + if cluster == XtreamCluster::Series { + if let Some(content) = xtream_repository::xtream_load_series_info(config, target.name.as_str(), pli.get_virtual_id()).await { + return Ok(content); + } + } else if cluster == XtreamCluster::Video { + if let Some(content) = xtream_repository::xtream_load_vod_info(config, target.name.as_str(), pli.get_virtual_id()).await { + return Ok(content); + } + } + + if let Ok(content) = get_xtream_stream_info_content(info_url, input).await { + return match cluster { + XtreamCluster::Live => Ok(content), + XtreamCluster::Video => xtream_repository::write_and_get_xtream_vod_info(config, target, pli, &content).await, + XtreamCluster::Series => xtream_repository::write_and_get_xtream_series_info(config, target, pli, &content).await, + }; + } + + Err(Error::new(ErrorKind::Other, format!("Cant find stream with id: {}/{}/{}", + target.name.replace(' ', "_").as_str(), &cluster, pli.get_virtual_id()))) +} + fn get_skip_cluster(input: &ConfigInput) -> Vec { let mut skip_cluster = vec![]; if let Some(input_options) = &input.options { @@ -135,10 +198,10 @@ pub async fn get_xtream_playlist(input: &ConfigInput, working_dir: &str) -> (Vec } Err(err) => errors.push(err) } - }, + } (Err(err1), Err(err2)) => { errors.extend([err1, err2]); - }, + } (_, Err(err)) | (Err(err), _) => errors.push(err), } }