diff --git a/src/model/playlist.rs b/src/model/playlist.rs index f817ea725..4ab5af07d 100644 --- a/src/model/playlist.rs +++ b/src/model/playlist.rs @@ -151,8 +151,6 @@ pub struct PlaylistItemHeader { pub additional_properties: Option, #[serde(default, skip_serializing, skip_deserializing)] pub item_type: PlaylistItemType, - #[serde(default, skip_serializing, skip_deserializing)] - pub series_fetched: bool, // only used for series_info #[serde(default)] pub category_id: u32, #[serde(default)] @@ -320,7 +318,6 @@ pub struct XtreamPlaylistItem { pub xtream_cluster: XtreamCluster, pub additional_properties: Option, pub item_type: PlaylistItemType, - pub series_fetched: bool, // only used for series_info pub category_id: u32, pub input_id: u16, } @@ -399,7 +396,6 @@ impl PlaylistItem { xtream_cluster: header.xtream_cluster, additional_properties: header.additional_properties.as_ref().and_then(|props| serde_json::to_string(props).ok()), item_type: header.item_type, - series_fetched: header.series_fetched, category_id: header.category_id, input_id: header.input_id, } diff --git a/src/model/xtream.rs b/src/model/xtream.rs index 8eeae1a74..5db8790ac 100644 --- a/src/model/xtream.rs +++ b/src/model/xtream.rs @@ -364,7 +364,7 @@ impl XtreamSeriesInfoEpisode { add_str_property_if_exists!(result, self.info.releasedate, "release_date"); add_str_property_if_exists!(result, self.title, "title"); add_i64_property_if_exists!(result, self.season, "season"); - add_str_property_if_exists!(result, series_info.info.youtube_trailer, "youtube_trailer"); + add_opt_i64_property_if_exists!(result, self.info.tmdb_id, "tmdb_id"); if result.is_empty() { None } else { Some(Value::Object(result)) } } } diff --git a/src/processing/xtream_parser.rs b/src/processing/xtream_parser.rs index d3bb45e5e..ce86a15f0 100644 --- a/src/processing/xtream_parser.rs +++ b/src/processing/xtream_parser.rs @@ -127,7 +127,6 @@ pub fn parse_xtream(input: &ConfigInput, item_type: PlaylistItemType::from(xtream_cluster), xtream_cluster, additional_properties: stream.get_additional_properties(), - series_fetched: false, category_id: 0, input_id, ..Default::default() diff --git a/src/processing/xtream_processor.rs b/src/processing/xtream_processor.rs index 77f8e722f..0e0d280be 100644 --- a/src/processing/xtream_processor.rs +++ b/src/processing/xtream_processor.rs @@ -1,15 +1,15 @@ 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::{FetchedPlaylist, PlaylistEntry, PlaylistItem, PlaylistItemType, XtreamCluster}; use crate::repository::storage::get_input_storage_path; use crate::utils::download; +use crate::{info_err, notify_err}; +use serde::{Deserialize, Serialize}; use serde_json::Value; use std::collections::HashMap; use std::fs::{File, OpenOptions}; use std::io::{BufWriter, Error, ErrorKind, Write}; use std::path::PathBuf; -use serde::{Deserialize, Serialize}; -use crate::{info_err, notify_err}; const FILE_SERIES_INFO: &str = "xtream_series_info"; const FILE_VOD_INFO: &str = "xtream_vod_info"; @@ -98,6 +98,17 @@ pub(in crate::processing) fn write_info_content_to_wal_file(writer: &mut BufWrit Ok(()) } +pub(in crate::processing) fn create_resolve_episode_wal_files(cfg: &Config, input: &ConfigInput) -> Option<(File, PathBuf)> { + match get_input_storage_path(input, &cfg.working_dir) { + Ok(storage_path) => { + let info_path = storage_path.join(format!("series_episode_record.{FILE_SUFFIX_WAL}")); + let info_file = OpenOptions::new().append(true).open(&info_path).ok()?; + Some((info_file, info_path)) + } + Err(_) => None + } +} + pub(in crate::processing) fn create_resolve_info_wal_files(cfg: &Config, input: &ConfigInput, cluster: XtreamCluster) -> Option<(File, File, PathBuf, PathBuf)> { match get_input_storage_path(input, &cfg.working_dir) { Ok(storage_path) => { @@ -119,7 +130,7 @@ pub(in crate::processing) fn create_resolve_info_wal_files(cfg: &Config, input: } -pub(in crate::processing) fn has_different_ts(ts: &u64, pli: &PlaylistItem, field: &str) -> bool { +pub(in crate::processing) fn has_different_ts(ts: u64, pli: &PlaylistItem, field: &str) -> bool { pli.header .borrow() .additional_properties @@ -128,7 +139,7 @@ pub(in crate::processing) fn has_different_ts(ts: &u64, pli: &PlaylistItem, fiel 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; + return update_ts != ts; } } true @@ -140,14 +151,14 @@ pub(in crate::processing) fn has_different_ts(ts: &u64, pli: &PlaylistItem, fiel 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) + timestamp.is_none() || has_different_ts(*timestamp.unwrap(), pli, field) } else { false } } pub(in crate::processing) async fn read_processed_info_ids(cfg: &Config, errors: &mut Vec, fpl: &FetchedPlaylist<'_>, - cluster: XtreamCluster, extract_ts: F) -> HashMap + item_type: PlaylistItemType, extract_ts: F) -> HashMap where F: Fn(&V) -> u64, V: Serialize + for<'de> Deserialize<'de> + Clone, @@ -155,7 +166,7 @@ where let mut processed_info_ids = HashMap::new(); let file_path = match get_input_storage_path(fpl.input, &cfg.working_dir) - .map(|storage_path| xtream_get_record_file_path(&storage_path, cluster)).and_then(|opt| opt.ok_or_else(|| Error::new(ErrorKind::Other, "Not supported".to_string()))) + .map(|storage_path| xtream_get_record_file_path(&storage_path, item_type)).and_then(|opt| opt.ok_or_else(|| Error::new(ErrorKind::Other, "Not supported".to_string()))) { Ok(file_path) => file_path, Err(err) => { diff --git a/src/processing/xtream_processor_series.rs b/src/processing/xtream_processor_series.rs index 56ac9632b..d410c2d6d 100644 --- a/src/processing/xtream_processor_series.rs +++ b/src/processing/xtream_processor_series.rs @@ -2,31 +2,24 @@ use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::{Config, ConfigTarget, InputType}; use crate::model::playlist::{FetchedPlaylist, PlaylistGroup, PlaylistItem, PlaylistItemType, XtreamCluster}; use crate::processing::playlist_processor::ProcessingPipe; -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, read_processed_info_ids, should_update_info, write_info_content_to_wal_file}; -use crate::repository::xtream_repository::{xtream_get_info_file_paths, xtream_update_input_info_file, xtream_update_input_series_record_from_wal_file}; +use crate::processing::xtream_parser::parse_xtream_series_info; +use crate::processing::xtream_processor::{create_resolve_episode_wal_files, create_resolve_info_wal_files, get_u32_from_serde_value, get_u64_from_serde_value, playlist_resolve_process_playlist_item, read_processed_info_ids, should_update_info, write_info_content_to_wal_file}; +use crate::repository::storage::get_input_storage_path; +use crate::repository::xtream_repository::{xtream_get_info_file_paths, xtream_update_input_info_file, xtream_update_input_series_episodes_record_from_wal_file, xtream_update_input_series_record_from_wal_file}; +use crate::repository::IndexedDocumentReader; use crate::{create_resolve_options_function_for_xtream_target, handle_error, handle_error_and_return, info_err, notify_err}; use serde_json::{Map, Value}; -use std::collections::{HashMap}; +use std::collections::HashMap; use std::fs::File; use std::io::{BufWriter, Write}; -use crate::processing::xtream_parser::parse_xtream_series_info; -use crate::repository::{IndexedDocumentReader}; -use crate::repository::storage::get_input_storage_path; const TAG_SERIES_INFO_SERIES_ID: &str = "series_id"; const TAG_SERIES_INFO_LAST_MODIFIED: &str = "last_modified"; create_resolve_options_function_for_xtream_target!(series); -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(()) -} - -async fn read_processed_series_info_ids(cfg: &Config, errors: &mut Vec, fpl: &FetchedPlaylist<'_>, - cluster: XtreamCluster) -> HashMap { - read_processed_info_ids(cfg, errors, fpl, cluster, |ts: &u64| *ts).await +async fn read_processed_series_info_ids(cfg: &Config, errors: &mut Vec, fpl: &FetchedPlaylist<'_>) -> HashMap { + read_processed_info_ids(cfg, errors, fpl, PlaylistItemType::SeriesInfo, |ts: &u64| *ts).await } fn extract_info_record_from_series_info(content: &str) -> Option<(u32, u64)> { @@ -46,19 +39,29 @@ fn write_series_info_record_to_wal_file( Ok(()) } +fn write_series_episode_record_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 should_update_series_info(pli: &PlaylistItem, processed_provider_ids: &HashMap) -> bool { should_update_info(pli, processed_provider_ids, "last_modified") } async fn playlist_resolve_series_info(cfg: &Config, errors: &mut Vec, - processed_fpl: &mut FetchedPlaylist<'_>, resolve_delay: u16) -> HashMap{ + processed_fpl: &mut FetchedPlaylist<'_>, resolve_delay: u16) -> HashMap { // 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_content_file, mut wal_record_file, wal_content_path, wal_record_path)) = create_resolve_info_wal_files(cfg, processed_fpl.input, XtreamCluster::Series) else { return HashMap::new(); }; - let mut processed_info_ids = read_processed_series_info_ids(cfg, errors, processed_fpl, XtreamCluster::Series).await; + let mut processed_info_ids = read_processed_series_info_ids(cfg, errors, processed_fpl).await; let mut content_writer = BufWriter::new(&wal_content_file); let mut record_writer = BufWriter::new(&wal_record_file); let mut content_updated = false; @@ -73,10 +76,10 @@ async fn playlist_resolve_series_info(cfg: &Config, errors: &mut Vec = vec![]; let input = provider_fpl.input; - let (info_path, idx_path) = match get_input_storage_path(input, &cfg.working_dir) + 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::Series)) - { - Ok(Some(paths)) => paths, - _ => { - errors.push(notify_err!("Failed to open input info file for series".to_string())); - return result; - } + else { + errors.push(notify_err!("Failed to open input info file for series".to_string())); + return result; }; - let _file_lock = match cfg.file_locks.read_lock(&info_path).await { - Ok(lock) => lock, - Err(err) => { - errors.push(notify_err!(format!("Could not lock input info file for series: {err}"))); - return result; - } + let Ok(_file_lock) = cfg.file_locks.read_lock(&info_path).await else { + errors.push(notify_err!("Could not lock input info file for series".to_string())); + return result; }; - let mut doc_reader = match IndexedDocumentReader::::new(&info_path, &idx_path) { - Ok(reader) => reader, - Err(_) => return result, + let Ok(mut doc_reader) = IndexedDocumentReader::::new(&info_path, &idx_path) else { return result; }; + + let Some((mut wal_file, wal_path)) = create_resolve_episode_wal_files(cfg, input) else { + errors.push(notify_err!("Could not create wal file for series episodes record".to_string())); + return result; }; + let mut wal_writer = BufWriter::new(&wal_file); for plg in provider_fpl .playlistgroups @@ -141,26 +141,30 @@ async fn process_series_info( .iter() .filter(|pli| pli.header.borrow().item_type == PlaylistItemType::SeriesInfo) { - if let Some(provider_id) = pli.header.borrow_mut().get_provider_id() { - if processed_ids.contains_key(&provider_id) { - if let Ok(content) = doc_reader.get(&provider_id) { - match serde_json::from_str::(&content) { - Ok(series_content) => { - if let Ok(Some(mut series)) = - parse_xtream_series_info(&series_content, pli.header.borrow().group.as_str(), input) - { - group_series.append(&mut series); - - TODO write tmdb ids of episoded to input for kodi export - } + let Some(provider_id) = pli.header.borrow_mut().get_provider_id() else { continue; }; + if !processed_ids.contains_key(&provider_id) { continue; } + let Ok(content) = doc_reader.get(&provider_id) else { continue; }; + match serde_json::from_str::(&content) { + Ok(series_content) => { + if let Ok(Some(mut series)) = + parse_xtream_series_info(&series_content, pli.header.borrow().group.as_str(), input) + { + for episode in &series { + let Some(provider_id) = &episode.header.borrow_mut().get_provider_id() else { continue; }; + let Some(Value::Object(props)) = &episode.header.borrow().additional_properties else { continue; }; + let Some(value) = props.get("tmdb_id") else { continue; }; + if let Some(tmdb_id) = get_u32_from_serde_value(value) { + println!("episode {tmdb_id} = {provider_id}"); + handle_error!(write_series_episode_record_to_wal_file(&mut wal_writer, *provider_id, tmdb_id), + |err| errors.push(info_err!(format!("Failed to write to series episode wal file: {err}")))); } - Err(err) => errors.push(info_err!(format!("Failed to parse JSON: {err}"))), } + group_series.append(&mut series); } } + Err(err) => errors.push(info_err!(format!("Failed to parse JSON: {err}"))), } } - if !group_series.is_empty() { result.push(PlaylistGroup { id: plg.id, @@ -171,6 +175,11 @@ async fn process_series_info( } } + handle_error!(wal_writer.flush(), + |err| errors.push(notify_err!(format!("Failed to resolve series episodes, could not write to wal file {err}")))); + drop(wal_writer); + handle_error!(xtream_update_input_series_episodes_record_from_wal_file(cfg, input, &mut wal_file, &wal_path).await, + |err| errors.push(err)); result } @@ -186,7 +195,7 @@ pub async fn playlist_resolve_series(cfg: &Config, target: &ConfigTarget, let processed_ids = playlist_resolve_series_info(cfg, errors, processed_fpl, resolve_delay).await; if processed_ids.is_empty() { return; } - let mut series_playlist = process_series_info(cfg, provider_fpl, errors, &processed_ids).await; + let mut series_playlist = process_series_info(cfg, provider_fpl, errors, &processed_ids).await; // original content saved into original list for plg in &series_playlist { provider_fpl.update_playlist(plg); @@ -202,8 +211,6 @@ pub async fn playlist_resolve_series(cfg: &Config, target: &ConfigTarget, for plg in &series_playlist { processed_fpl.update_playlist(plg); } - - } // diff --git a/src/processing/xtream_processor_vod.rs b/src/processing/xtream_processor_vod.rs index a759085be..c68b64b0d 100644 --- a/src/processing/xtream_processor_vod.rs +++ b/src/processing/xtream_processor_vod.rs @@ -1,6 +1,6 @@ use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::{Config, ConfigTarget, InputType}; -use crate::model::playlist::{FetchedPlaylist, PlaylistItem, XtreamCluster}; +use crate::model::playlist::{FetchedPlaylist, PlaylistItem, PlaylistItemType, 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, read_processed_info_ids, should_update_info, write_info_content_to_wal_file}; use crate::repository::xtream_repository::{xtream_update_input_info_file, xtream_update_input_vod_record_from_wal_file, InputVodInfoRecord}; use crate::{create_resolve_options_function_for_xtream_target, handle_error, handle_error_and_return, notify_err}; @@ -17,9 +17,8 @@ const TAG_VOD_INFO_ADDED: &str = "added"; create_resolve_options_function_for_xtream_target!(video); -async fn read_processed_vod_info_ids(cfg: &Config, errors: &mut Vec, fpl: &FetchedPlaylist<'_>, - cluster: XtreamCluster) -> HashMap { - read_processed_info_ids(cfg, errors, fpl, cluster, |record: &InputVodInfoRecord| record.ts).await +async fn read_processed_vod_info_ids(cfg: &Config, errors: &mut Vec, fpl: &FetchedPlaylist<'_>) -> HashMap { + read_processed_info_ids(cfg, errors, fpl, PlaylistItemType::Video, |record: &InputVodInfoRecord| record.ts).await } fn extract_info_record_from_vod_info(content: &str) -> Option<(u32, InputVodInfoRecord)> { @@ -71,7 +70,7 @@ pub async fn playlist_resolve_vod(cfg: &Config, target: &ConfigTarget, errors: & let Some((mut wal_content_file, mut wal_record_file, wal_content_path, wal_record_path)) = create_resolve_info_wal_files(cfg, fpl.input, XtreamCluster::Video) else { return; }; - let mut processed_info_ids = read_processed_vod_info_ids(cfg, errors, fpl, XtreamCluster::Video).await; + let mut processed_info_ids = read_processed_vod_info_ids(cfg, errors, fpl).await; let mut content_writer = BufWriter::new(&wal_content_file); let mut record_writer = BufWriter::new(&wal_record_file); let mut content_updated = false; diff --git a/src/repository/bplustree.rs b/src/repository/bplustree.rs index a93963118..2dc950a90 100644 --- a/src/repository/bplustree.rs +++ b/src/repository/bplustree.rs @@ -606,8 +606,7 @@ where V: Serialize + for<'de> Deserialize<'de> + Clone, { pub fn new(tree: &'a BPlusTree) -> Self { - let mut stack = Vec::new(); - stack.push(&tree.root); + let stack = vec![&tree.root]; Self { stack, current_keys: None, diff --git a/src/repository/indexed_document.rs b/src/repository/indexed_document.rs index 4e0a101e4..0a9f7ee52 100644 --- a/src/repository/indexed_document.rs +++ b/src/repository/indexed_document.rs @@ -57,7 +57,7 @@ impl IndexedDocument { //////////////////////////////////////////////////////// /// -/// IndexedDocumentWriter +/// `IndexedDocumentWriter` /// //////////////////////////////////////////////////////// /** @@ -232,7 +232,7 @@ where //////////////////////////////////////////////////////// /// -/// IndexedDocumentReader +/// `IndexedDocumentReader` /// //////////////////////////////////////////////////////// pub struct IndexedDocumentReader @@ -257,8 +257,8 @@ where .read(true) .write(false) .truncate(false) - .open(&main_path)?; - let index_tree = IndexedDocumentIndex::::load(&index_path)?; + .open(main_path)?; + let index_tree = IndexedDocumentIndex::::load(index_path)?; Ok(Self { main_file: BufReader::new(main_file), @@ -285,12 +285,10 @@ where //////////////////////////////////////////////////////// -// -/// IndexedDocumentIterator +/// `IndexedDocumentIterator` /// -/// Iterator | Sequential access with has_next / next +/// Iterator | Sequential access with `has_next` / next //////////////////////////////////////////////////////// - pub(in crate::repository) struct IndexedDocumentIterator { main_path: PathBuf, main_file: BufReader, @@ -391,7 +389,7 @@ where //////////////////////////////////////////////////////// /// -/// IndexedDocumentDirectAccess +/// `IndexedDocumentDirectAccess` /// //////////////////////////////////////////////////////// pub(in crate::repository) struct IndexedDocumentDirectAccess {} @@ -421,10 +419,9 @@ impl IndexedDocumentDirectAccess { //////////////////////////////////////////////////////// /// -/// IndexedDocumentGarbageCollector +/// `IndexedDocumentGarbageCollector` /// //////////////////////////////////////////////////////// - pub(in crate::repository) struct IndexedDocumentGarbageCollector { main_path: PathBuf, index_path: PathBuf, diff --git a/src/repository/kodi_repository.rs b/src/repository/kodi_repository.rs index c746a92fa..63da9d882 100644 --- a/src/repository/kodi_repository.rs +++ b/src/repository/kodi_repository.rs @@ -1,7 +1,7 @@ use crate::{create_m3u_filter_error_result, notify_err}; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::{Config, ConfigTarget}; -use crate::model::playlist::{PlaylistGroup, XtreamCluster}; +use crate::model::playlist::{PlaylistGroup, PlaylistItemType}; use crate::repository::bplustree::BPlusTree; use crate::repository::storage::get_input_storage_path; use crate::utils::file_lock_manager::FileReadGuard; @@ -11,6 +11,7 @@ use log::error; use std::collections::HashMap; use std::fs::File; use std::io::Write; +use std::rc::Rc; use std::sync::LazyLock; use crate::repository::xtream_repository::xtream_get_record_file_path; @@ -84,7 +85,7 @@ fn kodi_style_rename(name: &str, style: &KodiStyle) -> (Vec, String) { (vec![new_name.clone()], new_name) } -static KODY_STYLE: LazyLock = LazyLock::new(|| KodiStyle { +static KODI_STYLE: LazyLock = LazyLock::new(|| KodiStyle { season: regex::Regex::new(r"[Ss]\d\d").unwrap(), episode: regex::Regex::new(r"[Ee]\d\d").unwrap(), year: regex::Regex::new(r"\d\d\d\d").unwrap(), @@ -93,7 +94,7 @@ static KODY_STYLE: LazyLock = LazyLock::new(|| KodiStyle { type InputTmdbIndexMap = HashMap)>>; async fn get_tmdb_id(cfg: &Config, provider_id: Option, input_id: u16, - input_indexes: &mut InputTmdbIndexMap, cluster: XtreamCluster) -> Option { + input_indexes: &mut InputTmdbIndexMap, item_type: PlaylistItemType) -> Option { match provider_id { None => None, Some(pid) => { @@ -108,7 +109,8 @@ 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(Some(tmdb_path)) = get_input_storage_path(input, &cfg.working_dir) - .map(|storage_path| xtream_get_record_file_path(&storage_path, cluster)) { + + .map(|storage_path| xtream_get_record_file_path(&storage_path, item_type)) { 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(); @@ -148,18 +150,26 @@ pub async fn kodi_write_strm_playlist(target: &ConfigTarget, cfg: &Config, new_p for pg in new_playlist { for pli in &pg.channels { - let header = &mut pli.header.borrow_mut(); - let mut dir_path = path.join(sanitize_for_filename(&header.group, underscore_whitespace)); - let mut kodi_file_name = sanitize_for_filename(&header.title, underscore_whitespace); - let mut additional_info = String::new(); - if kodi_style { + let (group, title, item_type, provider_id, input_id, url) = { + let mut header = pli.header.borrow_mut(); + let group = Rc::clone(&header.group); + let title = Rc::clone(&header.title); + let item_type = header.item_type; let provider_id = header.get_provider_id(); let input_id = header.input_id; - let (kodi_file_dir_name, kodi_style_filename) = kodi_style_rename(&kodi_file_name, &KODY_STYLE); + let url = Rc::clone(&header.url); + (group, title, item_type, provider_id, input_id, url) + }; + + let mut dir_path = path.join(sanitize_for_filename(&group, underscore_whitespace)); + let mut kodi_file_name = sanitize_for_filename(&title, underscore_whitespace); + let mut additional_info = String::new(); + if kodi_style { + let (kodi_file_dir_name, kodi_style_filename) = kodi_style_rename(&kodi_file_name, &KODI_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, header.xtream_cluster).await; + let tmdb_id = get_tmdb_id(cfg, provider_id, input_id, &mut input_tmdb_indexes, item_type).await; additional_info = match tmdb_id { None => { String::new() } Some(id) => { format!(" {{tmdb={id}}}") } @@ -172,7 +182,7 @@ pub async fn kodi_write_strm_playlist(target: &ConfigTarget, cfg: &Config, new_p let file_path = dir_path.join(format!("{kodi_file_name}{additional_info}.strm")); match File::create(&file_path) { Ok(mut strm_file) => { - match file_utils::check_write(&strm_file.write_all(header.url.as_bytes())) { + match file_utils::check_write(&strm_file.write_all(url.as_bytes())) { Ok(()) => (), Err(e) => return create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "failed to write strm playlist: {}", e), } diff --git a/src/repository/playlist_repository.rs b/src/repository/playlist_repository.rs index c0a720096..4e1e36f65 100644 --- a/src/repository/playlist_repository.rs +++ b/src/repository/playlist_repository.rs @@ -49,7 +49,7 @@ pub async fn persist_playlist(playlist: &mut [PlaylistGroup], epg: Option<&Epg>, let result = match output.target { TargetType::M3u => m3u_write_playlist(target, cfg, &target_path, playlist).await, TargetType::Xtream => xtream_write_playlist(target, cfg, playlist).await, - TargetType::Strm => kodi_write_strm_playlist(target, cfg, playlist, output.filename.as_ref().map(|x| x.as_str())).await, + TargetType::Strm => kodi_write_strm_playlist(target, cfg, playlist, output.filename.as_deref()).await, }; if let Err(err) = result { diff --git a/src/repository/xtream_repository.rs b/src/repository/xtream_repository.rs index d84dc903b..13cceec98 100644 --- a/src/repository/xtream_repository.rs +++ b/src/repository/xtream_repository.rs @@ -27,6 +27,7 @@ const FILE_SERIES_EPISODES: &str = "series_episodes"; const FILE_VOD_INFO: &str = "vod_info"; const FILE_VOD_INFO_RECORD: &str = "vod_info_record"; const FILE_SERIES_INFO_RECORD: &str = "series_info_record"; +const FILE_SERIES_EPISODE_RECORD: &str = "series_episode_record"; const FILE_SERIES: &str = "series"; pub const FILE_EPG: &str = "epg.xml"; const PATH_XTREAM: &str = "xtream"; @@ -100,18 +101,14 @@ pub fn xtream_get_info_file_paths( None } -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}"))) - } +pub fn xtream_get_record_file_path(storage_path: &Path, item_type: PlaylistItemType) -> Option { + match item_type { + PlaylistItemType::Video => Some(storage_path.join(format!("{FILE_VOD_INFO_RECORD}.{FILE_SUFFIX_DB}"))), + PlaylistItemType::SeriesInfo | PlaylistItemType::Series => Some(storage_path.join(format!("{FILE_SERIES_INFO_RECORD}.{FILE_SUFFIX_DB}"))), + PlaylistItemType::SeriesEpisode => Some(storage_path.join(format!("{FILE_SERIES_EPISODE_RECORD}.{FILE_SUFFIX_DB}"))), + _ => None, } } - async fn write_playlists_to_file( cfg: &Config, storage_path: &Path, @@ -725,7 +722,7 @@ pub async fn xtream_update_input_vod_record_from_wal_file( wal_file: &mut File, wal_path: &Path, ) -> Result<(), M3uFilterError> { - let record_path = get_input_storage_path(input, &cfg.working_dir).map(|storage_path| xtream_get_record_file_path(&storage_path, XtreamCluster::Video)) + let record_path = get_input_storage_path(input, &cfg.working_dir).map(|storage_path| xtream_get_record_file_path(&storage_path, PlaylistItemType::Video)) .map_err(|err| notify_err!(format!("Error accessing storage path: {err}"))) .and_then(|opt| opt.ok_or_else(|| notify_err!(format!("Error accessing storage path for input: {}", input.name.clone().unwrap_or_else(|| input.id.to_string())))))?; @@ -769,12 +766,12 @@ pub async fn xtream_update_input_series_record_from_wal_file( wal_file: &mut File, wal_path: &Path, ) -> Result<(), M3uFilterError> { - let record_path = get_input_storage_path(input, &cfg.working_dir).map(|storage_path| xtream_get_record_file_path(&storage_path, XtreamCluster::Series)) + let record_path = get_input_storage_path(input, &cfg.working_dir).map(|storage_path| xtream_get_record_file_path(&storage_path, PlaylistItemType::SeriesInfo)) .map_err(|err| notify_err!(format!("Error accessing storage path: {err}"))) .and_then(|opt| opt.ok_or_else(|| notify_err!(format!("Error accessing storage path for input: {}", input.name.clone().unwrap_or_else(|| input.id.to_string())))))?; match cfg.file_locks.write_lock(&record_path).await { Ok(_file_lock) => { - wal_file.seek(SeekFrom::Start(0)).map_err(|err| notify_err!(format!("Could not read vod wal info {err}")))?; + wal_file.seek(SeekFrom::Start(0)).map_err(|err| notify_err!(format!("Could not read series wal info {err}")))?; let mut reader = BufReader::new(wal_file); let mut provider_id_bytes = [0u8; 4]; let mut ts_bytes = [0u8; 8]; @@ -793,7 +790,46 @@ pub async fn xtream_update_input_series_record_from_wal_file( tree_record_index.store(&record_path).map_err(|err| notify_err!(format!("Could not store series record info {err}")))?; drop(reader); if let Err(err) = fs::remove_file(wal_path) { - error!("Failed to delete recrd WAL file for series {err}"); + error!("Failed to delete record WAL file for series {err}"); + } + Ok(()) + } + + Err(err) => Err(info_err!(format!("{err}"))), + } +} + +pub async fn xtream_update_input_series_episodes_record_from_wal_file( + cfg: &Config, + input: &ConfigInput, + wal_file: &mut File, + wal_path: &Path, +) -> Result<(), M3uFilterError> { + let record_path = get_input_storage_path(input, &cfg.working_dir).map(|storage_path| xtream_get_record_file_path(&storage_path, PlaylistItemType::SeriesEpisode)) + .map_err(|err| notify_err!(format!("Error accessing storage path: {err}"))) + .and_then(|opt| opt.ok_or_else(|| notify_err!(format!("Error accessing storage path for input: {}", input.name.clone().unwrap_or_else(|| input.id.to_string())))))?; + match cfg.file_locks.write_lock(&record_path).await { + Ok(_file_lock) => { + wal_file.seek(SeekFrom::Start(0)).map_err(|err| notify_err!(format!("Could not read series episode 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 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 + } + 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_record_index.insert(provider_id, tmdb_id); + } + tree_record_index.store(&record_path).map_err(|err| notify_err!(format!("Could not store series episode record info {err}")))?; + drop(reader); + if let Err(err) = fs::remove_file(wal_path) { + error!("Failed to delete record WAL file for series episode {err}"); } Ok(()) } diff --git a/src/utils/download.rs b/src/utils/download.rs index d8292f8e0..13d947331 100644 --- a/src/utils/download.rs +++ b/src/utils/download.rs @@ -29,7 +29,7 @@ fn prepare_file_path(persist: Option<&str>, working_dir: &str, action: &str) -> pub async fn get_m3u_playlist(cfg: &Config, input: &ConfigInput, working_dir: &str) -> (Vec, Vec) { let url = input.url.clone(); - let persist_file_path = prepare_file_path(input.persist.as_ref().map(|x| x.as_str()), working_dir, ""); + let persist_file_path = prepare_file_path(input.persist.as_deref(), working_dir, ""); match request_utils::get_input_text_content(input, working_dir, &url, persist_file_path).await { Ok(text) => { (m3u_parser::parse_m3u(cfg, input, text.lines()), vec![]) @@ -198,8 +198,8 @@ pub async fn get_xtream_playlist(input: &ConfigInput, working_dir: &str) -> (Vec if !skip_cluster.contains(xtream_cluster) { let category_url = format!("{base_url}&action={category}"); let stream_url = format!("{base_url}&action={stream}"); - let category_file_path = prepare_file_path(input.persist.as_ref().map(|x| x.as_str()), working_dir, format!("{category}_").as_str()); - let stream_file_path = prepare_file_path(input.persist.as_ref().map(|x| x.as_str()), working_dir, format!("{stream}_").as_str()); + let category_file_path = prepare_file_path(input.persist.as_deref(), working_dir, format!("{category}_").as_str()); + let stream_file_path = prepare_file_path(input.persist.as_deref(), working_dir, format!("{stream}_").as_str()); match futures::join!( request_utils::get_input_json_content(input, category_url.as_str(), category_file_path), @@ -238,7 +238,7 @@ pub async fn get_xmltv(_cfg: &Config, input: &ConfigInput, working_dir: &str) -> None => (None, vec![]), Some(url) => { debug!("Getting epg file path for url: {}", url); - let persist_file_path = prepare_file_path(input.persist.as_ref().map(|x| x.as_str()), working_dir, "") + let persist_file_path = prepare_file_path(input.persist.as_deref(), working_dir, "") .map(|path| file_utils::add_prefix_to_filename(&path, "epg_", Some("xml"))); match request_utils::get_input_text_content_as_file(input, working_dir, url, persist_file_path).await { diff --git a/test/rest-api.http b/test/rest-api.http index 90a9c1c03..054a5eeae 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://10.41.41.41:8901/player_api.php?action=username=xt&password=xt&get_vod_categories +GET http://localhost:8901/player_api.php?action=username=xt&password=xt&get_vod_categories ### xtream vod info -GET http://10.41.41.41:8901/player_api.php?username=xt&password=xt&action=get_vod_streams&category_id=44 +GET http://localhost:8901/player_api.php?username=xt&password=xt&action=get_vod_streams&category_id=44 ### xtream vod info -GET http://10.41.41.41:8901/player_api.php?username=xt&password=xt&action=get_vod_info&vod_id=8051 +GET http://localhost:8901/player_api.php?username=xt&password=xt&action=get_vod_info&vod_id=8051 ### xtream series_categories -GET http://10.41.41.41:8901/player_api.php?username=xt&password=xt&action=get_series_categories +GET http://localhost:8901/player_api.php?username=xt&password=xt&action=get_series_categories ### xtream series -GET http://10.41.41.41:8901/player_api.php?username=xt&password=xt&action=get_series&category_id=56 +GET http://localhost:8901/player_api.php?username=xt&password=xt&action=get_series&category_id=56 ### xtream series info -GET http://10.41.41.41:8901/player_api.php?username=xt&password=xt&action=get_series_info&series_id=13923 +GET http://localhost: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