From ee689b5f0a11247ea2d227cc65179683b4c25de7 Mon Sep 17 00:00:00 2001 From: euzu Date: Sun, 22 Dec 2024 05:00:59 +0100 Subject: [PATCH] xtrem reolve video/series wip --- bin/{build_raspi.sh => build_aarch64.sh} | 2 - bin/build_armv7.sh | 3 + bin/release.sh | 5 + src/processing/xtream_processor.rs | 107 +-- src/processing/xtream_processor_series.rs | 136 ++-- src/processing/xtream_processor_vod.rs | 167 +++-- src/repository/kodi_repository.rs | 3 +- src/repository/xtream_repository.rs | 798 ++++++++++++---------- test/rest-api.http | 12 +- 9 files changed, 693 insertions(+), 540 deletions(-) rename bin/{build_raspi.sh => build_aarch64.sh} (52%) create mode 100755 bin/build_armv7.sh diff --git a/bin/build_raspi.sh b/bin/build_aarch64.sh similarity index 52% rename from bin/build_raspi.sh rename to bin/build_aarch64.sh index 1ac9c120f..b12e5e6d8 100755 --- a/bin/build_raspi.sh +++ b/bin/build_aarch64.sh @@ -1,5 +1,3 @@ #!/usr/bin/env bash -#echo "building binary for raspi armv7" -#env RUSTFLAGS="--remap-path-prefix $HOME=~" cross build --release --target armv7-unknown-linux-musleabihf echo "building binary for raspi aarch64" env RUSTFLAGS="--remap-path-prefix $HOME=~" cross build --release --target aarch64-unknown-linux-musl diff --git a/bin/build_armv7.sh b/bin/build_armv7.sh new file mode 100755 index 000000000..5214ba8c7 --- /dev/null +++ b/bin/build_armv7.sh @@ -0,0 +1,3 @@ +#!/usr/bin/env bash +echo "building binary for raspi armv7" +env RUSTFLAGS="--remap-path-prefix $HOME=~" cross build --release --target armv7-unknown-linux-musleabihf diff --git a/bin/release.sh b/bin/release.sh index 0bd1e452e..71717dc21 100755 --- a/bin/release.sh +++ b/bin/release.sh @@ -2,6 +2,11 @@ set -e set -o pipefail +rustup target add x86_64-unknown-linux-musl +rustup target add x86_64-pc-windows-gnu +rustup target add armv7-unknown-linux-musleabihf +rustup add aarch64-unknown-linux-musl + if ! command -v cargo-set-version &> /dev/null then echo "cargo-set-version could not be found. install it with 'cargo install cargo-edit'" diff --git a/src/processing/xtream_processor.rs b/src/processing/xtream_processor.rs index e6cf9261f..dcbcfe72c 100644 --- a/src/processing/xtream_processor.rs +++ b/src/processing/xtream_processor.rs @@ -1,20 +1,44 @@ use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::{Config, ConfigInput}; -use crate::model::playlist::{FetchedPlaylist, PlaylistEntry, PlaylistItem, XtreamCluster}; +use crate::model::playlist::{PlaylistEntry, PlaylistItem, XtreamCluster}; use crate::repository::storage::get_input_storage_path; -use crate::repository::xtream_repository::{xtream_get_info_file_paths}; -use crate::repository::IndexedDocumentIndex; use crate::utils::download; use serde_json::Value; -use std::collections::HashSet; +use std::collections::{HashMap}; use std::fs::{File, OpenOptions}; use std::io::{BufWriter, Error, ErrorKind, Write}; - 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! handle_error { + ($stmt:expr) => { + if let Err(err) = $stmt { + errors.push(err); + } + }; + + ($stmt:expr, $map_err:expr) => { + if let Err(err) = $stmt { + errors.push($map_err(err)); + } + }; +} +pub(in crate::processing) use handle_error; + +#[macro_export] +macro_rules! handle_error_and_return { + ($stmt:expr, $map_err:expr) => { + if let Err(err) = $stmt { + errors.push($map_err(err)); + return; + } + }; +} +pub(in crate::processing) use handle_error_and_return; + #[macro_export] macro_rules! create_resolve_options_function_for_xtream_input { ($cluster:ident) => { @@ -32,12 +56,12 @@ macro_rules! create_resolve_options_function_for_xtream_input { } -pub fn get_u32_from_serde_value(value: &Value) -> Option { +pub fn get_u64_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::Number(num_val) => num_val.as_u64(), Value::String(str_val) => { - match str_val.parse::() { - Ok(sid) => Some(sid), + match str_val.parse::() { + Ok(val) => Some(val), Err(_) => None } } @@ -45,6 +69,10 @@ pub fn get_u32_from_serde_value(value: &Value) -> Option { } } +pub fn get_u32_from_serde_value(value: &Value) -> Option { + get_u64_from_serde_value(value).and_then(|val| u32::try_from(val).ok()) +} + 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); @@ -63,32 +91,7 @@ pub(in crate::processing) async fn playlist_resolve_process_playlist_item(pli: & result } - -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, 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_info_ids.insert(*doc_id); } - }); - } - drop(file_lock); - } - Err(err) => errors.push(M3uFilterError::new(M3uFilterErrorKind::Info, format!("{err}"))), - } - } - Ok(None) => errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not create storage path for input {}", &fpl.input.name.as_ref().map_or("?", |v| v)))), - Err(err) => errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not create storage path for input {err}"))), - } - } - processed_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<()> { +pub(in crate::processing) fn write_info_content_to_wal_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())?; @@ -107,10 +110,10 @@ pub(in crate::processing) fn create_resolve_info_wal_files(cfg: &Config, input: 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 info_path = storage_path.join(format!("{file_prefix}_record.{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)); + let info_file = OpenOptions::new().append(true).open(info_path).ok()?; + return Some((content_file, info_file)); } None @@ -118,3 +121,31 @@ pub(in crate::processing) fn create_resolve_info_wal_files(cfg: &Config, input: Err(_) => None } } + + +pub(in crate::processing) fn has_different_ts(ts: &u64, pli: &PlaylistItem, field: &str) -> bool { + pli.header + .borrow() + .additional_properties + .as_ref() + .map_or(false, |v| match v { + Value::Object(map) => { + if let Some(updated) = map.get(field) { + if let Some(update_ts) = get_u64_from_serde_value(updated) { + return update_ts != *ts; + } + } + true + } + _ => true, + }) +} + +pub(in crate::processing) fn should_update_info(pli: &PlaylistItem, processed_provider_ids: &HashMap, field: &str) -> bool { + if let Some(provider_id) = pli.header.borrow_mut().get_provider_id() { + let timestamp = processed_provider_ids.get(&provider_id); + timestamp.is_none() || has_different_ts(timestamp.unwrap(), pli, field) + } else { + false + } +} diff --git a/src/processing/xtream_processor_series.rs b/src/processing/xtream_processor_series.rs index d336e1c89..0e22915e3 100644 --- a/src/processing/xtream_processor_series.rs +++ b/src/processing/xtream_processor_series.rs @@ -1,45 +1,68 @@ +use std::collections::HashMap; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::{Config, ConfigTarget, InputType}; -use crate::model::playlist::{FetchedPlaylist, PlaylistItemType, XtreamCluster}; +use crate::model::playlist::{FetchedPlaylist, PlaylistItem, 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::{create_resolve_options_function_for_xtream_input, handle_error, handle_error_and_return}; 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; +use crate::processing::xtream_processor::{create_resolve_info_wal_files, get_u32_from_serde_value, + get_u64_from_serde_value, playlist_resolve_process_playlist_item, + should_update_info, write_info_content_to_wal_file}; +use crate::repository::bplustree::BPlusTree; +use crate::repository::storage::get_input_storage_path; +use crate::repository::xtream_repository::{xtream_get_record_file_path, xtream_update_input_info_file, xtream_update_input_record_from_wal_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"; +const TAG_SERIES_INFO_SERIES_ID: &str = "series_id"; +const TAG_SERIES_INFO_LAST_MODIFIED: &str = "last_modified"; 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<()> { +fn write_series_info_tmdb_to_wal_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)); - } +async fn read_processed_series_info_ids(cfg: &Config, errors: &mut Vec, fpl: &FetchedPlaylist<'_>, + cluster: XtreamCluster) -> HashMap { + let mut processed_info_ids = HashMap::new(); + { + match get_input_storage_path(fpl.input, &cfg.working_dir) + .map(|storage_path| xtream_get_record_file_path(&storage_path, cluster)).await { + Ok(file_path) => { + match cfg.file_locks.read_lock(&file_path).await { + Ok(file_lock) => { + if let Ok(info_records) = BPlusTree::::load(&file_path) { + info_records.traverse(|keys, timestamps| { + for (provider_id, ts) in keys.iter().zip(timestamps.iter()) { + processed_info_ids.insert(*provider_id, *ts); + } + }); } + drop(file_lock); } - return Some((stream_id, 0)); + 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}"))), } } - None + processed_info_ids +} + +fn extract_info_record_from_series_info(content: &str) -> Option<(u32, u64)> { + let doc = serde_json::from_str::>(content).ok()?; + let provider_id = get_u32_from_serde_value(doc.get(TAG_SERIES_INFO_SERIES_ID)?)?; + let last_modified = get_u64_from_serde_value(doc.get(TAG_SERIES_INFO_LAST_MODIFIED)?) + .unwrap_or(0); + Some((provider_id, last_modified)) +} + +fn should_update_series_info(pli: &PlaylistItem, processed_provider_ids: &HashMap) -> bool { + should_update_info(pli, processed_provider_ids, "last_modified") } pub async fn playlist_resolve_series(cfg: &Config, target: &ConfigTarget, errors: &mut Vec, @@ -53,57 +76,42 @@ pub async fn playlist_resolve_series(cfg: &Config, target: &ConfigTarget, errors // 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 Some((mut wal_file_content, mut wal_file_record)) = create_resolve_info_wal_files(cfg, processed_fpl.input, XtreamCluster::Series) + else { return; }; + + let mut processed_info_ids = read_processed_series_info_ids(cfg, errors, processed_fpl, XtreamCluster::Series).await; + let mut content_writer = BufWriter::new(&wal_file_content); + let mut record_writer = BufWriter::new(&wal_file_record); + let mut content_updated = false; - 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) + .filter(|&plg| plg.xtream_cluster == XtreamCluster::Series) .flat_map(|plg| &plg.channels) - .filter(|pli| pli.header.borrow().item_type == PlaylistItemType::SeriesInfo) + .filter(|&pli| pli.header.borrow().item_type == PlaylistItemType::SeriesInfo) + .filter(|&pli| should_update_series_info(pli, &processed_info_ids)) { - 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; - } - } + // content contains all episodes + if let Some(content) = playlist_resolve_process_playlist_item(pli, processed_fpl.input, errors, resolve_delay, XtreamCluster::Series).await { + if let Some((provider_id, ts)) = extract_info_record_from_series_info(&content) { + handle_error_and_return!(write_info_content_to_wal_file(&mut content_writer, provider_id, &content), |err| M3uFilterError::new( M3uFilterErrorKind::Notify, format!("Failed to resolve series, could not write to content wal file {err}"))); + processed_info_ids.insert(provider_id, ts); + handle_error_and_return!(write_info_record_to_wal_file(&mut record_writer, provider_id, &info_record), |err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve series wal, could not write to record wal file {err}"))); + content_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); - } + if content_updated { + handle_error!(content_writer.flush(), |err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve vod, could not write to wal file {err}"))); + drop(content_writer); + handle_error!(record_writer.flush(), |err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve vod tmdb, could not write to wal file {err}"))); + drop(record_writer); + handle_error!(xtream_update_input_info_file(cfg, processed_fpl.input, &mut wal_file_content, XtreamCluster::Series).await); + handle_error!(xtream_update_input_record_from_wal_file(cfg, processed_fpl.input, &mut wal_file_record, XtreamCluster::Series).await); } + + + Now update provider_fpl + } // diff --git a/src/processing/xtream_processor_vod.rs b/src/processing/xtream_processor_vod.rs index 71aa8c61c..d415c3049 100644 --- a/src/processing/xtream_processor_vod.rs +++ b/src/processing/xtream_processor_vod.rs @@ -1,44 +1,101 @@ +use std::collections::HashMap; +use std::fs::File; +use crate::{create_resolve_options_function_for_xtream_input, handle_error, handle_error_and_return}; 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 crate::model::playlist::{FetchedPlaylist, PlaylistItem, XtreamCluster}; +use crate::processing::xtream_processor::{create_resolve_info_wal_files, get_u32_from_serde_value, + get_u64_from_serde_value, playlist_resolve_process_playlist_item, + should_update_info, write_info_content_to_wal_file}; +use crate::repository::xtream_repository::{xtream_get_record_file_path, xtream_update_input_info_file, + xtream_update_input_record_from_wal_file}; 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::repository::bplustree::BPlusTree; +use crate::repository::storage::get_input_storage_path; 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"; +const TAG_VOD_INFO_ADDED: &str = "added"; + +struct InputVodInfoRecord { + pub(crate) tmdb_id: u32, + pub(crate) ts: u64, +} 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<()> { +async fn read_processed_vod_info_ids(cfg: &Config, errors: &mut Vec, fpl: &FetchedPlaylist<'_>, + cluster: XtreamCluster) -> HashMap { + let mut processed_info_ids = HashMap::new(); + { + match get_input_storage_path(fpl.input, &cfg.working_dir) + .map(|storage_path| xtream_get_record_file_path(&storage_path, cluster)).await { + Ok(file_path) => { + match cfg.file_locks.read_lock(&file_path).await { + Ok(file_lock) => { + if let Ok(info_records) = BPlusTree::::load(&file_path) { + info_records.traverse(|keys, records| { + for (provider_id, record) in keys.iter().zip(records.iter()) { + processed_info_ids.insert(*provider_id, *record.ts); + } + }); + } + drop(file_lock); + } + Err(err) => errors.push(M3uFilterError::new(M3uFilterErrorKind::Info, format!("{err}"))), + } + } + Ok(None) => errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not create storage path for input {}", &fpl.input.name.as_ref().map_or("?", |v| v)))), + Err(err) => errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not create storage path for input {err}"))), + } + } + processed_info_ids +} + +fn extract_info_record_from_vod_info(content: &str) -> Option<(u32, InputVodInfoRecord)> { + let doc = serde_json::from_str::>(content).ok()?; + + let movie_data = doc.get(TAG_VOD_INFO_MOVIE_DATA)?.as_object()?; + let provider_id = get_u32_from_serde_value( + movie_data.get(TAG_VOD_INFO_STREAM_ID)?, + )?; + + let added = movie_data + .get(TAG_VOD_INFO_ADDED) + .and_then(get_u64_from_serde_value) + .unwrap_or(0); + + let info = doc.get(TAG_VOD_INFO_INFO)?.as_object(); + let tmdb_id = info + .and_then(|info| info.get(TAG_VOD_INFO_TMDB_ID)) + .and_then(get_u32_from_serde_value) + .unwrap_or(0); + + Some((provider_id, InputVodInfoRecord { + tmdb_id, + ts: added, + })) +} + + +fn write_vod_info_record_to_wal_file( + writer: &mut BufWriter<&File>, + provider_id: u32, + record: &InputVodInfoRecord, +) -> std::io::Result<()> { writer.write_all(&provider_id.to_le_bytes())?; - writer.write_all(&tmdb_id.to_le_bytes())?; + writer.write_all(&record.tmdb_id.to_le_bytes())?; + writer.write_all(&record.ts.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 + + +fn should_update_vod_info(pli: &PlaylistItem, processed_provider_ids: &HashMap) -> bool { + should_update_info(pli, processed_provider_ids, "added") } pub async fn playlist_resolve_vod(cfg: &Config, target: &ConfigTarget, errors: &mut Vec, fpl: &FetchedPlaylist<'_>) { @@ -47,52 +104,36 @@ pub async fn playlist_resolve_vod(cfg: &Config, target: &ConfigTarget, errors: & // 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 }; + // We collect the content into a wal file and write it once we collected everything. + let Some((mut wal_file_content, mut wal_file_record)) = 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 processed_info_ids = read_processed_vod_info_ids(cfg, errors, fpl, XtreamCluster::Video).await; + let mut content_writer = BufWriter::new(&wal_file_content); + let mut record_writer = BufWriter::new(&wal_file_record); 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 { + + for pli in fpl.playlistgroups.iter() + .flat_map(|plg| &plg.channels) + .filter(|&chan| chan.header.borrow().xtream_cluster == XtreamCluster::Video) + .filter(|&chan| should_update_vod_info(chan, &processed_info_ids)) { 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; - } + if let Some((provider_id, info_record)) = extract_info_record_from_vod_info(&content) { + let ts = info_record.ts; + handle_error_and_return!(write_info_content_to_wal_file(&mut content_writer, provider_id, &content), |err| M3uFilterError::new( M3uFilterErrorKind::Notify, format!("Failed to resolve vod, could not write to content wal file {err}"))); + processed_info_ids.insert(provider_id, ts); + handle_error_and_return!(write_vod_info_record_to_wal_file(&mut record_writer, provider_id, + &info_record), |err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve vod wal, could not write to record wal file {err}"))); 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}"))); - } + handle_error!(content_writer.flush(), |err| 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); - } + handle_error!(record_writer.flush(), |err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve vod tmdb, could not write to wal file {err}"))); + drop(record_writer); + handle_error!(xtream_update_input_info_file(cfg, fpl.input, &mut wal_file_content, XtreamCluster::Video).await); + handle_error!(xtream_update_input_record_from_wal_file(cfg, fpl.input, &mut wal_file_record, XtreamCluster::Video).await); } } diff --git a/src/repository/kodi_repository.rs b/src/repository/kodi_repository.rs index 88c7daefd..ff640f010 100644 --- a/src/repository/kodi_repository.rs +++ b/src/repository/kodi_repository.rs @@ -4,7 +4,6 @@ use crate::model::config::{Config, ConfigTarget}; use crate::model::playlist::PlaylistGroup; use crate::repository::bplustree::BPlusTree; use crate::repository::storage::get_input_storage_path; -use crate::repository::xtream_repository::xtream_get_vod_tmdb_file_path; use crate::utils::file_lock_manager::FileReadGuard; use crate::utils::file_utils; use chrono::Datelike; @@ -108,7 +107,7 @@ async fn get_tmdb_id(cfg: &Config, provider_id: Option, input_id: u16, std::collections::hash_map::Entry::Vacant(entry) => { if let Some(input) = cfg.get_input_by_id(input_id) { if let Ok(tmdb_path) = get_input_storage_path(input, &cfg.working_dir) - .map(|storage_path| xtream_get_vod_tmdb_file_path(&storage_path)) { + .map(|storage_path| xtream_get_record_file_path(&storage_path, cluster)) { if let Ok(file_lock) = cfg.file_locks.read_lock(&tmdb_path).await { if let Ok(tree) = BPlusTree::::load(&tmdb_path) { let tmdb_id = tree.query(&pid).copied(); diff --git a/src/repository/xtream_repository.rs b/src/repository/xtream_repository.rs index da6f67e76..2d272f5b2 100644 --- a/src/repository/xtream_repository.rs +++ b/src/repository/xtream_repository.rs @@ -1,8 +1,8 @@ -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::io::{Seek, SeekFrom}; use std::path::{Path, PathBuf}; use log::error; @@ -11,10 +11,10 @@ use serde_json::{json, Map, Value}; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::{Config, ConfigInput, ConfigTarget}; use crate::model::playlist::{ - PlaylistEntry, PlaylistGroup, PlaylistItem, PlaylistItemType, XtreamCluster, - XtreamPlaylistItem, + PlaylistEntry, PlaylistGroup, PlaylistItem, PlaylistItemType, XtreamCluster, XtreamPlaylistItem, }; use crate::model::xtream::XtreamMappingOptions; +use crate::repository::bplustree::BPlusTree; use crate::repository::indexed_document::{ IndexedDocumentGarbageCollector, IndexedDocumentQuery, IndexedDocumentReader, IndexedDocumentUpdate, IndexedDocumentWriter, @@ -27,14 +27,14 @@ 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}; -use crate::repository::bplustree::BPlusTree; 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_VOD_INFO_TMDB: &str = "vod_info_tmdb"; +const FILE_VOD_INFO_RECORD: &str = "vod_info_record"; +const FILE_SERIES_INFO_RECORD: &str = "series_info_record"; const FILE_SERIES: &str = "series"; pub const FILE_EPG: &str = "epg.xml"; const PATH_XTREAM: &str = "xtream"; @@ -86,7 +86,10 @@ 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}")); @@ -99,8 +102,16 @@ pub fn xtream_get_info_file_paths(storage_path: &Path, cluster: XtreamCluster) - None } -pub fn xtream_get_vod_tmdb_file_path(storage_path: &Path) -> PathBuf { - storage_path.join(format!("{FILE_VOD_INFO_TMDB}.{FILE_SUFFIX_DB}")) +pub fn xtream_get_record_file_path(storage_path: &Path, cluster: XtreamCluster) -> Option { + match cluster { + XtreamCluster::Live => None, + XtreamCluster::Video => { + Some(storage_path.join(format!("{FILE_VOD_INFO_RECORD}.{FILE_SUFFIX_DB}"))) + } + XtreamCluster::Series => { + Some(storage_path.join(format!("{FILE_SERIES_INFO_RECORD}.{FILE_SUFFIX_DB}"))) + } + } } async fn write_playlists_to_file( @@ -243,7 +254,7 @@ pub async fn xtream_write_playlist( XtreamCluster::Series => &mut cat_series_col, XtreamCluster::Video => &mut cat_vod_col, } - .push(json!({ + .push(json!({ TAG_CATEGORY_ID: format!("{}", &cat_id), TAG_CATEGORY_NAME: plg.title.clone(), TAG_PARENT_ID: 0 @@ -291,18 +302,26 @@ pub async fn xtream_write_playlist( match json_write_documents_to_file(&col_path, data) { Ok(()) => {} Err(err) => { - errors.push(format!("Persisting collection failed: {}: {}", &col_path.to_str().unwrap(), err)); + errors.push(format!( + "Persisting collection failed: {}: {}", + &col_path.to_str().unwrap(), + err + )); } } } - match write_playlists_to_file(cfg, &path, - vec![ - (XtreamCluster::Live, &mut live_col), - (XtreamCluster::Video, &mut vod_col), - (XtreamCluster::Series, &mut series_col), - ], - ).await { + match write_playlists_to_file( + cfg, + &path, + vec![ + (XtreamCluster::Live, &mut live_col), + (XtreamCluster::Video, &mut vod_col), + (XtreamCluster::Series, &mut series_col), + ], + ) + .await + { Ok(()) => { if let Err(err) = xtream_garbage_collect(cfg, &target.name).await { if err.kind() != ErrorKind::NotFound { @@ -316,7 +335,11 @@ pub async fn xtream_write_playlist( } if !errors.is_empty() { - return create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "{}", errors.join("\n")); + return create_m3u_filter_error_result!( + M3uFilterErrorKind::Notify, + "{}", + errors.join("\n") + ); } Ok(()) @@ -457,7 +480,7 @@ pub async fn xtream_get_item_for_stream_id( mapping.parent_virtual_id, &storage_path, ) - .await?; + .await?; item.provider_id = mapping.provider_id; Ok(item) } @@ -469,7 +492,7 @@ pub async fn xtream_get_item_for_stream_id( &storage_path, cluster, ) - .await?; + .await?; item.provider_id = mapping.provider_id; Ok(item) } @@ -486,7 +509,7 @@ pub async fn xtream_load_rewrite_playlist( config: &Config, target: &ConfigTarget, category_id: u32, -) -> Result>, M3uFilterError> { +) -> Result>, M3uFilterError> { Ok(Box::new( XtreamPlaylistIterator::new(cluster, config, target, category_id).await?, )) @@ -520,134 +543,58 @@ pub async fn xtream_write_series_info( { let target_id_mapping_file = get_target_id_mapping_file(&target_path); let _file_lock = config - .file_locks - .write_lock(&target_id_mapping_file) - .await?; - if let Ok(mut target_id_mapping) = - IndexedDocumentUpdate::::try_new(&target_id_mapping_file) - { - if let Some(record) = target_id_mapping.query(&series_info_id) { - let new_record = record.copy_update_timestamp(); - let _ = target_id_mapping.update(&series_info_id, new_record); - } - }; - } - - Ok(()) + .file_locks + .write_lock(&target_id_mapping_file) + .await?; + if let Ok(mut target_id_mapping) = + IndexedDocumentUpdate::::try_new(&target_id_mapping_file) + { + if let Some(record) = target_id_mapping.query(&series_info_id) { + let new_record = record.copy_update_timestamp(); + let _ = target_id_mapping.update(&series_info_id, new_record); + } + }; } - 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( + Ok(()) +} + +pub async fn xtream_write_vod_info( + config: &Config, + target_name: &str, + virtual_id: u32, + content: &str, +) -> Result<(), Error> { + let storage_path = try_option_ok!(xtream_get_storage_path(config, target_name)); + let (info_path, idx_path) = try_option_ok!(xtream_get_info_file_paths( &storage_path, XtreamCluster::Video )); - { - let _file_lock = config.file_locks.write_lock(&info_path).await?; - let mut writer = IndexedDocumentWriter::new_append(info_path, idx_path)?; - writer.write_doc(virtual_id, content).map_err(|_| { - Error::new( - ErrorKind::Other, - format!("failed to write xtream vod info for target {target_name}"), - ) - })?; + { + let _file_lock = config.file_locks.write_lock(&info_path).await?; + let mut writer = IndexedDocumentWriter::new_append(info_path, idx_path)?; + writer.write_doc(virtual_id, content).map_err(|_| { + Error::new( + ErrorKind::Other, + format!("failed to write xtream vod info for target {target_name}"), + ) + })?; - writer.store()?; - } - Ok(()) + 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)?; +// Reads the series info entry if exists +pub async fn xtream_load_series_info( + config: &Config, + target_name: &str, + series_id: u32, +) -> Option { + let target_path = get_target_storage_path(config, target_name)?; + let storage_path = xtream_get_storage_path(config, target_name)?; - { - let target_id_mapping_file = get_target_id_mapping_file(&target_path); - let _file_lock = config - .file_locks - .read_lock(&target_id_mapping_file) - .await - .map_err(|err| { - error!( - "Could not lock id mapping for target {target_name}: {}", - err - ); - Error::new( - ErrorKind::Other, - format!("ID mapping load error for target {target_name}"), - ) - }) - .ok()?; - let mut target_id_mapping = - IndexedDocumentQuery::::try_new(&target_id_mapping_file) - .map_err(|err| { - error!( - "Could not load id mapping for target {target_name}: {}", - err - ); - Error::new( - ErrorKind::Other, - format!("ID mapping load error for target {target_name}"), - ) - }) - .ok()?; - - if let Some(id_record) = target_id_mapping.query(&series_id) { - if id_record.is_expired() { - return None; - } - } - } - - let (info_path, idx_path) = xtream_get_info_file_paths(&storage_path, XtreamCluster::Series)?; - - if info_path.exists() && idx_path.exists() { - { - let _file_lock = config - .file_locks - .read_lock(&info_path) - .await - .map_err(|err| { - error!("Could not lock document {:?}: {}", info_path, err); - Error::new( - ErrorKind::Other, - format!("Document Reader error for target {target_name}"), - ) - }) - .ok()?; - return match IndexedDocumentReader::::read_indexed_item( - &info_path, &idx_path, &series_id, - ) { - Ok(content) => Some(content), - Err(err) => { - error!( - "Failed to read series info for id {series_id} for {target_name}: {}", - err - ); - None - } - }; - } - } - None - } - - pub async fn xtream_get_vod_info_mapping( - config: &Config, - target_name: &str, - vod_id: u32, - ) -> Option { - let target_path = get_target_storage_path(config, target_name)?; + { let target_id_mapping_file = get_target_id_mapping_file(&target_path); let _file_lock = config .file_locks @@ -655,9 +602,9 @@ pub async fn xtream_write_series_info( .await .map_err(|err| { error!( - "Could not lock id mapping for target {target_name}: {}", - err - ); + "Could not lock id mapping for target {target_name}: {}", + err + ); Error::new( ErrorKind::Other, format!("ID mapping load error for target {target_name}"), @@ -668,9 +615,9 @@ pub async fn xtream_write_series_info( IndexedDocumentQuery::::try_new(&target_id_mapping_file) .map_err(|err| { error!( - "Could not load id mapping for target {target_name}: {}", - err - ); + "Could not load id mapping for target {target_name}: {}", + err + ); Error::new( ErrorKind::Other, format!("ID mapping load error for target {target_name}"), @@ -678,211 +625,289 @@ pub async fn xtream_write_series_info( }) .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) => { - // 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 - } - }; + if let Some(id_record) = target_id_mapping.query(&series_id) { + if id_record.is_expired() { + return None; } } - None } - pub async fn write_and_get_xtream_vod_info

( - config: &Config, - target: &ConfigTarget, - pli: &P, - content: &str, - ) -> Result - where - P: PlaylistEntry, - { - let mut doc = serde_json::from_str::>(content) - .map_err(|_| Error::new(ErrorKind::Other, "Failed to parse JSON content"))?; + let (info_path, idx_path) = xtream_get_info_file_paths(&storage_path, XtreamCluster::Series)?; - // 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"))?; - - let target_path = get_target_storage_path(config, target.name.as_str()).ok_or_else(|| { - Error::new( - ErrorKind::Other, - format!("Could not find path for target {}", target.name), - ) - })?; - - let episodes = doc - .get_mut("episodes") - .and_then(Value::as_object_mut) - .ok_or_else(|| Error::new(ErrorKind::Other, "No episodes found in content"))?; - - let virtual_id = pli_series_info.get_virtual_id(); + if info_path.exists() && idx_path.exists() { { - let target_id_mapping_file = get_target_id_mapping_file(&target_path); let _file_lock = config .file_locks - .write_lock(&target_id_mapping_file) + .read_lock(&info_path) .await .map_err(|err| { + error!("Could not lock document {:?}: {}", info_path, err); Error::new( ErrorKind::Other, - format!( - "Could not load id mapping for target {} err:{err}", - target.name - ), + format!("Document Reader error for target {target_name}"), ) - })?; - let mut target_id_mapping = TargetIdMapping::new(&target_id_mapping_file); - let options = XtreamMappingOptions::from_target_options(target.options.as_ref()); - - let provider_url = pli_series_info.get_provider_url(); - for episode_list in episodes.values_mut().filter_map(Value::as_array_mut) { - for episode in episode_list.iter_mut().filter_map(Value::as_object_mut) { - if let Some(provider_id) = episode - .get("id") - .and_then(Value::as_str) - .and_then(|id| id.parse::().ok()) - { - let uuid = hash_string(&format!("{provider_url}/{provider_id}")); - let episode_virtual_id = target_id_mapping.insert_entry( - uuid, - provider_id, - PlaylistItemType::SeriesEpisode, - virtual_id, - ); - episode.insert( - "id".to_string(), - Value::String(episode_virtual_id.to_string()), - ); - } - if options.skip_series_direct_source { - episode.insert(TAG_DIRECT_SOURCE.to_string(), Value::String(String::new())); - } + }) + .ok()?; + return match IndexedDocumentReader::::read_indexed_item( + &info_path, &idx_path, &series_id, + ) { + Ok(content) => Some(content), + Err(err) => { + error!( + "Failed to read series info for id {series_id} for {target_name}: {}", + err + ); + None } - } - - drop(target_id_mapping); + }; } - let result = serde_json::to_string(&doc) - .map_err(|_| Error::new(ErrorKind::Other, "Failed to serialize updated series info"))?; - xtream_write_series_info(config, target.name.as_str(), virtual_id, &result) + } + None +} + +pub async fn xtream_get_vod_info_mapping( + config: &Config, + target_name: &str, + vod_id: u32, +) -> Option { + let target_path = get_target_storage_path(config, target_name)?; + let target_id_mapping_file = get_target_id_mapping_file(&target_path); + let _file_lock = config + .file_locks + .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) => { + // 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 + } + }; + } + } + None +} + +pub async fn write_and_get_xtream_vod_info

( + config: &Config, + target: &ConfigTarget, + pli: &P, + content: &str, +) -> Result +where + P: PlaylistEntry, +{ + let mut doc = serde_json::from_str::>(content) + .map_err(|_| Error::new(ErrorKind::Other, "Failed to parse JSON content"))?; + + // 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"))?; + + let target_path = get_target_storage_path(config, target.name.as_str()).ok_or_else(|| { + Error::new( + ErrorKind::Other, + format!("Could not find path for target {}", target.name), + ) + })?; + + let episodes = doc + .get_mut("episodes") + .and_then(Value::as_object_mut) + .ok_or_else(|| Error::new(ErrorKind::Other, "No episodes found in content"))?; + + let virtual_id = pli_series_info.get_virtual_id(); + { + let target_id_mapping_file = get_target_id_mapping_file(&target_path); + let _file_lock = config + .file_locks + .write_lock(&target_id_mapping_file) .await - .ok(); + .map_err(|err| { + Error::new(ErrorKind::Other, format!("Could not load id mapping for target {} err:{err}", target.name )) + })?; + let mut target_id_mapping = TargetIdMapping::new(&target_id_mapping_file); + let options = XtreamMappingOptions::from_target_options(target.options.as_ref()); - Ok(result) - } - - pub async fn xtream_get_input_vod_info( - cfg: &Config, - input: &ConfigInput, - provider_id: u32, - ) -> Option { - if let Ok(Some((info_path, idx_path))) = get_input_storage_path(input, &cfg.working_dir) - .map(|storage_path| xtream_get_info_file_paths(&storage_path, XtreamCluster::Video)) - { - if let Ok(_file_lock) = cfg.file_locks.read_lock(&info_path).await { - if let Ok(content) = IndexedDocumentReader::::read_indexed_item( - &info_path, &idx_path, &provider_id, - ) { - return Some(content); + let provider_url = pli_series_info.get_provider_url(); + for episode_list in episodes.values_mut().filter_map(Value::as_array_mut) { + for episode in episode_list.iter_mut().filter_map(Value::as_object_mut) { + if let Some(provider_id) = episode + .get("id") + .and_then(Value::as_str) + .and_then(|id| id.parse::().ok()) + { + let uuid = hash_string(&format!("{provider_url}/{provider_id}")); + let episode_virtual_id = target_id_mapping.insert_entry( + uuid, + provider_id, + PlaylistItemType::SeriesEpisode, + virtual_id, + ); + episode.insert( + "id".to_string(), + Value::String(episode_virtual_id.to_string()), + ); + } + if options.skip_series_direct_source { + episode.insert(TAG_DIRECT_SOURCE.to_string(), Value::String(String::new())); } } } - None - } - pub async fn xtream_update_input_info_file( - cfg: &Config, - input: &ConfigInput, - 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, cluster)) - { - Ok(Some((info_path, idx_path))) => { - match cfg.file_locks.write_lock(&info_path).await { - Ok(_file_lock) => { - 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) { + drop(target_id_mapping); + } + let result = serde_json::to_string(&doc) + .map_err(|_| Error::new(ErrorKind::Other, "Failed to serialize updated series info"))?; + xtream_write_series_info(config, target.name.as_str(), virtual_id, &result) + .await + .ok(); + + Ok(result) +} + +pub async fn xtream_get_input_vod_info( + cfg: &Config, + input: &ConfigInput, + provider_id: u32, +) -> Option { + if let Ok(Some((info_path, idx_path))) = get_input_storage_path(input, &cfg.working_dir) + .map(|storage_path| xtream_get_info_file_paths(&storage_path, XtreamCluster::Video)) + { + if let Ok(_file_lock) = cfg.file_locks.read_lock(&info_path).await { + if let Ok(content) = IndexedDocumentReader::::read_indexed_item( + &info_path, + &idx_path, + &provider_id, + ) { + return Some(content); + } + } + } + None +} + +pub async fn xtream_update_input_info_file( + cfg: &Config, + input: &ConfigInput, + 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, cluster)) + { + Ok(Some((info_path, idx_path))) => { + match cfg.file_locks.write_lock(&info_path).await { + Ok(_file_lock) => { + 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]; let mut length_bytes = [0u8; 4]; @@ -907,48 +932,91 @@ pub async fn xtream_write_series_info( } 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}"))), } + Err(err) => Err(M3uFilterError::new( + M3uFilterErrorKind::Info, + format!("{err}"), + )), } - Ok(None) => Err(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not create storage path for input {}", &input.name.as_ref().map_or("?", |v| v)))), - Err(err) => Err(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not create storage path for input {err}"))), } + Ok(None) => Err(M3uFilterError::new( + M3uFilterErrorKind::Notify, + format!( + "Could not create storage path for input {}", + &input.name.as_ref().map_or("?", |v| v) + ), + )), + Err(err) => Err(M3uFilterError::new( + M3uFilterErrorKind::Notify, + format!("Could not create storage path for input {err}"), + )), } +} - pub async fn xtream_update_input_vod_tmdb_file( - cfg: &Config, - input: &ConfigInput, - temp_file: &mut File, - ) -> Result<(), M3uFilterError> { - match get_input_storage_path(input, &cfg.working_dir) - .map(|storage_path| xtream_get_vod_tmdb_file_path(&storage_path)) - { - Ok(tmdb_path) => { - match cfg.file_locks.write_lock(&tmdb_path).await { - Ok(_file_lock) => { - temp_file.seek(SeekFrom::Start(0)).map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not read vod tmdb info {err}")))?; - let mut reader = BufReader::new(temp_file); - let mut provider_id_bytes = [0u8; 4]; - let mut tmdb_id_bytes = [0u8; 4]; - let mut tree_tmdb_index: BPlusTree:: = BPlusTree::load(&tmdb_path).unwrap_or_else(|_| BPlusTree::new()); - loop { - if reader.read_exact(&mut provider_id_bytes).is_err() { - break; // End of file - } - let provider_id = u32::from_le_bytes(provider_id_bytes); - if reader.read_exact(&mut tmdb_id_bytes).is_err() { - break; // End of file - } - let tmdb_id = u32::from_le_bytes(tmdb_id_bytes); - tree_tmdb_index.insert(provider_id, tmdb_id); +pub async fn xtream_update_input_record_from_wal_file( + cfg: &Config, + input: &ConfigInput, + wal_file: &mut File, + cluster: XtreamCluster, +) -> Result<(), M3uFilterError> { + match get_input_storage_path(input, &cfg.working_dir) + .map(|storage_path| xtream_get_vod_record_file_path(&storage_path)) + { + Ok(record_path) => { + match cfg.file_locks.write_lock(&record_path).await { + Ok(_file_lock) => { + wal_file.seek(SeekFrom::Start(0)).map_err(|err| { + M3uFilterError::new( + M3uFilterErrorKind::Notify, + format!("Could not read vod wal info {err}"), + ) + })?; + let mut reader = BufReader::new(wal_file); + let mut provider_id_bytes = [0u8; 4]; + let mut tmdb_id_bytes = [0u8; 4]; + let mut ts_bytes = [0u8; 8]; + let mut tree_record_index: BPlusTree = + BPlusTree::load(&record_path).unwrap_or_else(|_| BPlusTree::new()); + loop { + if reader.read_exact(&mut provider_id_bytes).is_err() { + break; // End of file } - tree_tmdb_index.store(&tmdb_path).map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not store vod tmdb info {err}")))?; - Ok(()) + let provider_id = u32::from_le_bytes(provider_id_bytes); + if reader.read_exact(&mut tmdb_id_bytes).is_err() { + break; // End of file + } + let tmdb_id = u32::from_le_bytes(tmdb_id_bytes); + if reader.read_exact(&mut ts_bytes).is_err() { + break; // End of file + } + let ts = u64::from_le_bytes(tmdb_id_bytes); + tree_record_index.insert( + provider_id, + InputInfoRecord { + provider_id, + tmdb_id, + ts, + }, + ); } - Err(err) => Err(M3uFilterError::new(M3uFilterErrorKind::Info, format!("{err}"))), + tree_record_index.store(&record_path).map_err(|err| { + M3uFilterError::new( + M3uFilterErrorKind::Notify, + format!("Could not store vod record info {err}"), + ) + })?; + Ok(()) } + + Err(err) => Err(M3uFilterError::new( + M3uFilterErrorKind::Info, + format!("{err}"), + )), } - Err(err) => Err(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not create storage path for input {err}"))), } - } \ No newline at end of file + Err(err) => Err(M3uFilterError::new( + M3uFilterErrorKind::Notify, + format!("Could not create storage path for input {err}"), + )), + } +} diff --git a/test/rest-api.http b/test/rest-api.http index 34617a712..90a9c1c03 100644 --- a/test/rest-api.http +++ b/test/rest-api.http @@ -39,22 +39,22 @@ GET http://localhost:8901/player_api.php?action=get_live_categories&username=xt& GET http://localhost:8901/player_api.php?action=get_live_streams&username=xt&password=xt&category_id=102 ### xtream vod_categories -GET http://localhost:8901/player_api.php?action=get_vod_categories&username=xt&password=xt +GET http://10.41.41.41:8901/player_api.php?action=username=xt&password=xt&get_vod_categories ### xtream vod info -GET http://localhost:8901/player_api.php?action=get_vod_streams&category_id=3&username=xt&password=xt +GET http://10.41.41.41:8901/player_api.php?username=xt&password=xt&action=get_vod_streams&category_id=44 ### xtream vod info -GET http://localhost:8901/player_api.php?action=get_vod_info&vod_id=2&username=xt&password=xt +GET http://10.41.41.41:8901/player_api.php?username=xt&password=xt&action=get_vod_info&vod_id=8051 ### xtream series_categories -GET http://localhost:8901/player_api.php?action=get_series_categories&username=xt&password=xt +GET http://10.41.41.41:8901/player_api.php?username=xt&password=xt&action=get_series_categories ### xtream series -GET http://localhost:8901/player_api.php?action=get_series&username=xt&password=xt&category_id=120 +GET http://10.41.41.41:8901/player_api.php?username=xt&password=xt&action=get_series&category_id=56 ### xtream series info -GET http://localhost:8901/player_api.php?action=get_series_info&username=xt&password=xt&series_id=16467 +GET http://10.41.41.41:8901/player_api.php?username=xt&password=xt&action=get_series_info&series_id=13923 ### xtream series stream GET http://localhost:8901/series/xt/xt/18137