From 69e8ad4fd1fa0f69f27c2465cc92d0c47694e2fe Mon Sep 17 00:00:00 2001 From: euzu Date: Thu, 19 Dec 2024 13:02:10 +0100 Subject: [PATCH] resolve series wip --- Cargo.lock | 2 +- Cargo.toml | 2 +- frontend/package.json | 3 +- .../src/component/form-view/from-view.tsx | 2 +- .../playlist-tree/playlist-tree.scss | 20 +- .../component/playlist-tree/playlist-tree.tsx | 16 +- .../src/component/tag-input/tag-input.tsx | 5 +- .../src/component/user-view/user-view.scss | 81 +++--- frontend/src/index.tsx | 8 - frontend/src/model/server-config.ts | 7 + src/model/xtream.rs | 2 +- src/processing/mod.rs | 4 +- src/processing/playlist_processor.rs | 3 +- src/processing/xtream_processor.rs | 250 +++++------------- src/processing/xtream_processor_series.rs | 145 ++++++++++ src/processing/xtream_processor_vod.rs | 98 +++++++ src/repository/kodi_repository.rs | 5 +- src/repository/xtream_repository.rs | 25 +- 18 files changed, 418 insertions(+), 260 deletions(-) create mode 100644 src/processing/xtream_processor_series.rs create mode 100644 src/processing/xtream_processor_vod.rs diff --git a/Cargo.lock b/Cargo.lock index f1a8f6e4c..229a14810 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1780,6 +1780,7 @@ dependencies = [ "log", "mime", "openssl", + "paste", "path-clean", "pest", "pest_derive", @@ -1793,7 +1794,6 @@ dependencies = [ "serde", "serde_json", "serde_yaml", - "tempfile", "time", "tokio", "tokio-stream", diff --git a/Cargo.toml b/Cargo.toml index db68fef6c..70092ae3c 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -53,4 +53,4 @@ bytes = "1.8.0" async-std = "1.13" tokio-stream = { version = "0.1", features = ["sync"] } tokio = "1.41" -tempfile = "3.14" +paste = "1.0" diff --git a/frontend/package.json b/frontend/package.json index 5de936e89..e90cf70e5 100755 --- a/frontend/package.json +++ b/frontend/package.json @@ -5,6 +5,7 @@ "version": "2.0.10", "private": true, "homepage": "./", + "license": "MIT", "scripts": { "react-start": "react-scripts start", "react-build": "react-scripts build", @@ -21,7 +22,7 @@ "notistack": "3.0.1", "react": "18.3.1", "react-dom": "18.3.1", - "react-tag-input-component": "latest", + "react-tag-input-component": "2.0.2", "rxjs": "7.8.1" }, "devDependencies": { diff --git a/frontend/src/component/form-view/from-view.tsx b/frontend/src/component/form-view/from-view.tsx index 2e4bfe1fe..6c5d3ce9e 100644 --- a/frontend/src/component/form-view/from-view.tsx +++ b/frontend/src/component/form-view/from-view.tsx @@ -70,7 +70,7 @@ export default function FormView(props: FormViewProps) { case FormFieldType.MAP: return
case FormFieldType.TAGS: - return + return case FormFieldType.NUMBER: case FormFieldType.TEXT: default: diff --git a/frontend/src/component/playlist-tree/playlist-tree.scss b/frontend/src/component/playlist-tree/playlist-tree.scss index 7574df2ad..4d222fc70 100644 --- a/frontend/src/component/playlist-tree/playlist-tree.scss +++ b/frontend/src/component/playlist-tree/playlist-tree.scss @@ -21,14 +21,14 @@ } .tree-group { - .tree-group-header { + &__header { display: flex; flex-flow: row; align-items: center; height: 1.3rem; white-space: nowrap; - .tree-group-header-content { + &-content { display: flex; flex-flow: row; align-items: center; @@ -39,14 +39,14 @@ accent-color: var(--checkbox-checked-background-color); } - .tree-group-count { + &__count { margin-left: 10px; color: var(--tree-count-color); } } - .tree-group-childs { - .tree-channel { + &__childs { + .tree-group__channel { display: flex; flex-flow: row; align-items: center; @@ -54,7 +54,7 @@ height: 1.3rem; gap: 3px; - .tree-channel-tools { + &-tools { display: flex; flex-flow: row; align-items: center; @@ -62,14 +62,14 @@ } - .tree-channel-content { + &-content { display: flex; flex-flow: row; align-items: center; white-space: nowrap; } - .tree-channel-nr { + &-nr { min-width: 2rem; font-size: 1rem; margin-right: 8px; @@ -83,14 +83,14 @@ } } - .tree-group-header-content, .tree-channel-content { + &__header-content, &__channel-content { cursor: pointer; &:hover { background-color: var(--tree-hover-background-color); color: var(--tree-hover-color); - .tree-group-count { + .tree-group__count { color: var(--tree-hover-color); } } diff --git a/frontend/src/component/playlist-tree/playlist-tree.tsx b/frontend/src/component/playlist-tree/playlist-tree.tsx index eb9c20888..95bc954e0 100644 --- a/frontend/src/component/playlist-tree/playlist-tree.tsx +++ b/frontend/src/component/playlist-tree/playlist-tree.tsx @@ -118,8 +118,8 @@ export default function PlaylistTree(props: PlaylistTreeProps) { }, [videoExtensions]); const renderEntry = useCallback((entry: PlaylistItem, index: number): React.ReactNode => { - return
-
+ return
+
{getIconByName('LinkRounded')}
@@ -139,26 +139,26 @@ export default function PlaylistTree(props: PlaylistTreeProps) { }
-
-
{index + 1}
+
+
{index + 1}
{entry.header.name}
}, [handleClipboardUrl, handlePlayUrl, handleDownloadUrl, isVideoFile, handleWebSearch, serverConfig]); const renderGroup = useCallback((group: PlaylistGroup): React.ReactNode => { return
-
+
{getIconByName(expanded.current[group.id] ? 'ExpandMore' : 'ChevronRight')}
-
+
{group.title} -
({group.channels.length})
+
({group.channels.length})
{expanded.current[group.id] && ( -
+
{group.channels.map(renderEntry)}
)}
; diff --git a/frontend/src/component/tag-input/tag-input.tsx b/frontend/src/component/tag-input/tag-input.tsx index f45e60600..57d0ab649 100644 --- a/frontend/src/component/tag-input/tag-input.tsx +++ b/frontend/src/component/tag-input/tag-input.tsx @@ -6,10 +6,11 @@ interface TagInputProps { name: string; values: string[]; onChange: (name: string, values: string[]) => void; + placeHolder?: string; } export default function TagInput(props: TagInputProps) { - const { name, values, onChange } = props; + const { name, values, onChange, placeHolder } = props; const handleTagsChange = (newTags: string[]) => { onChange(name, newTags); @@ -21,7 +22,7 @@ export default function TagInput(props: TagInputProps) { value={values} onChange={handleTagsChange} name={name} - placeHolder="Add tags..." + placeHolder={placeHolder ?? "Add tags..."} />
); diff --git a/frontend/src/component/user-view/user-view.scss b/frontend/src/component/user-view/user-view.scss index 282a17eb9..a19dc7ebe 100644 --- a/frontend/src/component/user-view/user-view.scss +++ b/frontend/src/component/user-view/user-view.scss @@ -6,19 +6,22 @@ } .user { - @include preferences.prefsRoot(); - @include preferences.prefsToolbar(); + & { + min-height: 400px; + max-height: 480px; - min-height: 400px; - max-height: 480px; + @include preferences.prefsRoot(); + @include preferences.prefsToolbar(); + } &__content { - display: flex; - flex-flow: column; - flex: 1 1 auto; - justify-content: flex-start; - overflow: hidden; - + & { + display: flex; + flex-flow: column; + flex: 1 1 auto; + justify-content: flex-start; + overflow: hidden; + } label { font-weight: bold; @@ -43,31 +46,37 @@ } &__target { - display: flex; - flex-flow: column; - border: 1px solid var(--border-color); - border-radius: var(--border-radius); - overflow: hidden; - background-color: var(--background-color); - gap: 12px; + & { + display: flex; + flex-flow: column; + border: 1px solid var(--border-color); + border-radius: var(--border-radius); + overflow: hidden; + background-color: var(--background-color); + gap: 12px; + } &-target { - font-size: 1.2rem; - align-items: center; - display: flex; - flex-flow: row; - padding: 8px; + & { + font-size: 1.2rem; + align-items: center; + display: flex; + flex-flow: row; + padding: 8px; + } label { flex: 1 1 0; } &-toolbar { - flex: 0 0 auto; + & { + flex: 0 0 auto; + padding-right: 8px; + } button { @include common.iconButton(); } - padding-right: 8px; } } @@ -99,22 +108,26 @@ } &-col { - display: table-cell; - width: auto; - vertical-align: bottom; - padding-left: 4px; - padding-right: 4px; + & { + display: table-cell; + width: auto; + vertical-align: bottom; + padding-left: 4px; + padding-right: 4px; + } input { border: none !important; height: 2rem; } &-toolbar { - display: flex; - flex-flow: row; - align-items: flex-end; - justify-content: flex-end; - padding-left: 4px; + & { + display: flex; + flex-flow: row; + align-items: flex-end; + justify-content: flex-end; + padding-left: 4px; + } span { padding: 4px; transform: scale(0.7) translateY(5px); diff --git a/frontend/src/index.tsx b/frontend/src/index.tsx index 542cd1cfb..1c874a037 100755 --- a/frontend/src/index.tsx +++ b/frontend/src/index.tsx @@ -14,11 +14,3 @@ root.render( ); - -// ReactDOM.render( -// -// -// -// -// -// , document.getElementById('root')); \ No newline at end of file diff --git a/frontend/src/model/server-config.ts b/frontend/src/model/server-config.ts index 72cac010a..3788051fa 100644 --- a/frontend/src/model/server-config.ts +++ b/frontend/src/model/server-config.ts @@ -44,7 +44,14 @@ export interface TargetConfig { kodi_style: boolean, xtream_skip_live_direct_source: boolean, xtream_skip_video_direct_source: boolean, + xtream_skip_series_direct_source: boolean, xtream_resolve_series: boolean, + xtream_resolve_series_delay: number, + xtream_resolve_video: boolean, + xtream_resolve_video_delay: number, + m3u_include_type_in_url: boolean, + m3u_mask_redirect_url: boolean, + share_live_streams: boolean, }, sort: { match_as_ascii: boolean, diff --git a/src/model/xtream.rs b/src/model/xtream.rs index b2ce01232..d0bf4accd 100644 --- a/src/model/xtream.rs +++ b/src/model/xtream.rs @@ -298,7 +298,7 @@ pub struct XtreamSeriesInfoInfo { #[derive(Debug, Clone, Serialize, Deserialize)] pub struct XtreamSeriesInfoEpisodeInfo { - pub tmdb_id: u32, + pub tmdb_id: Option, pub releasedate: String, pub plot: String, pub duration_secs: u32, diff --git a/src/processing/mod.rs b/src/processing/mod.rs index 52cb7fd0c..a0db237a1 100644 --- a/src/processing/mod.rs +++ b/src/processing/mod.rs @@ -4,4 +4,6 @@ pub mod playlist_processor; pub mod xmltv_parser; mod playlist_watch; mod xtream_processor; -mod affix_processor; \ No newline at end of file +mod affix_processor; +mod xtream_processor_vod; +mod xtream_processor_series; \ No newline at end of file diff --git a/src/processing/playlist_processor.rs b/src/processing/playlist_processor.rs index 299a3ce76..3736eeadf 100644 --- a/src/processing/playlist_processor.rs +++ b/src/processing/playlist_processor.rs @@ -26,7 +26,8 @@ use crate::model::stats::{InputStats, PlaylistStats}; use crate::processing::affix_processor::apply_affixes; use crate::processing::playlist_watch::process_group_watch; use crate::processing::xmltv_parser::flatten_tvguide; -use crate::processing::xtream_processor::{playlist_resolve_series, playlist_resolve_vod}; +use crate::processing::xtream_processor_series::playlist_resolve_series; +use crate::processing::xtream_processor_vod::playlist_resolve_vod; use crate::repository::playlist_repository::persist_playlist; use crate::utils::default_utils::default_as_default; use crate::utils::download; diff --git a/src/processing/xtream_processor.rs b/src/processing/xtream_processor.rs index 55b492f88..e6cf9261f 100644 --- a/src/processing/xtream_processor.rs +++ b/src/processing/xtream_processor.rs @@ -1,63 +1,54 @@ 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::model::config::{Config, ConfigInput}; +use crate::model::playlist::{FetchedPlaylist, PlaylistEntry, PlaylistItem, XtreamCluster}; use crate::repository::storage::get_input_storage_path; -use crate::repository::xtream_repository::{xtream_get_info_file_paths, xtream_update_input_vod_info_file, xtream_update_input_vod_tmdb_file}; +use crate::repository::xtream_repository::{xtream_get_info_file_paths}; use crate::repository::IndexedDocumentIndex; use crate::utils::download; -use serde_json::{Map, Value}; +use serde_json::Value; use std::collections::HashSet; -use std::fs::File; +use std::fs::{File, OpenOptions}; use std::io::{BufWriter, Error, ErrorKind, Write}; -use std::rc::Rc; -const TAG_VOD_INFO_INFO: &str = "info"; -const TAG_VOD_INFO_MOVIE_DATA: &str = "movie_data"; -const TAG_VOD_INFO_TMDB_ID: &str = "tmdb_id"; -const TAG_VOD_INFO_STREAM_ID: &str = "stream_id"; -pub async fn playlist_resolve_series(target: &ConfigTarget, errors: &mut Vec, - pipe: &ProcessingPipe, - provider_fpl: &mut FetchedPlaylist<'_>, - processed_fpl: &mut FetchedPlaylist<'_>) { - let (resolve_series, resolve_series_delay) = - if let Some(options) = &target.options { - (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 { - // 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 { - provider_fpl.update_playlist(plg); - } - // run processing pipe over new items - for f in pipe { - let r = f(&mut series_playlist, target); - if let Some(v) = r { - series_playlist = v; +const FILE_SERIES_INFO:&str = "xtream_series_info"; +const FILE_VOD_INFO:&str = "xtream_vod_info"; +const FILE_SUFFIX_WAL:&str = "wal"; + +#[macro_export] +macro_rules! create_resolve_options_function_for_xtream_input { + ($cluster:ident) => { + paste::paste! { // Paste Makro benötigt! + fn [](target: &ConfigTarget, fpl: &FetchedPlaylist) -> (bool, u16) { + let (resolve, resolve_delay) = + target.options.as_ref().map_or((false, 0), |opt| { + (opt.[] && fpl.input.input_type == InputType::Xtream, + opt.[]) + }); + (resolve, resolve_delay) } } - // assign new items to the new playlist - for plg in &series_playlist { - processed_fpl.update_playlist(plg); + }; +} + + +pub fn get_u32_from_serde_value(value: &Value) -> Option { + match value { + Value::Number(num_val) => num_val.as_u64().and_then(|val| u32::try_from(val).ok()), + Value::String(str_val) => { + match str_val.parse::() { + Ok(sid) => Some(sid), + Err(_) => None + } } + _ => None, } } -async fn playlist_resolve_vod_process_playlist_item(pli: &PlaylistItem, input: &ConfigInput, errors: &mut Vec, resolve_delay: u16) -> Option { +pub(in crate::processing) async fn playlist_resolve_process_playlist_item(pli: &PlaylistItem, input: &ConfigInput, errors: &mut Vec, resolve_delay: u16, cluster: XtreamCluster) -> 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) { + if let Some(info_url) = download::get_xtream_player_api_info_url(input, cluster, provider_id) { result = match download::get_xtream_stream_info_content(&info_url, input).await { Ok(content) => Some(content), Err(err) => { @@ -72,145 +63,17 @@ async fn playlist_resolve_vod_process_playlist_item(pli: &PlaylistItem, input: & result } -fn write_vod_info_content_to_temp_file(writer: &mut BufWriter<&File>, provider_id: u32, content: &str) -> std::io::Result<()> { - let length = u32::try_from(content.len()).map_err(|err| Error::new(ErrorKind::Other, err))?; - if length > 0 { - writer.write_all(&provider_id.to_le_bytes())?; - writer.write_all(&length.to_le_bytes())?; - writer.write_all(content.as_bytes())?; - } - Ok(()) -} -fn write_vod_info_tmdb_to_temp_file(writer: &mut BufWriter<&File>, provider_id: u32, tmdb_id: u32) -> std::io::Result<()> { - writer.write_all(&provider_id.to_le_bytes())?; - writer.write_all(&tmdb_id.to_le_bytes())?; - Ok(()) -} - -fn get_resolve_video_options(target: &ConfigTarget, fpl: &FetchedPlaylist) -> (bool, u16) { - let (resolve_movies, resolve_delay) = - target.options.as_ref().map_or((false, 0), |opt| (opt.xtream_resolve_video && fpl.input.input_type == InputType::Xtream, opt.xtream_resolve_video_delay)); - (resolve_movies, resolve_delay) -} - -fn get_u32_from_serde_value(value: &Value) -> Option { - match value { - Value::Number(num_val) => num_val.as_u64().and_then(|val| u32::try_from(val).ok()), - Value::String(str_val) => { - match str_val.parse::() { - Ok(sid) => Some(sid), - Err(_) => None - } - } - _ => None, - } -} - -fn extract_provider_id_and_tmdb_id_from_vod_info(content: &str) -> Option<(u32, u32)> { - if let Ok(mut doc) = serde_json::from_str::>(content) { - if let Some(Value::Object(movie_data)) = doc.get_mut(TAG_VOD_INFO_MOVIE_DATA) { - if let Some(stream_id_value) = movie_data.get(TAG_VOD_INFO_STREAM_ID) { - if let Some(stream_id) = get_u32_from_serde_value(stream_id_value) { - if let Some(Value::Object(info)) = doc.get_mut(TAG_VOD_INFO_INFO) { - if let Some(tmdb_id_value) = info.get(TAG_VOD_INFO_TMDB_ID) { - if let Some(tmdb_id) = get_u32_from_serde_value(tmdb_id_value) { - return Some((stream_id, tmdb_id)); - } - } - } - return Some((stream_id, 0)); - } - } - } - } - None -} - -fn create_resolve_vod_info_temp_files(errors: &mut Vec) -> Option<(File, File)> { - let temp_file_info = match tempfile::tempfile() { - Ok(value) => value, - Err(err) => { - errors.push(M3uFilterError::new(M3uFilterErrorKind::Info, format!("Cant resolve vod, could not create temporary file {err}"))); - return None; - } - }; - let temp_file_tmdb = match tempfile::tempfile() { - Ok(value) => value, - Err(err) => { - errors.push(M3uFilterError::new(M3uFilterErrorKind::Info, format!("Cant resolve vod tmdb, could not create temporary file {err}"))); - return None; - } - }; - Some((temp_file_info, temp_file_tmdb)) -} - -pub async fn playlist_resolve_vod(cfg: &Config, target: &ConfigTarget, errors: &mut Vec, fpl: &FetchedPlaylist<'_>) { - let (resolve_movies, resolve_delay) = get_resolve_video_options(target, fpl); - if !resolve_movies { return; } - - // we cant write to the indexed-document directly because of the write lock and time-consuming operation. - // All readers would be waiting for the lock and the app would be unresponsive. - // We collect the content into a temp file and write it once we collected everything. - let Some((mut temp_file_info, mut temp_file_tmdb)) = create_resolve_vod_info_temp_files(errors) else { return }; - - let mut processed_vod_ids = read_processed_vod_info_ids(cfg, errors, fpl).await; - let mut info_writer = BufWriter::new(&temp_file_info); - let mut tmdb_writer = BufWriter::new(&temp_file_tmdb); - let mut info_updated = false; - let mut tmdb_updated = false; - for pli in fpl.playlistgroups.iter().flat_map(|plg| &plg.channels) { - let a = pli.header.borrow_mut().get_provider_id().as_ref().map_or(false, |pid| processed_vod_ids.contains(pid)); - if !a { - if let Some(content) = playlist_resolve_vod_process_playlist_item(pli, fpl.input, errors, resolve_delay).await { - if let Some((provider_id, tmdb_id)) = extract_provider_id_and_tmdb_id_from_vod_info(&content) { - if let Err(err) = write_vod_info_content_to_temp_file(&mut info_writer, provider_id, &content) { - errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve vod, could not write to temporary file {err}"))); - return; - } - info_updated = true; - processed_vod_ids.insert(provider_id); - if tmdb_id > 0 { - if let Err(err) = write_vod_info_tmdb_to_temp_file(&mut tmdb_writer, provider_id, tmdb_id) { - errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve vod tmdb, could not write to temporary file {err}"))); - return; - } - tmdb_updated = true; - } - } - } - } - } - if info_updated { - if let Err(err) = info_writer.flush() { - errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve vod, could not write to temporary file {err}"))); - } - drop(info_writer); - if let Err(err) = xtream_update_input_vod_info_file(cfg, fpl.input, &mut temp_file_info).await { - errors.push(err); - } - } - if tmdb_updated { - if let Err(err) = tmdb_writer.flush() { - errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve vod tmdb, could not write to temporary file {err}"))); - } - drop(tmdb_writer); - if let Err(err) = xtream_update_input_vod_tmdb_file(cfg, fpl.input, &mut temp_file_tmdb).await { - errors.push(err); - } - } -} - -async fn read_processed_vod_info_ids(cfg: &Config, errors: &mut Vec, fpl: &FetchedPlaylist<'_>) -> HashSet { - let mut processed_vod_ids = HashSet::new(); +pub(in crate::processing) async fn read_processed_info_ids(cfg: &Config, errors: &mut Vec, fpl: &FetchedPlaylist<'_>, cluster: XtreamCluster) -> HashSet { + let mut processed_info_ids = HashSet::new(); { - match get_input_storage_path(fpl.input, &cfg.working_dir).map(|storage_path| xtream_get_info_file_paths(&storage_path, XtreamCluster::Video)) { + match get_input_storage_path(fpl.input, &cfg.working_dir).map(|storage_path| xtream_get_info_file_paths(&storage_path, cluster)) { Ok(Some((file_path, idx_path))) => { match cfg.file_locks.read_lock(&file_path).await { Ok(file_lock) => { if let Ok(info_id_mapping) = IndexedDocumentIndex::::load(&idx_path) { info_id_mapping.traverse(|keys, _| { - for doc_id in keys { processed_vod_ids.insert(*doc_id); } + for doc_id in keys { processed_info_ids.insert(*doc_id); } }); } drop(file_lock); @@ -222,5 +85,36 @@ async fn read_processed_vod_info_ids(cfg: &Config, errors: &mut Vec errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not create storage path for input {err}"))), } } - processed_vod_ids -} \ No newline at end of file + processed_info_ids +} + +pub(in crate::processing) fn write_info_content_to_temp_file(writer: &mut BufWriter<&File>, provider_id: u32, content: &str) -> std::io::Result<()> { + let length = u32::try_from(content.len()).map_err(|err| Error::new(ErrorKind::Other, err))?; + if length > 0 { + writer.write_all(&provider_id.to_le_bytes())?; + writer.write_all(&length.to_le_bytes())?; + writer.write_all(content.as_bytes())?; + } + Ok(()) +} + +pub(in crate::processing) fn create_resolve_info_wal_files(cfg: &Config, input: &ConfigInput, cluster: XtreamCluster) -> Option<(File, File)> { + match get_input_storage_path(input, &cfg.working_dir) { + Ok(storage_path) => { + if let Some(file_prefix) = match cluster { + XtreamCluster::Live => None, + XtreamCluster::Video => Some(FILE_SERIES_INFO), + XtreamCluster::Series => Some(FILE_VOD_INFO) + } { + let content_path = storage_path.join(format!("{file_prefix}_content.{FILE_SUFFIX_WAL}")); + let tmdb_path = storage_path.join(format!("{file_prefix}_tmdb.{FILE_SUFFIX_WAL}")); + let content_file = OpenOptions::new().append(true).open(content_path).ok()?; + let tmdb_file = OpenOptions::new().append(true).open(tmdb_path).ok()?; + return Some((content_file, tmdb_file)); + } + None + + } + Err(_) => None + } +} diff --git a/src/processing/xtream_processor_series.rs b/src/processing/xtream_processor_series.rs new file mode 100644 index 000000000..d336e1c89 --- /dev/null +++ b/src/processing/xtream_processor_series.rs @@ -0,0 +1,145 @@ +use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; +use crate::model::config::{Config, ConfigTarget, InputType}; +use crate::model::playlist::{FetchedPlaylist, PlaylistItemType, XtreamCluster}; +use serde_json::{Map, Value}; +use std::fs::File; +use std::io::{BufWriter, Write}; +use crate::create_resolve_options_function_for_xtream_input; +use crate::processing::playlist_processor::ProcessingPipe; +use crate::processing::xtream_processor::{create_resolve_info_wal_files, playlist_resolve_process_playlist_item, read_processed_info_ids, write_info_content_to_temp_file}; +use crate::repository::xtream_repository::xtream_update_input_info_file; + +const TAG_series_INFO_INFO: &str = "info"; +const TAG_series_INFO_MOVIE_DATA: &str = "movie_data"; +const TAG_series_INFO_TMDB_ID: &str = "tmdb_id"; +const TAG_series_INFO_STREAM_ID: &str = "stream_id"; + +create_resolve_options_function_for_xtream_input!(series); + +fn write_series_info_tmdb_to_temp_file(writer: &mut BufWriter<&File>, provider_id: u32, tmdb_id: u32) -> std::io::Result<()> { + writer.write_all(&provider_id.to_le_bytes())?; + writer.write_all(&tmdb_id.to_le_bytes())?; + Ok(()) +} + +fn extract_provider_id_and_tmdb_id_from_series_info(content: &str) -> Option<(u32, u32)> { + if let Ok(mut doc) = serde_json::from_str::>(content) { + if let Some(Value::Object(movie_data)) = doc.get_mut(TAG_series_INFO_MOVIE_DATA) { + if let Some(stream_id_value) = movie_data.get(TAG_series_INFO_STREAM_ID) { + if let Some(stream_id) = crate::processing::xtream_processor::get_u32_from_serde_value(stream_id_value) { + if let Some(Value::Object(info)) = doc.get_mut(TAG_series_INFO_INFO) { + if let Some(tmdb_id_value) = info.get(TAG_series_INFO_TMDB_ID) { + if let Some(tmdb_id) = crate::processing::xtream_processor::get_u32_from_serde_value(tmdb_id_value) { + return Some((stream_id, tmdb_id)); + } + } + } + return Some((stream_id, 0)); + } + } + } + } + None +} + +pub async fn playlist_resolve_series(cfg: &Config, target: &ConfigTarget, errors: &mut Vec, + pipe: &ProcessingPipe, + provider_fpl: &mut FetchedPlaylist<'_>, + processed_fpl: &mut FetchedPlaylist<'_> +) { + let (resolve_series, resolve_delay) = get_resolve_series_options(target, processed_fpl); + if !resolve_series { return; } + + // we cant write to the indexed-document directly because of the write lock and time-consuming operation. + // All readers would be waiting for the lock and the app would be unresponsive. + // We collect the content into a wal file and write it once we collected everything. + let Some((mut wal_file_info, mut wal_file_tmdb)) = create_resolve_info_wal_files(cfg, processed_fpl.input, XtreamCluster::Series) else { return }; + + let mut processed_series_ids = read_processed_info_ids(cfg, errors, processed_fpl, XtreamCluster::Series).await; + let mut info_writer = BufWriter::new(&wal_file_info); + let mut tmdb_writer = BufWriter::new(&wal_file_tmdb); + let mut info_updated = false; + let mut tmdb_updated = false; + for pli in 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) + { + let processed_entry = pli.header.borrow_mut().get_provider_id().as_ref().map_or(false, |pid| processed_series_ids.contains(pid)); + if !processed_entry { + if let Some(content) = playlist_resolve_process_playlist_item(pli, processed_fpl.input, errors, resolve_delay, XtreamCluster::Series).await { + if let Some((provider_id, tmdb_id)) = extract_provider_id_and_tmdb_id_from_series_info(&content) { + if let Err(err) = write_info_content_to_temp_file(&mut info_writer, provider_id, &content) { + errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve series, could not write to temporary file {err}"))); + return; + } + info_updated = true; + processed_series_ids.insert(provider_id); + if tmdb_id > 0 { + if let Err(err) = write_series_info_tmdb_to_temp_file(&mut tmdb_writer, provider_id, tmdb_id) { + errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve series tmdb, could not write to temporary file {err}"))); + return; + } + tmdb_updated = true; + } + } + } + } + } + if info_updated { + if let Err(err) = info_writer.flush() { + errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve series, could not write to temporary file {err}"))); + } + drop(info_writer); + if let Err(err) = xtream_update_input_info_file(cfg, processed_fpl.input, &mut wal_file_info, XtreamCluster::Series).await { + errors.push(err); + } + } + if tmdb_updated { + if let Err(err) = tmdb_writer.flush() { + errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve series tmdb, could not write to temporary file {err}"))); + } + drop(tmdb_writer); + if let Err(err) = xtream_update_input_series_tmdb_file(cfg, processed_fpl.input, &mut wal_file_tmdb).await { + errors.push(err); + } + } +} + +// +// pub(in crate::processing) async fn playlist_resolve_series(target: &ConfigTarget, errors: &mut Vec, +// pipe: &ProcessingPipe, +// provider_fpl: &mut FetchedPlaylist<'_>, +// processed_fpl: &mut FetchedPlaylist<'_>) { +// let (resolve_series, resolve_series_delay) = +// if let Some(options) = &target.options { +// (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 { +// // 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 { +// provider_fpl.update_playlist(plg); +// } +// // run processing pipe over new items +// for f in pipe { +// let r = f(&mut series_playlist, target); +// if let Some(v) = r { +// series_playlist = v; +// } +// } +// // assign new items to the new playlist +// for plg in &series_playlist { +// processed_fpl.update_playlist(plg); +// } +// } +// } \ No newline at end of file diff --git a/src/processing/xtream_processor_vod.rs b/src/processing/xtream_processor_vod.rs new file mode 100644 index 000000000..71aa8c61c --- /dev/null +++ b/src/processing/xtream_processor_vod.rs @@ -0,0 +1,98 @@ +use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; +use crate::model::config::{Config, ConfigTarget, InputType}; +use crate::model::playlist::{FetchedPlaylist, XtreamCluster}; +use crate::processing::xtream_processor::{create_resolve_info_wal_files, playlist_resolve_process_playlist_item, read_processed_info_ids, write_info_content_to_temp_file}; +use crate::repository::xtream_repository::{xtream_update_input_info_file, xtream_update_input_vod_tmdb_file}; +use serde_json::{Map, Value}; +use std::fs::File; +use std::io::{BufWriter, Write}; +use crate::create_resolve_options_function_for_xtream_input; + +const TAG_VOD_INFO_INFO: &str = "info"; +const TAG_VOD_INFO_MOVIE_DATA: &str = "movie_data"; +const TAG_VOD_INFO_TMDB_ID: &str = "tmdb_id"; +const TAG_VOD_INFO_STREAM_ID: &str = "stream_id"; + +create_resolve_options_function_for_xtream_input!(video); + +fn write_vod_info_tmdb_to_temp_file(writer: &mut BufWriter<&File>, provider_id: u32, tmdb_id: u32) -> std::io::Result<()> { + writer.write_all(&provider_id.to_le_bytes())?; + writer.write_all(&tmdb_id.to_le_bytes())?; + Ok(()) +} + +fn extract_provider_id_and_tmdb_id_from_vod_info(content: &str) -> Option<(u32, u32)> { + if let Ok(mut doc) = serde_json::from_str::>(content) { + if let Some(Value::Object(movie_data)) = doc.get_mut(TAG_VOD_INFO_MOVIE_DATA) { + if let Some(stream_id_value) = movie_data.get(TAG_VOD_INFO_STREAM_ID) { + if let Some(stream_id) = crate::processing::xtream_processor::get_u32_from_serde_value(stream_id_value) { + if let Some(Value::Object(info)) = doc.get_mut(TAG_VOD_INFO_INFO) { + if let Some(tmdb_id_value) = info.get(TAG_VOD_INFO_TMDB_ID) { + if let Some(tmdb_id) = crate::processing::xtream_processor::get_u32_from_serde_value(tmdb_id_value) { + return Some((stream_id, tmdb_id)); + } + } + } + return Some((stream_id, 0)); + } + } + } + } + None +} + +pub async fn playlist_resolve_vod(cfg: &Config, target: &ConfigTarget, errors: &mut Vec, fpl: &FetchedPlaylist<'_>) { + let (resolve_movies, resolve_delay) = get_resolve_video_options(target, fpl); + if !resolve_movies { return; } + + // we cant write to the indexed-document directly because of the write lock and time-consuming operation. + // All readers would be waiting for the lock and the app would be unresponsive. + // We collect the content into a temp file and write it once we collected everything. + let Some((mut wal_file_config, mut wal_file_tmdb)) = create_resolve_info_wal_files(cfg, fpl.input, XtreamCluster::Video) else { return }; + + let mut processed_vod_ids = read_processed_info_ids(cfg, errors, fpl, XtreamCluster::Video).await; + let mut content_writer = BufWriter::new(&wal_file_config); + let mut tmdb_writer = BufWriter::new(&wal_file_tmdb); + let mut content_updated = false; + let mut tmdb_updated = false; + for pli in fpl.playlistgroups.iter().flat_map(|plg| &plg.channels).filter(|chan| chan.header.borrow().xtream_cluster == XtreamCluster::Video) { + let processed_entry = pli.header.borrow_mut().get_provider_id().as_ref().map_or(false, |pid| processed_vod_ids.contains(pid)); + if !processed_entry { + if let Some(content) = playlist_resolve_process_playlist_item(pli, fpl.input, errors, resolve_delay, XtreamCluster::Video).await { + if let Some((provider_id, tmdb_id)) = extract_provider_id_and_tmdb_id_from_vod_info(&content) { + if let Err(err) = write_info_content_to_temp_file(&mut content_writer, provider_id, &content) { + errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve vod, could not write to temporary file {err}"))); + return; + } + content_updated = true; + processed_vod_ids.insert(provider_id); + if tmdb_id > 0 { + if let Err(err) = write_vod_info_tmdb_to_temp_file(&mut tmdb_writer, provider_id, tmdb_id) { + errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve vod tmdb, could not write to temporary file {err}"))); + return; + } + tmdb_updated = true; + } + } + } + } + } + if content_updated { + if let Err(err) = content_writer.flush() { + errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve vod, could not write to wal file {err}"))); + } + drop(content_writer); + if let Err(err) = xtream_update_input_info_file(cfg, fpl.input, &mut wal_file_config, XtreamCluster::Video).await { + errors.push(err); + } + } + if tmdb_updated { + if let Err(err) = tmdb_writer.flush() { + errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve vod tmdb, could not write to wal file {err}"))); + } + drop(tmdb_writer); + if let Err(err) = xtream_update_input_vod_tmdb_file(cfg, fpl.input, &mut wal_file_tmdb).await { + errors.push(err); + } + } +} diff --git a/src/repository/kodi_repository.rs b/src/repository/kodi_repository.rs index dbb83b573..88c7daefd 100644 --- a/src/repository/kodi_repository.rs +++ b/src/repository/kodi_repository.rs @@ -155,8 +155,8 @@ pub async fn kodi_write_strm_playlist(target: &ConfigTarget, cfg: &Config, new_p if kodi_style { let provider_id = header.get_provider_id(); let input_id = header.input_id; - let (kodi_file_dir_name, file_name) = kodi_style_rename(&kodi_file_name, &KODY_STYLE); - kodi_file_name = file_name; + let (kodi_file_dir_name, kodi_style_filename) = kodi_style_rename(&kodi_file_name, &KODY_STYLE); + kodi_file_name = kodi_style_filename; kodi_file_dir_name.iter().for_each(|p| dir_path = dir_path.join(p)); let tmdb_id = get_tmdb_id(cfg, provider_id, input_id, &mut input_tmdb_indexes).await; @@ -164,7 +164,6 @@ pub async fn kodi_write_strm_playlist(target: &ConfigTarget, cfg: &Config, new_p None => { String::new() } Some(id) => { format!(" {{tmdb={id}}}") } }; - } if let Err(e) = std::fs::create_dir_all(&dir_path) { error!("cant create directory: {:?}", &dir_path); diff --git a/src/repository/xtream_repository.rs b/src/repository/xtream_repository.rs index d67f00f9c..da6f67e76 100644 --- a/src/repository/xtream_repository.rs +++ b/src/repository/xtream_repository.rs @@ -1,5 +1,6 @@ use std::io::{Seek, SeekFrom}; use std::collections::HashMap; +use std::fs; use std::fs::File; use std::io::{BufReader, Error, ErrorKind, Read}; use std::path::{Path, PathBuf}; @@ -714,7 +715,7 @@ pub async fn xtream_write_series_info( &info_path, &idx_path, &vod_id, ) { Ok(content) => Some(content), - Err(err) => { + Err(_err) => { // this is not an error, it means the info is not indexed // error!("Failed to read vod info for id {vod_id} for {target_name}: {}",err); None @@ -867,19 +868,20 @@ pub async fn xtream_write_series_info( None } - pub async fn xtream_update_input_vod_info_file( + pub async fn xtream_update_input_info_file( cfg: &Config, input: &ConfigInput, - temp_file: &mut File, + wal_file: &mut File, + cluster: XtreamCluster ) -> Result<(), M3uFilterError> { match get_input_storage_path(input, &cfg.working_dir) - .map(|storage_path| xtream_get_info_file_paths(&storage_path, XtreamCluster::Video)) + .map(|storage_path| xtream_get_info_file_paths(&storage_path, cluster)) { Ok(Some((info_path, idx_path))) => { match cfg.file_locks.write_lock(&info_path).await { Ok(_file_lock) => { - temp_file.seek(SeekFrom::Start(0)).map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not read vod info {err}")))?; - let mut reader = BufReader::new(temp_file); + wal_file.seek(SeekFrom::Start(0)).map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not read {cluster} info {err}")))?; + let mut reader = BufReader::new(wal_file); match IndexedDocumentWriter::::new_append(info_path, idx_path) { Ok(mut writer) => { let mut provider_id_bytes = [0u8; 4]; @@ -889,18 +891,21 @@ pub async fn xtream_write_series_info( break; // End of file } let provider_id = u32::from_le_bytes(provider_id_bytes); - reader.read_exact(&mut length_bytes).map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not read temporary vod info {err}")))?; + reader.read_exact(&mut length_bytes).map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not read temporary {cluster} info {err}")))?; let length = u32::from_le_bytes(length_bytes) as usize; let mut buffer = vec![0u8; length]; - reader.read_exact(&mut buffer).map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not read temporary vod info {err}")))?; + reader.read_exact(&mut buffer).map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not read temporary {cluster} info {err}")))?; if let Ok(content) = String::from_utf8(buffer) { let _ = writer.write_doc(provider_id, &content); } } - writer.store().map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not store vod info {err}")))?; + writer.store().map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not store {cluster} info {err}")))?; + if let Err(err) = fs::remove_file(wal_file) { + error!("Failed to delete WAL file for {cluster}"); + } 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::Notify, format!("Could not create create indexed document writer for {cluster} info {err}"))), } } Err(err) => Err(M3uFilterError::new(M3uFilterErrorKind::Info, format!("{err}"))),