diff --git a/src/api/xtream_api.rs b/src/api/xtream_api.rs index dc7a03643..c7d921e8a 100644 --- a/src/api/xtream_api.rs +++ b/src/api/xtream_api.rs @@ -13,12 +13,12 @@ use futures::Stream; use log::{debug, error, warn, Level}; use serde_json::{Map, Value}; +use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::api::api_utils::{get_user_server_info, get_user_target, get_user_target_by_credentials, is_stream_share_enabled, serve_file, stream_response}; use crate::api::model::app_state::AppState; use crate::api::model::request::UserApiRequest; use crate::api::model::xtream::XtreamAuthorizationResponse; -use crate::debug_if_enabled; -use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; +use crate::{debug_if_enabled, info_err}; use crate::model::api_proxy::{ProxyType, ProxyUserCredentials}; use crate::model::config::TargetType; use crate::model::config::{Config, ConfigInput, ConfigTarget}; @@ -445,7 +445,7 @@ async fn xtream_player_api( skip_flag_optional!(skip_vod, xtream_repository::xtream_load_rewrite_playlist(XtreamCluster::Video, &app_state.config, target, category_id).await), ACTION_GET_SERIES => skip_flag_optional!(skip_series, xtream_repository::xtream_load_rewrite_playlist(XtreamCluster::Series, &app_state.config, target, category_id).await), - _ => Some(Err(M3uFilterError::new(M3uFilterErrorKind::Info, format!("Cant find action: {action} for target: {}", &target.name)) + _ => Some(Err(info_err!(format!("Cant find action: {action} for target: {}", &target.name)) )), }; diff --git a/src/filter.rs b/src/filter.rs index 7b584a20e..1ae3747cc 100644 --- a/src/filter.rs +++ b/src/filter.rs @@ -13,7 +13,7 @@ use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::ItemField; use crate::model::playlist::{PlaylistItem, PlaylistItemType}; use crate::utils::directed_graph::DirectedGraph; -use crate::{create_m3u_filter_error_result, exit}; +use crate::{create_m3u_filter_error_result, exit, info_err}; pub fn get_field_value(pli: &PlaylistItem, field: &ItemField) -> Rc { let header = pli.header.borrow(); @@ -432,7 +432,7 @@ pub fn get_filter(filter_text: &str, templates: Option<&Vec>) - if !errors.is_empty() { errors.push(format!("Unable to parse filter: {}", &filter_text)); - return Err(M3uFilterError::new(M3uFilterErrorKind::Info, errors.join("\n"))); + return Err(info_err!(errors.join("\n"))); } result.map_or_else(|| create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "Unable to parse filter: {}", &filter_text), Ok) diff --git a/src/m3u_filter_error.rs b/src/m3u_filter_error.rs index cb05b2e8e..671e8c51c 100644 --- a/src/m3u_filter_error.rs +++ b/src/m3u_filter_error.rs @@ -21,6 +21,20 @@ macro_rules! get_errors_notify_message { } } +#[macro_export] +macro_rules! notify_err { + ($text:expr) => { + M3uFilterError::new(M3uFilterErrorKind::Notify, $text) + }; +} + +#[macro_export] +macro_rules! info_err { + ($text:expr) => { + M3uFilterError::new(M3uFilterErrorKind::Info, $text) + }; +} + #[derive(Debug, PartialEq, Eq)] pub enum M3uFilterErrorKind { // do not send with messaging diff --git a/src/model/api_proxy.rs b/src/model/api_proxy.rs index 05cf2b55c..162ee9763 100644 --- a/src/model/api_proxy.rs +++ b/src/model/api_proxy.rs @@ -4,7 +4,7 @@ use std::str::FromStr; use enum_iterator::Sequence; use log::debug; -use crate::create_m3u_filter_error_result; +use crate::{create_m3u_filter_error_result, info_err}; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::utils::config_reader; @@ -269,10 +269,7 @@ impl ApiProxyConfig { if errors.is_empty() { Ok(()) } else { - Err(M3uFilterError::new( - M3uFilterErrorKind::Info, - errors.join("\n"), - )) + Err(info_err!(errors.join("\n"))) } } diff --git a/src/model/config.rs b/src/model/config.rs index fab932e04..5e24ae7a4 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -15,6 +15,7 @@ use path_clean::PathClean; use url::Url; use crate::filter::{get_filter, prepare_templates, Filter, MockValueProcessor, PatternTemplate, ValueProvider}; +use crate::info_err; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::messaging::MsgKind; use crate::model::api_proxy::{ApiProxyConfig, ProxyUserCredentials}; @@ -359,7 +360,7 @@ impl ConfigTarget { pub fn prepare(&mut self, id: u16, templates: Option<&Vec>) -> Result<(), M3uFilterError> { self.id = id; if self.output.is_empty() { - return Err(M3uFilterError::new(M3uFilterErrorKind::Info, format!("Missing output format for {}", self.name))); + return Err(info_err!(format!("Missing output format for {}", self.name))); } let mut m3u_cnt = 0; let mut strm_cnt = 0; @@ -562,7 +563,7 @@ impl ConfigInput { pub fn prepare(&mut self, id: u16) -> Result<(), M3uFilterError> { self.id = id; if self.url.trim().is_empty() { - return Err(M3uFilterError::new(M3uFilterErrorKind::Info, "url for input is mandatory".to_string())); + return Err(info_err!("url for input is mandatory".to_string())); } if let Some(user_name) = &self.username { if user_name.trim().is_empty() { @@ -582,7 +583,7 @@ impl ConfigInput { } InputType::Xtream => { if self.username.is_none() || self.password.is_none() { - return Err(M3uFilterError::new(M3uFilterErrorKind::Info, "for input type xtream: username and password are mandatory".to_string())); + return Err(info_err!("for input type xtream: username and password are mandatory".to_string())); } } } diff --git a/src/model/mapping.rs b/src/model/mapping.rs index d0420ebd0..24282bf24 100644 --- a/src/model/mapping.rs +++ b/src/model/mapping.rs @@ -15,7 +15,7 @@ use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::{ItemField, AFFIX_FIELDS, COUNTER_FIELDS, MAPPER_ATTRIBUTE_FIELDS}; use crate::model::playlist::{FieldAccessor, PlaylistItem}; use crate::utils::string_utils::Capitalize; -use crate::{create_m3u_filter_error_result, handle_m3u_filter_error_result, valid_property}; +use crate::{create_m3u_filter_error_result, handle_m3u_filter_error_result, info_err, valid_property}; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default)] pub struct MappingTag { @@ -209,25 +209,25 @@ impl Mapper { pub fn prepare(&mut self, templates: Option<&Vec>, tags: Option<&Vec>) -> Result<(), M3uFilterError> { for key in self.attributes.keys() { if !valid_property!(key.as_str(), MAPPER_ATTRIBUTE_FIELDS) { - return Err(M3uFilterError::new(M3uFilterErrorKind::Info, format!("Invalid mapper attribute field {key}"))); + return Err(info_err!(format!("Invalid mapper attribute field {key}"))); } } for key in self.suffix.keys() { if !valid_property!(key.as_str(), AFFIX_FIELDS) { - return Err(M3uFilterError::new(M3uFilterErrorKind::Info, format!("Invalid mapper suffix field {key}"))); + return Err(info_err!(format!("Invalid mapper suffix field {key}"))); } } for key in self.prefix.keys() { if !valid_property!(key.as_str(), AFFIX_FIELDS) { - return Err(M3uFilterError::new(M3uFilterErrorKind::Info, format!("Invalid mapper prefix field {key}"))); + return Err(info_err!(format!("Invalid mapper prefix field {key}"))); } } for (key, value) in &self.assignments { if !valid_property!(key.as_str(), MAPPER_ATTRIBUTE_FIELDS) { - return Err(M3uFilterError::new(M3uFilterErrorKind::Info, format!("Invalid mapper assignment field {key}"))); + return Err(info_err!(format!("Invalid mapper assignment field {key}"))); } if !valid_property!(value.as_str(), MAPPER_ATTRIBUTE_FIELDS) { - return Err(M3uFilterError::new(M3uFilterErrorKind::Info, format!("Invalid mapper assignment field {value}"))); + return Err(info_err!(format!("Invalid mapper assignment field {value}"))); } } @@ -237,7 +237,7 @@ impl Mapper { for t in transforms { let field = t.field.as_str(); if !valid_property!(field, AFFIX_FIELDS) { - return Err(M3uFilterError::new(M3uFilterErrorKind::Info, format!("Invalid mapper transform field {field}"))); + return Err(info_err!(format!("Invalid mapper transform field {field}"))); } t.prepare(templates)?; } @@ -452,7 +452,7 @@ impl Mapping { let mut counters = vec![]; for def in counter_def_list { if !valid_property!(def.field.as_str(), COUNTER_FIELDS) { - return Err(M3uFilterError::new(M3uFilterErrorKind::Info, format!("Invalid counter field {}", def.field))); + return Err(info_err!(format!("Invalid counter field {}", def.field))); } match get_filter(&def.filter, templates) { Ok(flt) => { @@ -464,7 +464,7 @@ impl Mapping { value: Arc::new(AtomicU32::new(def.value)), }); } - Err(e) => return Err(M3uFilterError::new(M3uFilterErrorKind::Info, e.to_string())) + Err(e) => return Err(info_err!(e.to_string())) } } self.t_counter = Some(counters); diff --git a/src/model/xtream.rs b/src/model/xtream.rs index d0bf4accd..8eeae1a74 100644 --- a/src/model/xtream.rs +++ b/src/model/xtream.rs @@ -349,7 +349,7 @@ impl XtreamSeriesInfoEpisode { let mut result = Map::new(); let bdpath = &series_info.info.backdrop_path; if !bdpath.is_empty() { - result.insert(String::from("backdrop_path"), Value::Array(Vec::from([Value::String(String::from(bdpath.first().unwrap()))]))); + result.insert(String::from("backdrop_path"), Value::Array(Vec::from([Value::String(String::from(bdpath.first()?))]))); } add_str_property_if_exists!(result, self.added.as_str(), "added"); add_str_property_if_exists!(result, series_info.info.cast.as_str(), "cast"); diff --git a/src/processing/playlist_processor.rs b/src/processing/playlist_processor.rs index 5abd5b427..fb66930e7 100644 --- a/src/processing/playlist_processor.rs +++ b/src/processing/playlist_processor.rs @@ -1,6 +1,5 @@ extern crate unidecode; -use crate::debug_if_enabled; use async_std::sync::Mutex; use core::cmp::Ordering; use std::cell::RefCell; @@ -32,7 +31,7 @@ use crate::repository::playlist_repository::persist_playlist; use crate::utils::default_utils::default_as_default; use crate::utils::download; use crate::utils::request_utils::mask_sensitive_info; -use crate::{get_errors_notify_message, model::config, Config}; +use crate::{debug_if_enabled, get_errors_notify_message, model::config, Config, notify_err}; fn is_valid(pli: &PlaylistItem, target: &ConfigTarget) -> bool { let provider = ValueProvider { pli: RefCell::new(pli) }; @@ -319,7 +318,7 @@ async fn process_source(cfg: Arc, source_idx: usize, user_targets: Arc

, source_idx: usize, user_targets: Arc

()); for target in &source.targets { diff --git a/src/processing/xtream_processor.rs b/src/processing/xtream_processor.rs index 18d46ef0b..77f8e722f 100644 --- a/src/processing/xtream_processor.rs +++ b/src/processing/xtream_processor.rs @@ -9,6 +9,7 @@ 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"; @@ -76,7 +77,7 @@ pub(in crate::processing) async fn playlist_resolve_process_playlist_item(pli: & result = match download::get_xtream_stream_info_content(&info_url, input).await { Ok(content) => Some(content), Err(err) => { - errors.push(M3uFilterError::new(M3uFilterErrorKind::Info, format!("{err}"))); + errors.push(info_err!(format!("{err}"))); None } }; @@ -159,7 +160,7 @@ where Ok(file_path) => file_path, Err(err) => { let fpl_name = fpl.input.name.as_ref().map_or("?", String::as_str); - errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not create storage path for input {fpl_name}: {err}"))); + errors.push(notify_err!(format!("Could not create storage path for input {fpl_name}: {err}"))); return processed_info_ids; } }; @@ -173,7 +174,7 @@ where } drop(file_lock); } - Err(err) => errors.push(M3uFilterError::new(M3uFilterErrorKind::Info, format!("{err}"))), + Err(err) => errors.push(info_err!(format!("{err}"))), } processed_info_ids } diff --git a/src/processing/xtream_processor_series.rs b/src/processing/xtream_processor_series.rs index 6d2180b21..56ac9632b 100644 --- a/src/processing/xtream_processor_series.rs +++ b/src/processing/xtream_processor_series.rs @@ -4,14 +4,13 @@ use crate::model::playlist::{FetchedPlaylist, PlaylistGroup, PlaylistItem, Playl 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::{create_resolve_options_function_for_xtream_target, handle_error, handle_error_and_return}; +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 log::error; -use crate::repository::bplustree::BPlusTree; -use crate::repository::IndexedDocumentIndex; +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"; @@ -73,25 +72,108 @@ async fn playlist_resolve_series_info(cfg: &Config, errors: &mut Vec, + errors: &mut Vec, + processed_ids: &HashMap, +) -> Vec { + let mut result: Vec = vec![]; + let input = provider_fpl.input; + + let (info_path, idx_path) = match 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; + } + }; + + 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 mut doc_reader = match IndexedDocumentReader::::new(&info_path, &idx_path) { + Ok(reader) => reader, + Err(_) => return result, + }; + + for plg in provider_fpl + .playlistgroups + .iter_mut() + .filter(|plg| plg.xtream_cluster == XtreamCluster::Series) + { + let mut group_series = vec![]; + + for pli in plg + .channels + .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 + } + } + Err(err) => errors.push(info_err!(format!("Failed to parse JSON: {err}"))), + } + } + } + } + } + + if !group_series.is_empty() { + result.push(PlaylistGroup { + id: plg.id, + title: plg.title.clone(), + channels: group_series, + xtream_cluster: XtreamCluster::Series, + }); + } + } + + result +} + pub async fn playlist_resolve_series(cfg: &Config, target: &ConfigTarget, errors: &mut Vec, @@ -104,38 +186,22 @@ 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; } - - if let Ok(Some((info_path, idx_path))) = get_input_storage_path(provider_fpl.input, &cfg.working_dir).map(|storage_path| xtream_get_info_file_paths(&storage_path, XtreamCluster::Series)) { - match cfg.file_locks.read_lock(&info_path).await { - Ok(_file_lock) => { - let index = IndexedDocumentIndex::::load(&idx_path).unwrap_or_else(|err| { - error!("Failed to load index {idx_path:?}: {err}"); - IndexedDocumentIndex::::new() - }); - let mut result: Vec = vec![]; - for plg in &mut provider_fpl.playlistgroups.iter() - .filter(|&plg| plg.xtream_cluster == XtreamCluster::Series) - { - let mut group_series: Vec = vec![]; - for pli in plg.channels.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 Some(offset) = index.query(&provider_id) { - IndexedDocumentReader:: - } - } - } - } - } - }, - Err(err) => errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, "Could not lock input info file for series".to_string())) - } - } else { - errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, "Failed to open input info file for series".to_string())); + 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); + } + // run processing pipe over new items + for f in pipe { + let r = f(&mut series_playlist, target); + if let Some(v) = r { + series_playlist = v; + } + } + // assign new items to the new playlist + for plg in &series_playlist { + processed_fpl.update_playlist(plg); } - - - // Now update provider_fpl } diff --git a/src/processing/xtream_processor_vod.rs b/src/processing/xtream_processor_vod.rs index bd2b163af..a759085be 100644 --- a/src/processing/xtream_processor_vod.rs +++ b/src/processing/xtream_processor_vod.rs @@ -3,7 +3,7 @@ use crate::model::config::{Config, ConfigTarget, InputType}; 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, 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}; +use crate::{create_resolve_options_function_for_xtream_target, handle_error, handle_error_and_return, notify_err}; use serde_json::{Map, Value}; use std::collections::HashMap; use std::fs::File; @@ -84,21 +84,26 @@ pub async fn playlist_resolve_vod(cfg: &Config, target: &ConfigTarget, errors: & if let Some(content) = playlist_resolve_process_playlist_item(pli, fpl.input, errors, resolve_delay, XtreamCluster::Video).await { 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| errors.push(M3uFilterError::new( M3uFilterErrorKind::Notify, format!("Failed to resolve vod, could not write to content wal file {err}")))); + handle_error_and_return!(write_info_content_to_wal_file(&mut content_writer, provider_id, &content), + |err| errors.push(notify_err!(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| errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve vod wal, could not write to record wal file {err}")))); + handle_error_and_return!(write_vod_info_record_to_wal_file(&mut record_writer, provider_id, &info_record), + |err| errors.push(notify_err!(format!("Failed to resolve vod wal, could not write to record wal file {err}")))); content_updated = true; } } } } if content_updated { - handle_error!(content_writer.flush(), |err| errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve vod, could not write to wal file {err}")))); + handle_error!(content_writer.flush(), + |err| errors.push(notify_err!(format!("Failed to resolve vod, could not write to wal file {err}")))); drop(content_writer); - handle_error!(record_writer.flush(), |err| errors.push(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to resolve vod tmdb, could not write to wal file {err}")))); + handle_error!(record_writer.flush(), + |err| errors.push(notify_err!(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_content_file, &wal_content_path, XtreamCluster::Video).await, |err| errors.push(err)); - handle_error!(xtream_update_input_vod_record_from_wal_file(cfg, fpl.input, &mut wal_record_file, &wal_record_path).await, |err| errors.push(err)); + handle_error!(xtream_update_input_info_file(cfg, fpl.input, &mut wal_content_file, &wal_content_path, XtreamCluster::Video).await, + |err| errors.push(err)); + handle_error!(xtream_update_input_vod_record_from_wal_file(cfg, fpl.input, &mut wal_record_file, &wal_record_path).await, + |err| errors.push(err)); } } diff --git a/src/repository/bplustree.rs b/src/repository/bplustree.rs index 0bbbbc8fe..a93963118 100644 --- a/src/repository/bplustree.rs +++ b/src/repository/bplustree.rs @@ -482,8 +482,10 @@ where // } // } +/// /// `BPlusTreeQuery` can be used to query the `BPlusTree` on-disk. /// If you intend to do frequent queries then use `BPlusTree` instead which loads the tree into memory. +/// pub struct BPlusTreeQuery { file: BufReader, _marker_k: PhantomData, diff --git a/src/repository/epg_repository.rs b/src/repository/epg_repository.rs index 4c21f5cf8..92cfeea92 100644 --- a/src/repository/epg_repository.rs +++ b/src/repository/epg_repository.rs @@ -3,7 +3,7 @@ use std::io::{Cursor, Write}; use std::path::{Path}; use log::{Level}; use quick_xml::{Writer}; -use crate::debug_if_enabled; +use crate::{debug_if_enabled, notify_err}; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::{Config, ConfigTarget, TargetOutput}; use crate::model::config::TargetType; @@ -20,23 +20,19 @@ fn epg_write_file(target: &ConfigTarget, epg: &Epg, path: &Path) -> Result<(), M Ok(mut epg_file) => { match epg_file.write_all("".as_bytes()) { Ok(()) => {} - Err(err) => return Err(M3uFilterError::new( - M3uFilterErrorKind::Notify, format!("failed to write epg: {} - {}", path.to_str().unwrap_or("?"), err))), + Err(err) => return Err(notify_err!(format!("failed to write epg: {} - {}", path.to_str().unwrap_or("?"), err))), } match epg_file.write_all(&result) { Ok(()) => { debug_if_enabled!("Epg for target {} written to {}", target.name, path.to_str().unwrap_or("?")); } - Err(err) => return Err(M3uFilterError::new( - M3uFilterErrorKind::Notify, format!("failed to write epg: {} - {}", path.to_str().unwrap_or("?"), err))), + Err(err) => return Err(notify_err!(format!("failed to write epg: {} - {}", path.to_str().unwrap_or("?"), err))), } } - Err(err) => return Err(M3uFilterError::new( - M3uFilterErrorKind::Notify, format!("failed to write epg: {} - {}", path.to_str().unwrap_or("?"), err))), + Err(err) => return Err(notify_err!(format!("failed to write epg: {} - {}", path.to_str().unwrap_or("?"), err))), } } - Err(err) => return Err(M3uFilterError::new( - M3uFilterErrorKind::Notify, format!("failed to write epg: {} - {}", path.to_str().unwrap_or("?"), err))), + Err(err) => return Err(notify_err!(format!("failed to write epg: {} - {}", path.to_str().unwrap_or("?"), err))), } Ok(()) } @@ -46,9 +42,7 @@ pub fn epg_write(target: &ConfigTarget, cfg: &Config, target_path: &Path, epg: O match &output.target { TargetType::M3u => { if output.filename.is_none() { - return Err(M3uFilterError::new( - M3uFilterErrorKind::Notify, - format!("write epg for target {} failed: No filename set", target.name))); + return Err(notify_err!(format!("write epg for target {} failed: No filename set", target.name))); } let path = m3u_get_epg_file_path(target_path); debug_if_enabled!("writing m3u epg to {}", path.to_str().unwrap_or("?")); @@ -61,9 +55,7 @@ pub fn epg_write(target: &ConfigTarget, cfg: &Config, target_path: &Path, epg: O debug_if_enabled!("writing xtream epg to {}", epg_path.to_str().unwrap_or("?")); epg_write_file(target, epg_data, &epg_path)?; } - None => return Err(M3uFilterError::new( - M3uFilterErrorKind::Notify, - format!("failed to serialize epg for target: {}, storage path not found", target.name))), + None => return Err(notify_err!(format!("failed to serialize epg for target: {}, storage path not found", target.name))), } } TargetType::Strm => {} diff --git a/src/repository/indexed_document.rs b/src/repository/indexed_document.rs index 21b7fde4f..4e0a101e4 100644 --- a/src/repository/indexed_document.rs +++ b/src/repository/indexed_document.rs @@ -4,7 +4,7 @@ use std::io::{BufReader, Error, ErrorKind, Read, Seek, SeekFrom, Write}; use std::marker::PhantomData; use std::path::{Path, PathBuf}; -use crate::repository::bplustree::{BPlusTree, BPlusTreeQuery, BPlusTreeUpdate}; +use crate::repository::bplustree::{BPlusTree, BPlusTreeQuery}; use crate::utils::file_utils; use log::error; use serde::{Deserialize, Serialize}; @@ -15,6 +15,8 @@ const LEN_SIZE: usize = 4; pub(in crate::repository) type OffsetPointer = u32; type SizeType = u32; +pub type IndexedDocumentIndex = BPlusTree; + pub(in crate::repository) struct IndexedDocument {} impl IndexedDocument { @@ -54,12 +56,10 @@ impl IndexedDocument { } //////////////////////////////////////////////////////// -// -// IndexedDocumentWriter -// +/// +/// IndexedDocumentWriter +/// //////////////////////////////////////////////////////// - - /** * Creates two files, * - content @@ -78,7 +78,7 @@ where index_path: PathBuf, main_file: File, main_offset: OffsetPointer, - index_tree: BPlusTree, + index_tree: IndexedDocumentIndex, dirty: bool, fragmented: bool, } @@ -120,12 +120,12 @@ where // Initialize the index tree (BPlusTree) - either by deserializing an existing one or creating a new one let index_tree = if append_mode && index_path.exists() { - BPlusTree::::load(&index_path).unwrap_or_else(|err| { + IndexedDocumentIndex::::load(&index_path).unwrap_or_else(|err| { error!("Failed to load index {:?}: {}", index_path, err); - BPlusTree::::new() + IndexedDocumentIndex::::new() }) } else { - BPlusTree::::new() + IndexedDocumentIndex::::new() }; Ok(Self { @@ -229,13 +229,69 @@ where } } + +//////////////////////////////////////////////////////// +/// +/// IndexedDocumentReader +/// +//////////////////////////////////////////////////////// +pub struct IndexedDocumentReader +where + T: serde::de::DeserializeOwned, + K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug, +{ + main_file: BufReader, + index_tree: IndexedDocumentIndex, + t_type: PhantomData, +} + + +impl IndexedDocumentReader +where + T: serde::de::DeserializeOwned, + K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug, +{ + pub fn new(main_path: &Path, index_path: &Path) -> Result { + if main_path.exists() && index_path.exists() { + let main_file = OpenOptions::new() + .read(true) + .write(false) + .truncate(false) + .open(&main_path)?; + let index_tree = IndexedDocumentIndex::::load(&index_path)?; + + Ok(Self { + main_file: BufReader::new(main_file), + index_tree, + t_type: PhantomData, + }) + } else { + Err(Error::new(ErrorKind::NotFound, format!("File not found {}", main_path.to_str().unwrap()))) + } + } + pub fn get(&mut self, doc_id: &K) -> Result { + if let Some(offset) = self.index_tree.query(doc_id) { + self.main_file.seek(SeekFrom::Start(u64::from(*offset)))?; + let buf_size = IndexedDocument::read_content_size(&mut self.main_file)?; + let mut buffer: Vec = vec![0; buf_size]; + self.main_file.read_exact(&mut buffer)?; + if let Ok(item) = bincode::deserialize::(&buffer) { + return Ok(item); + } + } + Err(Error::new(ErrorKind::NotFound, format!("Entry not found {doc_id:?}"))) + } +} + + //////////////////////////////////////////////////////// // -// IndexedDocumentReader -// +/// IndexedDocumentIterator +/// +/// Iterator | Sequential access with has_next / next //////////////////////////////////////////////////////// -pub(in crate::repository) struct IndexedDocumentReader { +pub(in crate::repository) struct IndexedDocumentIterator { main_path: PathBuf, main_file: BufReader, offsets: Vec, @@ -246,18 +302,17 @@ pub(in crate::repository) struct IndexedDocumentReader { k_type: PhantomData, } -impl IndexedDocumentReader +impl IndexedDocumentIterator where + T: serde::de::DeserializeOwned, K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug, { pub fn new(main_path: &Path, index_path: &Path) -> Result { if main_path.exists() && index_path.exists() { let mut offsets = Vec::::new(); { - let index_tree = BPlusTree::::load(index_path)?; - index_tree.traverse(|_, values| { - offsets.extend(values); - }); + let index_tree = IndexedDocumentIndex::::load(index_path)?; + index_tree.traverse(|_, values| offsets.extend(values)); offsets.sort_unstable(); } match File::open(main_path) { @@ -276,12 +331,10 @@ where Err(e) => Err(e) } } else { - Err(Error::new(ErrorKind::NotFound, format!("File not found {}", - main_path.to_str().unwrap()))) + Err(Error::new(ErrorKind::NotFound, format!("File not found {}", main_path.to_str().unwrap()))) } } - pub fn get_path(&self) -> &Path { &self.main_path } @@ -317,37 +370,9 @@ where } } } - - pub(in crate::repository) fn read_indexed_item(main_path: &Path, index_path: &Path, doc_id: &K) -> Result - { - if main_path.exists() && index_path.exists() { - // get the offset from index - let offset = IndexedDocument::get_offset(index_path, doc_id)?; - let mut main_file = File::open(main_path)?; - main_file.seek(SeekFrom::Start(offset))?; - let buf_size = IndexedDocument::read_content_size(&mut main_file)?; - let mut buffer: Vec = vec![0; buf_size]; - main_file.read_exact(&mut buffer)?; - if let Ok(item) = bincode::deserialize::(&buffer) { - return Ok(item); - } - } - Err(Error::new(ErrorKind::Other, format!("Failed to read item for id {:?} - {}", doc_id, main_path.to_str().unwrap()))) - } - - pub(in crate::repository) fn read_indexed_item_at_offset(main_file: &mut File, offset: u64, doc_id: &K) -> Result { - main_file.seek(SeekFrom::Start(offset))?; - let buf_size = IndexedDocument::read_content_size(main_file)?; - let mut buffer: Vec = vec![0; buf_size]; - main_file.read_exact(&mut buffer)?; - if let Ok(item) = bincode::deserialize::(&buffer) { - return Ok(item); - } - Err(Error::new(ErrorKind::Other, format!("Failed to read item for id {:?}", doc_id))) - } } -impl Iterator for IndexedDocumentReader +impl Iterator for IndexedDocumentIterator where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug, { @@ -365,16 +390,46 @@ where } //////////////////////////////////////////////////////// -// -// IndexedDocumentGarbageCollector -// +/// +/// IndexedDocumentDirectAccess +/// +//////////////////////////////////////////////////////// +pub(in crate::repository) struct IndexedDocumentDirectAccess {} + +impl IndexedDocumentDirectAccess { + pub(in crate::repository) fn read_indexed_item(main_path: &Path, index_path: &Path, doc_id: &K) -> Result + where + K: Ord + Serialize + for<'de> Deserialize<'de> + Clone + Debug, + T: serde::de::DeserializeOwned, + { + if main_path.exists() && index_path.exists() { + // get the offset from index + let offset = IndexedDocument::get_offset(index_path, doc_id)?; + let mut main_file = File::open(main_path)?; + main_file.seek(SeekFrom::Start(offset))?; + let buf_size = IndexedDocument::read_content_size(&mut main_file)?; + let mut buffer: Vec = vec![0; buf_size]; + main_file.read_exact(&mut buffer)?; + if let Ok(item) = bincode::deserialize::(&buffer) { + return Ok(item); + } + } + Err(Error::new(ErrorKind::Other, format!("Failed to read item for id {:?} - {}", doc_id, main_path.to_str().unwrap()))) + } +} + + +//////////////////////////////////////////////////////// +/// +/// IndexedDocumentGarbageCollector +/// //////////////////////////////////////////////////////// pub(in crate::repository) struct IndexedDocumentGarbageCollector { main_path: PathBuf, index_path: PathBuf, main_file: File, - index_tree: BPlusTree, + index_tree: IndexedDocumentIndex, } impl IndexedDocumentGarbageCollector @@ -400,7 +455,7 @@ where } // Initialize the index tree (BPlusTree) - by deserializing an existing one - let index_tree = BPlusTree::::load(&index_path)?; + let index_tree = IndexedDocumentIndex::::load(&index_path)?; Ok(Self { main_path, @@ -480,7 +535,7 @@ mod tests { use serde::{Deserialize, Serialize}; - use crate::repository::indexed_document::{IndexedDocumentGarbageCollector, IndexedDocumentReader, IndexedDocumentWriter}; + use crate::repository::indexed_document::{IndexedDocumentGarbageCollector, IndexedDocumentIterator, IndexedDocumentWriter}; // Example usage with a simple struct #[derive(Serialize, Deserialize, Debug, Clone, PartialEq, Eq)] @@ -539,7 +594,7 @@ mod tests { assert!(size_main_file_5 < size_main_file_4, "Failed, the filesize should be less"); } { - let reader = IndexedDocumentReader::::new(&main_path, &index_path)?; + let reader = IndexedDocumentIterator::::new(&main_path, &index_path)?; let mut i = 0; for doc in reader { assert_eq!(doc.id, i, "Wrong id"); @@ -552,7 +607,3 @@ mod tests { Ok(()) } } - -pub type IndexedDocumentQuery = BPlusTreeQuery; -pub type IndexedDocumentUpdate = BPlusTreeUpdate; -pub type IndexedDocumentIndex = BPlusTree; \ No newline at end of file diff --git a/src/repository/kodi_repository.rs b/src/repository/kodi_repository.rs index 4abf2f94e..c746a92fa 100644 --- a/src/repository/kodi_repository.rs +++ b/src/repository/kodi_repository.rs @@ -1,4 +1,4 @@ -use crate::create_m3u_filter_error_result; +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}; @@ -129,7 +129,7 @@ async fn get_tmdb_id(cfg: &Config, provider_id: Option, input_id: u16, pub async fn kodi_write_strm_playlist(target: &ConfigTarget, cfg: &Config, new_playlist: &[PlaylistGroup], filename: Option<&str>) -> Result<(), M3uFilterError> { if !new_playlist.is_empty() { if filename.is_none() { - return Err(M3uFilterError::new(M3uFilterErrorKind::Notify, "write strm playlist failed: ".to_string())); + return Err(notify_err!("write strm playlist failed: ".to_string())); } let underscore_whitespace = target.options.as_ref().is_some_and(|o| o.underscore_whitespace); let cleanup = target.options.as_ref().is_some_and(|o| o.cleanup); diff --git a/src/repository/m3u_playlist_iterator.rs b/src/repository/m3u_playlist_iterator.rs index 3c7e0a480..51c9b8d89 100644 --- a/src/repository/m3u_playlist_iterator.rs +++ b/src/repository/m3u_playlist_iterator.rs @@ -1,9 +1,10 @@ use crate::api::api_utils::get_user_server_info; +use crate::info_err; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::api_proxy::{ProxyType, ProxyUserCredentials}; use crate::model::config::{Config, ConfigTarget, ConfigTargetOptions}; use crate::model::playlist::{M3uPlaylistItem, PlaylistItemType}; -use crate::repository::indexed_document::IndexedDocumentReader; +use crate::repository::indexed_document::IndexedDocumentIterator; use crate::repository::m3u_repository::m3u_get_file_paths; use crate::repository::storage::ensure_target_storage_path; use crate::utils::file_lock_manager::FileReadGuard; @@ -11,7 +12,7 @@ use crate::utils::file_lock_manager::FileReadGuard; pub const M3U_STREAM_PATH: &str = "m3u-stream"; pub struct M3uPlaylistIterator { - reader: IndexedDocumentReader, + reader: IndexedDocumentIterator, base_url: String, username: String, password: String, @@ -32,15 +33,12 @@ impl M3uPlaylistIterator { let target_path = ensure_target_storage_path(cfg, target.name.as_str())?; let (m3u_path, idx_path) = m3u_get_file_paths(&target_path); - let file_lock = cfg.file_locks.read_lock(&m3u_path).await.map_err(|err| { - M3uFilterError::new( - M3uFilterErrorKind::Info, - format!("Could not lock document {m3u_path:?}: {err}"), - ) - })?; + let file_lock = cfg.file_locks.read_lock(&m3u_path).await + .map_err(|err| info_err!(format!("Could not lock document {m3u_path:?}: {err}")))?; let reader = - IndexedDocumentReader::::new(&m3u_path, &idx_path).map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Info,format!("Could not deserialize file {m3u_path:?} - {err}")))?; + IndexedDocumentIterator::::new(&m3u_path, &idx_path) + .map_err(|err| info_err!(format!("Could not deserialize file {m3u_path:?} - {err}")))?; let target_options = target.options.as_ref(); let include_type_in_url = target_options.is_some_and( |opts| opts.m3u_include_type_in_url); @@ -101,15 +99,9 @@ impl Iterator for M3uPlaylistIterator { let stream_url = match m3u_pli.item_type { PlaylistItemType::LiveHls => None, _ => match &self.proxy_type { - ProxyType::Reverse => Some(self.get_stream_url( - &m3u_pli, - self.include_type_in_url, - )), + ProxyType::Reverse => Some(self.get_stream_url(&m3u_pli, self.include_type_in_url)), ProxyType::Redirect => if self.mask_redirect_url { - Some(self.get_stream_url( - &m3u_pli, - self.include_type_in_url, - )) + Some(self.get_stream_url(&m3u_pli,self.include_type_in_url)) } else { None } diff --git a/src/repository/m3u_repository.rs b/src/repository/m3u_repository.rs index 7e0af03df..01c927a77 100644 --- a/src/repository/m3u_repository.rs +++ b/src/repository/m3u_repository.rs @@ -3,12 +3,12 @@ use std::io::{BufWriter, Error, ErrorKind, Write}; use std::path::{Path, PathBuf}; use log::error; -use crate::create_m3u_filter_error; +use crate::{create_m3u_filter_error, info_err}; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::api_proxy::{ProxyUserCredentials}; use crate::model::config::{Config, ConfigTarget}; use crate::model::playlist::{M3uPlaylistItem, PlaylistGroup, PlaylistItem, PlaylistItemType}; -use crate::repository::indexed_document::{IndexedDocumentReader, IndexedDocumentWriter}; +use crate::repository::indexed_document::{IndexedDocumentDirectAccess, IndexedDocumentWriter}; use crate::repository::m3u_playlist_iterator::M3uPlaylistIterator; use crate::repository::storage::{FILE_SUFFIX_DB, FILE_SUFFIX_INDEX}; use crate::utils::file_utils; @@ -61,7 +61,7 @@ pub async fn m3u_write_playlist(target: &ConfigTarget, cfg: &Config, target_path persist_m3u_playlist_as_text(target, cfg, &m3u_playlist); { - let _file_lock = cfg.file_locks.write_lock(&m3u_path).await.map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Info, format!("{err}")))?; + let _file_lock = cfg.file_locks.write_lock(&m3u_path).await.map_err(|err| info_err!(format!("{err}")))?; match IndexedDocumentWriter::new(m3u_path.clone(), idx_path) { Ok(mut writer) => { for m3u in m3u_playlist { @@ -94,6 +94,6 @@ pub async fn m3u_get_item_for_stream_id(cfg: &Config, stream_id: u32, m3u_path: } { let _file_lock = cfg.file_locks.read_lock(m3u_path).await?; - IndexedDocumentReader::::read_indexed_item(m3u_path, idx_path, &stream_id) + IndexedDocumentDirectAccess::read_indexed_item::(m3u_path, idx_path, &stream_id) } } \ No newline at end of file diff --git a/src/repository/mod.rs b/src/repository/mod.rs index d78c54e5d..b678bdc99 100644 --- a/src/repository/mod.rs +++ b/src/repository/mod.rs @@ -2,7 +2,7 @@ pub mod storage; pub mod target_id_mapping; pub mod bplustree; mod indexed_document; -pub use indexed_document::IndexedDocumentIndex; +pub use indexed_document::IndexedDocumentReader; pub mod playlist_repository; pub mod m3u_repository; pub mod xtream_repository; diff --git a/src/repository/playlist_repository.rs b/src/repository/playlist_repository.rs index ccf6b7679..c0a720096 100644 --- a/src/repository/playlist_repository.rs +++ b/src/repository/playlist_repository.rs @@ -1,3 +1,4 @@ +use crate::info_err; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::{Config, ConfigTarget, TargetType}; use crate::model::playlist::PlaylistItemType::LiveUnknown; @@ -23,7 +24,7 @@ pub async fn persist_playlist(playlist: &mut [PlaylistGroup], epg: Option<&Epg>, let _file_lock = match cfg.file_locks.write_lock(&target_id_mapping_file).await { Ok(lock) => lock, Err(err) => { - errors.push(M3uFilterError::new(M3uFilterErrorKind::Info, err.to_string())); + errors.push(info_err!(err.to_string())); return Err(errors); } }; @@ -55,7 +56,7 @@ pub async fn persist_playlist(playlist: &mut [PlaylistGroup], epg: Option<&Epg>, errors.push(err); } else { if let Err(err) = target_id_mapping.persist() { - errors.push(M3uFilterError::new(M3uFilterErrorKind::Info, err.to_string())); + errors.push(info_err!(err.to_string())); } if !playlist.is_empty() { if let Err(err) = epg_write(target, cfg, &target_path, epg, output) { @@ -66,7 +67,7 @@ pub async fn persist_playlist(playlist: &mut [PlaylistGroup], epg: Option<&Epg>, } if let Err(err) = target_id_mapping.persist() { - errors.push(M3uFilterError::new(M3uFilterErrorKind::Info, err.to_string())); + errors.push(info_err!(err.to_string())); } if errors.is_empty() { Ok(()) } else { Err(errors) } diff --git a/src/repository/storage.rs b/src/repository/storage.rs index 0285e43d2..5c20a033f 100644 --- a/src/repository/storage.rs +++ b/src/repository/storage.rs @@ -3,6 +3,7 @@ use std::fmt::Write; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::{Config, ConfigInput}; use crate::model::playlist::UUIDType; +use crate::notify_err; use crate::utils::file_utils; pub(in crate::repository) const FILE_SUFFIX_DB: &str = "db"; @@ -36,12 +37,12 @@ pub fn ensure_target_storage_path(cfg: &Config, target_name: &str) -> Result, + reader: IndexedDocumentIterator, options: XtreamMappingOptions, category_id: u32, _file_lock: FileReadGuard, @@ -23,12 +24,11 @@ impl XtreamPlaylistIterator { ) -> Result { if let Some(storage_path) = xtream_get_storage_path(config, target.name.as_str()) { let (xtream_path, idx_path) = xtream_get_file_paths(&storage_path, cluster); - let file_lock = config.file_locks.read_lock(&xtream_path).await.map_err(|err| - M3uFilterError::new(M3uFilterErrorKind::Info, format!("Could not lock document {xtream_path:?}: {err}")) - )?; + let file_lock = config.file_locks.read_lock(&xtream_path).await + .map_err(|err| info_err!(format!("Could not lock document {xtream_path:?}: {err}")))?; - let reader = IndexedDocumentReader::::new(&xtream_path, &idx_path) - .map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Info, format!("Could not deserialize file {} - {}", &xtream_path.to_str().unwrap(), err)))?; + let reader = IndexedDocumentIterator::::new(&xtream_path, &idx_path) + .map_err(|err| info_err!(format!("Could not deserialize file {} - {}", &xtream_path.to_str().unwrap(), err)))?; let options = XtreamMappingOptions::from_target_options(target.options.as_ref()); @@ -39,7 +39,7 @@ impl XtreamPlaylistIterator { _file_lock: file_lock, }) } else { - Err(M3uFilterError::new(M3uFilterErrorKind::Info, format!("Failed to find xtream storage for target {}", &target.name))) + Err(info_err!(format!("Failed to find xtream storage for target {}", &target.name))) } } } diff --git a/src/repository/xtream_repository.rs b/src/repository/xtream_repository.rs index 7d2a1de4b..d84dc903b 100644 --- a/src/repository/xtream_repository.rs +++ b/src/repository/xtream_repository.rs @@ -12,13 +12,13 @@ use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::{Config, ConfigInput, ConfigTarget}; use crate::model::playlist::{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}; -use crate::repository::storage::{get_input_storage_path, get_target_id_mapping_file, get_target_storage_path, hash_string,FILE_SUFFIX_DB, FILE_SUFFIX_INDEX}; +use crate::repository::bplustree::{BPlusTree, BPlusTreeQuery, BPlusTreeUpdate}; +use crate::repository::indexed_document::{IndexedDocumentGarbageCollector, IndexedDocumentWriter, IndexedDocumentDirectAccess}; +use crate::repository::storage::{get_input_storage_path, get_target_id_mapping_file, get_target_storage_path, hash_string, FILE_SUFFIX_DB, FILE_SUFFIX_INDEX}; use crate::repository::target_id_mapping::{TargetIdMapping, VirtualIdRecord}; use crate::repository::xtream_playlist_iterator::XtreamPlaylistIterator; use crate::utils::json_utils::{json_iter_array, json_write_documents_to_file}; -use crate::{create_m3u_filter_error, create_m3u_filter_error_result}; +use crate::{create_m3u_filter_error, create_m3u_filter_error_result, notify_err, info_err}; pub static COL_CAT_LIVE: &str = "cat_live"; pub static COL_CAT_SERIES: &str = "cat_series"; @@ -75,12 +75,12 @@ fn ensure_xtream_storage_path(cfg: &Config, target_name: &str) -> Result { for item in playlist { @@ -351,11 +351,7 @@ async fn xtream_read_item_for_stream_id( let (xtream_path, idx_path) = xtream_get_file_paths(storage_path, cluster); { let _file_lock = cfg.file_locks.read_lock(&xtream_path).await?; - IndexedDocumentReader::::read_indexed_item( - &xtream_path, - &idx_path, - &stream_id, - ) + IndexedDocumentDirectAccess::read_indexed_item::(&xtream_path, &idx_path, &stream_id) } } @@ -367,7 +363,7 @@ async fn xtream_read_series_item_for_stream_id( let (xtream_path, idx_path) = xtream_get_file_paths_for_series(storage_path); { let _file_lock = cfg.file_locks.read_lock(&xtream_path).await?; - IndexedDocumentReader::::read_indexed_item(&xtream_path, &idx_path, &stream_id) + IndexedDocumentDirectAccess::read_indexed_item::(&xtream_path, &idx_path, &stream_id) } } @@ -391,7 +387,7 @@ pub async fn xtream_get_item_for_stream_id( let target_id_mapping_file = get_target_id_mapping_file(&target_path); let _file_lock = config.file_locks.read_lock(&target_id_mapping_file).await.map_err(|err| Error::new(ErrorKind::Other, format!("Could not get lock for id mapping for target {} err:{err}", target.name)))?; - let mut target_id_mapping = IndexedDocumentQuery::::try_new(&target_id_mapping_file).map_err(|err| Error::new(ErrorKind::Other, format!("Could not load id mapping for target {} err:{err}", target.name)))?; + let mut target_id_mapping = BPlusTreeQuery::::try_new(&target_id_mapping_file).map_err(|err| Error::new(ErrorKind::Other, format!("Could not load id mapping for target {} err:{err}", target.name)))?; let mapping = target_id_mapping.query(&virtual_id).ok_or_else(|| Error::new(ErrorKind::Other, format!("Could not find mapping for target {} and id {}", target.name, virtual_id)))?; match mapping.item_type { PlaylistItemType::SeriesInfo => { @@ -447,7 +443,7 @@ 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 Ok(mut target_id_mapping) = BPlusTreeUpdate::::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); @@ -491,7 +487,7 @@ async fn xtream_get_info_mapping(config: &Config, target_name: &str, info_id: u3 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()?; - IndexedDocumentQuery::::try_new(&target_id_mapping_file).map_err(|err| { + BPlusTreeQuery::::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().map(|mut tree| tree.query(&info_id))? @@ -515,7 +511,7 @@ pub async fn xtream_load_series_info( error!("Could not lock document {:?}: {}", info_path, err); Error::new(ErrorKind::Other, format!("Document Reader error for target {target_name}")) }).ok()?; - return match IndexedDocumentReader::::read_indexed_item(&info_path, &idx_path, &series_id) { + return match IndexedDocumentDirectAccess::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); @@ -556,9 +552,7 @@ pub async fn xtream_load_vod_info( 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, - ) { + return match IndexedDocumentDirectAccess::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 @@ -669,11 +663,7 @@ pub async fn xtream_get_input_info( 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, cluster)) { 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, - ) { + if let Ok(content) = IndexedDocumentDirectAccess::read_indexed_item::(&info_path, &idx_path, &provider_id) { return Some(content); } } @@ -692,7 +682,7 @@ pub async fn xtream_update_input_info_file( 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}")))?; + wal_file.seek(SeekFrom::Start(0)).map_err(|err| notify_err!(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) => { @@ -703,30 +693,29 @@ pub async fn xtream_update_input_info_file( break; // End of file } let provider_id = u32::from_le_bytes(provider_id_bytes); - reader.read_exact(&mut length_bytes).map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not read temporary {cluster} info {err}")))?; + reader.read_exact(&mut length_bytes).map_err(|err| notify_err!(format!("Could not read temporary {cluster} info {err}")))?; let length = u32::from_le_bytes(length_bytes) as usize; let mut buffer = vec![0u8; length]; - reader.read_exact(&mut buffer).map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not read temporary {cluster} info {err}")))?; + reader.read_exact(&mut buffer).map_err(|err| notify_err!(format!("Could not read temporary {cluster} info {err}")))?; if let Ok(content) = String::from_utf8(buffer) { let _ = writer.write_doc(provider_id, &content); } } - writer.store().map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not store {cluster} info {err}")))?; + writer.store().map_err(|err| notify_err!(format!("Could not store {cluster} info {err}")))?; drop(reader); if let Err(err) = fs::remove_file(wal_path) { error!("Failed to delete WAL file for {cluster} {err}"); } Ok(()) } - Err(err) => Err(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not create create indexed document writer for {cluster} info {err}"))), + Err(err) => Err(notify_err!(format!("Could not create create indexed document writer for {cluster} info {err}"))), } } - Err(err) => Err(M3uFilterError::new(M3uFilterErrorKind::Info, format!("{err}"))), + Err(err) => Err(info_err!(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(notify_err!(format!("Could not create storage path for input {}", &input.name.as_ref().map_or("?", |v| v)))), + Err(err) => Err(notify_err!(format!("Could not create storage path for input {err}"))), } } @@ -737,12 +726,12 @@ pub async fn xtream_update_input_vod_record_from_wal_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)) - .map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Error accessing storage path: {err}"))) - .and_then(|opt| opt.ok_or_else(|| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Error accessing storage path for input: {}", input.name.clone().unwrap_or_else(|| input.id.to_string())))))?; + .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| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not read vod wal info {err}")))?; + wal_file.seek(SeekFrom::Start(0)).map_err(|err| notify_err!(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]; @@ -763,14 +752,14 @@ pub async fn xtream_update_input_vod_record_from_wal_file( let ts = u64::from_le_bytes(ts_bytes); tree_record_index.insert(provider_id, InputVodInfoRecord { tmdb_id, ts }); } - tree_record_index.store(&record_path).map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not store vod record info {err}")))?; + tree_record_index.store(&record_path).map_err(|err| notify_err!(format!("Could not store vod record info {err}")))?; drop(reader); if let Err(err) = fs::remove_file(wal_path) { error!("Failed to delete record WAL file for vod {err}"); } Ok(()) } - Err(err) => Err(M3uFilterError::new(M3uFilterErrorKind::Info, format!("{err}"))), + Err(err) => Err(info_err!(format!("{err}"))), } } @@ -781,11 +770,11 @@ pub async fn xtream_update_input_series_record_from_wal_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)) - .map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Error accessing storage path: {err}"))) - .and_then(|opt| opt.ok_or_else(|| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Error accessing storage path for input: {}", input.name.clone().unwrap_or_else(|| input.id.to_string())))))?; + .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| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not read vod wal info {err}")))?; + wal_file.seek(SeekFrom::Start(0)).map_err(|err| notify_err!(format!("Could not read vod wal info {err}")))?; let mut reader = BufReader::new(wal_file); let mut provider_id_bytes = [0u8; 4]; let mut ts_bytes = [0u8; 8]; @@ -801,7 +790,7 @@ pub async fn xtream_update_input_series_record_from_wal_file( let ts = u64::from_le_bytes(ts_bytes); tree_record_index.insert(provider_id, ts); } - tree_record_index.store(&record_path).map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Could not store series record info {err}")))?; + 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}"); @@ -809,6 +798,6 @@ pub async fn xtream_update_input_series_record_from_wal_file( Ok(()) } - Err(err) => Err(M3uFilterError::new(M3uFilterErrorKind::Info, format!("{err}"))), + Err(err) => Err(info_err!(format!("{err}"))), } } diff --git a/src/utils/config_reader.rs b/src/utils/config_reader.rs index fc78f3d88..d6117a2df 100644 --- a/src/utils/config_reader.rs +++ b/src/utils/config_reader.rs @@ -7,7 +7,7 @@ use log::{debug, error, info, warn}; use regex::Regex; use serde::Serialize; -use crate::{create_m3u_filter_error_result, handle_m3u_filter_error_result}; +use crate::{create_m3u_filter_error_result, handle_m3u_filter_error_result, info_err}; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::api_proxy::ApiProxyConfig; use crate::model::config::{Config, ConfigDto}; @@ -78,7 +78,7 @@ pub fn read_mapping(mapping_file: &str) -> Result, M3uFilterErr return Ok(Some(result)); }, Err(err) => { - return Err(M3uFilterError::new(M3uFilterErrorKind::Info, err.to_string())); + return Err(info_err!(err.to_string())); } } } diff --git a/src/utils/download.rs b/src/utils/download.rs index 40d8a3d0e..d8292f8e0 100644 --- a/src/utils/download.rs +++ b/src/utils/download.rs @@ -1,18 +1,15 @@ use crate::m3u_filter_error::M3uFilterError; use crate::model::config::{Config, ConfigInput, ConfigTarget}; -use crate::model::playlist::{FetchedPlaylist, PlaylistEntry, PlaylistGroup, PlaylistItem, PlaylistItemType, UUIDType, XtreamCluster}; +use crate::model::playlist::{PlaylistEntry, PlaylistGroup, XtreamCluster}; use crate::model::xmltv::TVGuide; -use crate::processing::xtream_parser::parse_xtream_series_info; use crate::processing::{m3u_parser, xtream_parser}; use crate::repository::xtream_repository::{xtream_get_input_info}; use crate::repository::xtream_repository; use crate::utils::{file_utils, request_utils}; use log::{debug, info}; use std::cmp::Ordering; -use std::collections::HashSet; use std::io::{Error, ErrorKind}; use std::path::PathBuf; -use std::rc::Rc; const ACTION_GET_SERIES_INFO: &str = "get_series_info"; const ACTION_GET_VOD_INFO: &str = "get_vod_info"; @@ -40,52 +37,52 @@ pub async fn get_m3u_playlist(cfg: &Config, input: &ConfigInput, working_dir: &s Err(err) => (vec![], vec![err]) } } - -pub async fn get_xtream_playlist_series(fpl: &mut FetchedPlaylist<'_>, process_uuids: HashSet>, errors: &mut Vec, resolve_delay: u16) -> Vec { - let input = fpl.input; - let mut result: Vec = vec![]; - for plg in &mut fpl.playlistgroups { - let mut group_series: Vec = vec![]; - for pli in &plg.channels { - let (fetch_series, series_info_url) = { - let mut header = pli.header.borrow_mut(); - let fetch_series = !header.series_fetched && header.item_type == PlaylistItemType::SeriesInfo && process_uuids.contains(header.get_uuid()); - if fetch_series { - header.series_fetched = true; - } - (fetch_series, header.url.to_string()) - }; - if fetch_series { - match request_utils::get_input_json_content(fpl.input, series_info_url.as_str(), None).await { - Ok(series_content) => { - match parse_xtream_series_info(&series_content, pli.header.borrow().group.as_str(), input) { - Ok(series_info) => { - if let Some(mut series) = series_info { - group_series.append(&mut series); - } - } - Err(err) => errors.push(err), - } - } - Err(err) => errors.push(err) - }; - if resolve_delay > 0 { - actix_web::rt::time::sleep(std::time::Duration::new(u64::from(resolve_delay), 0)).await; - } - } - } - if !group_series.is_empty() { - let group = PlaylistGroup { - id: plg.id, - title: plg.title.clone(), - channels: group_series, - xtream_cluster: XtreamCluster::Series, - }; - result.push(group); - } - } - result -} +// +// pub async fn get_xtream_playlist_series(fpl: &mut FetchedPlaylist<'_>, process_uuids: HashSet>, errors: &mut Vec, resolve_delay: u16) -> Vec { +// let input = fpl.input; +// let mut result: Vec = vec![]; +// for plg in &mut fpl.playlistgroups { +// let mut group_series: Vec = vec![]; +// for pli in &plg.channels { +// let (fetch_series, series_info_url) = { +// let mut header = pli.header.borrow_mut(); +// let fetch_series = !header.series_fetched && header.item_type == PlaylistItemType::SeriesInfo && process_uuids.contains(header.get_uuid()); +// if fetch_series { +// header.series_fetched = true; +// } +// (fetch_series, header.url.to_string()) +// }; +// if fetch_series { +// match request_utils::get_input_json_content(fpl.input, series_info_url.as_str(), None).await { +// Ok(series_content) => { +// match parse_xtream_series_info(&series_content, pli.header.borrow().group.as_str(), input) { +// Ok(series_info) => { +// if let Some(mut series) = series_info { +// group_series.append(&mut series); +// } +// } +// Err(err) => errors.push(err), +// } +// } +// Err(err) => errors.push(err) +// }; +// if resolve_delay > 0 { +// actix_web::rt::time::sleep(std::time::Duration::new(u64::from(resolve_delay), 0)).await; +// } +// } +// } +// if !group_series.is_empty() { +// let group = PlaylistGroup { +// id: plg.id, +// title: plg.title.clone(), +// channels: group_series, +// xtream_cluster: XtreamCluster::Series, +// }; +// result.push(group); +// } +// } +// result +// } pub fn get_xtream_player_api_action_url(input: &ConfigInput, action: &str) -> Option {