From e5173ad9031f0bf6af892dcbc764b71b47c7af5d Mon Sep 17 00:00:00 2001 From: euzu Date: Wed, 30 Oct 2024 11:14:55 +0100 Subject: [PATCH] Some optimizations, optional value compression for BPlustree --- CHANGELOG.md | 1 + src/api/m3u_api.rs | 15 ++- src/api/xmltv_api.rs | 12 +-- src/api/xtream_api.rs | 110 ++++++++++--------- src/auth/user.rs | 1 + src/filter.rs | 3 +- src/model/config.rs | 2 +- src/model/mapping.rs | 19 ++-- src/model/playlist.rs | 22 ++-- src/processing/m3u_parser.rs | 6 +- src/processing/playlist_processor.rs | 75 +++++++------ src/repository/bplustree.rs | 123 ++++++++++++---------- src/repository/epg_repository.rs | 9 +- src/repository/indexed_document_reader.rs | 23 ++-- src/repository/indexed_document_writer.rs | 74 +++++++------ src/repository/m3u_repository.rs | 116 ++++++++++---------- src/repository/playlist_repository.rs | 2 +- src/repository/target_id_mapping.rs | 19 ++-- src/repository/xtream_repository.rs | 58 +++++----- src/utils/config_reader.rs | 2 +- src/utils/download.rs | 7 +- src/utils/request_utils.rs | 2 +- test/rest-api.http | 4 +- 23 files changed, 333 insertions(+), 372 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 0e38ffc5e..a07995f74 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,6 +1,7 @@ # Changelog # 2.0.6 (2024-10-*) - breaking change virtual_id handling. You need to clear the data directory. +- fixed schedular implementation # 2.0.5(2024-10-16) - input url supports now scheme `file://...` (which is not necessary because file paths are supported). Gzip files are also supported. diff --git a/src/api/m3u_api.rs b/src/api/m3u_api.rs index fab63837f..2ad54a9cf 100644 --- a/src/api/m3u_api.rs +++ b/src/api/m3u_api.rs @@ -40,14 +40,13 @@ async fn m3u_api_stream( if target.has_output(&TargetType::M3u) { match get_target_storage_path(&app_state.config, target.name.as_str()) { Some(target_path) => { - if let Some((m3u_path, idx_path)) = m3u_get_file_paths(&target_path) { - match m3u_get_item_for_stream_id(&app_state.config, m3u_stream_id, &m3u_path, &idx_path) { - Ok(m3u_item) => { - return stream_response(m3u_item.url.as_str(), &req, None).await; - } - Err(err) => { - error!("Failed to get m3u url: {}", err); - } + let (m3u_path, idx_path) = m3u_get_file_paths(&target_path); + match m3u_get_item_for_stream_id(&app_state.config, m3u_stream_id, &m3u_path, &idx_path) { + Ok(m3u_item) => { + return stream_response(m3u_item.url.as_str(), &req, None).await; + } + Err(err) => { + error!("Failed to get m3u url: {}", err); } } } diff --git a/src/api/xmltv_api.rs b/src/api/xmltv_api.rs index ad310cff9..6b0791669 100644 --- a/src/api/xmltv_api.rs +++ b/src/api/xmltv_api.rs @@ -41,13 +41,11 @@ fn time_correct(date_time: &str, correction: &TimeDelta) -> String { } } -fn get_epg_path_for_target_of_type(target_name: &str, file_path: Option) -> Option { - if let Some(epg_path) = file_path { - if file_utils::path_exists(&epg_path) { - return Some(epg_path); - } - info!("Cant find epg file for {target_name} target: {}", epg_path.to_str().unwrap_or("?")); +fn get_epg_path_for_target_of_type(target_name: &str, epg_path: PathBuf) -> Option { + if file_utils::path_exists(&epg_path) { + return Some(epg_path); } + info!("Cant find epg file for {target_name} target: {}", epg_path.to_str().unwrap_or("?")); None } @@ -62,7 +60,7 @@ fn get_epg_path_for_target(config: &Config, target: &ConfigTarget) -> Option { if let Some(storage_path) = xtream_get_storage_path(config, &target.name) { - return get_epg_path_for_target_of_type(&target.name, Some(xtream_get_epg_file_path(&storage_path))); + return get_epg_path_for_target_of_type(&target.name, xtream_get_epg_file_path(&storage_path)); } } TargetType::Strm => {} diff --git a/src/api/xtream_api.rs b/src/api/xtream_api.rs index e095affaf..a2a16013e 100644 --- a/src/api/xtream_api.rs +++ b/src/api/xtream_api.rs @@ -23,7 +23,6 @@ use crate::repository::target_id_mapping::TargetIdMapping; use crate::repository::xtream_repository; use crate::utils::{json_utils, request_utils}; - macro_rules! try_option_bad_request { ($option:expr, $msg_is_error:expr, $msg:expr) => { match $option { @@ -200,7 +199,7 @@ async fn xtream_player_api_stream( } let (action_stream_id, stream_ext) = xtream_api_request_separate_number_and_rest(stream_req.stream_id); let virtual_id: u32 = try_result_bad_request!(action_stream_id.trim().parse()); - let pli = try_result_bad_request!(xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, None), true, format!("Failed to read xtream item for stream id {}", virtual_id)); + let pli = try_result_bad_request!(xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, None), true, format!("Failed to read xtream item for stream id {}", virtual_id)); let input = try_option_bad_request!(app_state.config.get_input_by_id(pli.input_id), true, format!("Cant find input for target {target_name}, context {}, stream_id {virtual_id}", stream_req.context)); let query_path = if stream_req.action_path.is_empty() { @@ -430,7 +429,7 @@ async fn xtream_get_catchup_response(app_state: &AppState, target: &ConfigTarget // TODO epg_id if let Some(catchup_provider_id) = epg_list_item.get("id").and_then(Value::as_str).and_then(|id| id.parse::().ok()) { let uuid = hash_string(&format!("{}/{}", pli.url, catchup_provider_id)); - let virtual_id = target_id_mapping.insert_entry(uuid, catchup_provider_id, &PlaylistItemType::Catchup, pli.provider_id); + let virtual_id = target_id_mapping.insert_entry(uuid, catchup_provider_id, PlaylistItemType::Catchup, pli.provider_id); epg_list_item.insert("id".to_string(), Value::String(virtual_id.to_string())); } } @@ -450,67 +449,64 @@ async fn xtream_player_api( api_req: UserApiRequest, app_state: &web::Data, ) -> HttpResponse { - match get_user_target(&api_req, app_state) { - Some((user, target)) => { - let action = api_req.action.trim(); - let target_name = &target.name; - if target.has_output(&TargetType::Xtream) { - if action.is_empty() { - return HttpResponse::Ok().json(get_user_info(&user, &app_state.config)); - } + if let Some((user, target)) = get_user_target(&api_req, app_state) { + let action = api_req.action.trim(); + let target_name = &target.name; + if target.has_output(&TargetType::Xtream) { + if action.is_empty() { + return HttpResponse::Ok().json(get_user_info(&user, &app_state.config)); + } - match action { - "get_series_info" => { - xtream_get_stream_info_response(app_state, &user, target, - api_req.series_id.trim(), - XtreamCluster::Series).await - } - "get_vod_info" => { - xtream_get_stream_info_response(app_state, &user, target, - api_req.vod_id.trim(), - XtreamCluster::Video).await - } - "get_epg" | - "get_short_epg" => { - xtream_get_short_epg(app_state, &user, target, - api_req.stream_id.trim(), - api_req.limit.trim()).await - } - "get_simple_data_table" => { - xtream_get_catchup_response(app_state, target, - api_req.stream_id.trim(), - api_req.start.trim(), - api_req.end.trim()).await - } - _ => { - let category_id = api_req.category_id.as_str().trim(); - if let Some(response) = xtream_player_api_handle_content_action(&app_state.config, target_name, action, category_id, req).await { - response - } else { - let cat_id = if category_id.is_empty() { 0 } else { category_id.parse::().unwrap_or(0) }; - match match action { - "get_live_streams" => xtream_repository::xtream_load_rewrite_playlist(&XtreamCluster::Live, &app_state.config, target, cat_id), - "get_vod_streams" => xtream_repository::xtream_load_rewrite_playlist(&XtreamCluster::Video, &app_state.config, target, cat_id), - "get_series" => xtream_repository::xtream_load_rewrite_playlist(&XtreamCluster::Series, &app_state.config, target, cat_id), - _ => Err(Error::new(ErrorKind::Unsupported, format!("Cant find action: {action} for target: {target_name}"))), - } { - Ok(payload) => HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body(payload), - Err(err) => { - error!("Could not create response for xtream target: {} action: {} err: {}", target_name, action, err); - HttpResponse::NoContent().finish() - } + match action { + "get_series_info" => { + xtream_get_stream_info_response(app_state, &user, target, + api_req.series_id.trim(), + XtreamCluster::Series).await + } + "get_vod_info" => { + xtream_get_stream_info_response(app_state, &user, target, + api_req.vod_id.trim(), + XtreamCluster::Video).await + } + "get_epg" | + "get_short_epg" => { + xtream_get_short_epg(app_state, &user, target, + api_req.stream_id.trim(), + api_req.limit.trim()).await + } + "get_simple_data_table" => { + xtream_get_catchup_response(app_state, target, + api_req.stream_id.trim(), + api_req.start.trim(), + api_req.end.trim()).await + } + _ => { + let category_id = api_req.category_id.as_str().trim(); + if let Some(response) = xtream_player_api_handle_content_action(&app_state.config, target_name, action, category_id, req).await { + response + } else { + let cat_id = if category_id.is_empty() { 0 } else { category_id.parse::().unwrap_or(0) }; + match match action { + "get_live_streams" => xtream_repository::xtream_load_rewrite_playlist(XtreamCluster::Live, &app_state.config, target, cat_id), + "get_vod_streams" => xtream_repository::xtream_load_rewrite_playlist(XtreamCluster::Video, &app_state.config, target, cat_id), + "get_series" => xtream_repository::xtream_load_rewrite_playlist(XtreamCluster::Series, &app_state.config, target, cat_id), + _ => Err(Error::new(ErrorKind::Unsupported, format!("Cant find action: {action} for target: {target_name}"))), + } { + Ok(payload) => HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body(payload), + Err(err) => { + error!("Could not create response for xtream target: {} action: {} err: {}", target_name, action, err); + HttpResponse::NoContent().finish() } } } } - } else { - HttpResponse::Ok().json(get_user_info(&user, &app_state.config)) } + } else { + HttpResponse::Ok().json(get_user_info(&user, &app_state.config)) } - _ => { - debug!("{}", if api_req.action.is_empty() { "Paremeter action is empty!" } else { "cant find user!" }); - HttpResponse::BadRequest().finish() - } + } else { + debug!("{}", if api_req.action.is_empty() { "Paremeter action is empty!" } else { "cant find user!" }); + HttpResponse::BadRequest().finish() } } diff --git a/src/auth/user.rs b/src/auth/user.rs index ae6875b09..59adfce2f 100644 --- a/src/auth/user.rs +++ b/src/auth/user.rs @@ -1,5 +1,6 @@ use std::ptr; +#[allow(clippy::unsafe_derive_deserialize)] #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub(crate) struct UserCredential { pub username: String, diff --git a/src/filter.rs b/src/filter.rs index 0e6b1141e..27365b726 100644 --- a/src/filter.rs +++ b/src/filter.rs @@ -195,8 +195,7 @@ impl std::fmt::Display for Filter { write!(f, "{} = {}", field, match item_type { PlaylistItemType::Live => "live", PlaylistItemType::Video => "vod", - PlaylistItemType::Series => "series", - PlaylistItemType::SeriesInfo => "series", // yes series-info is handled as series in filter + PlaylistItemType::Series | PlaylistItemType::SeriesInfo => "series", // yes series-info is handled as series in filter _ => "unsupported" }) } diff --git a/src/model/config.rs b/src/model/config.rs index 1e9f84b4f..1bab40697 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -873,7 +873,7 @@ impl Config { None } - pub(crate) fn set_mappings(&mut self, mappings_cfg: Mappings) { + pub(crate) fn set_mappings(&mut self, mappings_cfg: &Mappings) { for source in &mut self.sources { for target in &mut source.targets { if let Some(mapping_ids) = &target.mapping { diff --git a/src/model/mapping.rs b/src/model/mapping.rs index 41e2b8019..1862d46b6 100644 --- a/src/model/mapping.rs +++ b/src/model/mapping.rs @@ -158,7 +158,7 @@ impl MapperTransform { new_pattern = apply_templates_to_pattern(pattern, template_list); } } - let re = Regex::new(&*new_pattern); + let re = Regex::new(&new_pattern); if re.is_err() { return create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "cant parse regex: {}", new_pattern); } @@ -196,17 +196,17 @@ pub(crate) struct Mapper { impl Mapper { pub fn prepare(&mut self, templates: Option<&Vec>, tags: Option<&Vec>) -> Result<(), M3uFilterError> { - for (key, _) in &self.attributes { + 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}"))); } } - for (key, _) in &self.suffix { + 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}"))); } } - for (key, _) in &self.prefix { + 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}"))); } @@ -228,9 +228,7 @@ impl Mapper { if !valid_property!(field, AFFIX_FIELDS) { return Err(M3uFilterError::new(M3uFilterErrorKind::Info, format!("Invalid mapper transform field {field}"))); } - if let Err(err) = t.prepare(templates) { - return Err(err); - } + t.prepare(templates)?; } } } @@ -456,12 +454,11 @@ impl Mapping { Ok(flt) => { counters.push(MappingCounter { filter: flt, - field: def.field.to_owned(), - concat: def.concat.to_owned(), + field: def.field.clone(), + concat: def.concat.clone(), modifier: def.modifier.clone(), value: Arc::new(Mutex::new(def.value)), - } - ) + }); } Err(e) => return Err(M3uFilterError::new(M3uFilterErrorKind::Info, e.to_string())) } diff --git a/src/model/playlist.rs b/src/model/playlist.rs index 692841cf4..8beb2ac00 100644 --- a/src/model/playlist.rs +++ b/src/model/playlist.rs @@ -35,9 +35,10 @@ impl FetchedPlaylist<'_> { } } -#[derive(Debug, Copy, Clone, Eq, Hash, PartialEq, Serialize, Deserialize)] +#[derive(Debug, Copy, Clone, Eq, Hash, PartialEq, Serialize, Deserialize, Default)] #[repr(u8)] pub(crate) enum XtreamCluster { + #[default] Live = 1, Video = 2, Series = 3, @@ -54,12 +55,6 @@ impl XtreamCluster { } } -impl Default for XtreamCluster { - fn default() -> Self { - XtreamCluster::Live - } -} - impl Display for XtreamCluster { fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { write!(f, "{}", self.as_str()) @@ -78,9 +73,10 @@ impl TryFrom for XtreamCluster { } } -#[derive(Debug, Copy, Clone, Eq, Hash, PartialEq, Serialize, Deserialize)] +#[derive(Debug, Copy, Clone, Eq, Hash, PartialEq, Serialize, Deserialize, Default)] #[repr(u8)] pub(crate) enum PlaylistItemType { + #[default] Live = 1, Video = 2, Series = 3, // xtream series description @@ -89,12 +85,6 @@ pub(crate) enum PlaylistItemType { Catchup = 6, } -impl Default for PlaylistItemType { - fn default() -> Self { - PlaylistItemType::Live - } -} - impl From for PlaylistItemType { fn from(xtream_cluster: XtreamCluster) -> Self { match xtream_cluster { @@ -154,7 +144,7 @@ pub(crate) struct PlaylistItemHeader { impl PlaylistItemHeader { pub(crate) fn gen_uuid(&mut self) { - self.uuid = Rc::new(hash_string(&self.url)) + self.uuid = Rc::new(hash_string(&self.url)); } pub(crate) fn get_uuid(&self) -> &Rc<[u8; 32]> { &self.uuid @@ -343,7 +333,7 @@ impl PlaylistItem { Err(_) => None } }, - item_type: header.item_type.clone(), + item_type: header.item_type, series_fetched: header.series_fetched, category_id: header.category_id, input_id: header.input_id, diff --git a/src/processing/m3u_parser.rs b/src/processing/m3u_parser.rs index e366f5d7f..5457b3d58 100644 --- a/src/processing/m3u_parser.rs +++ b/src/processing/m3u_parser.rs @@ -9,7 +9,7 @@ use crate::utils::string_utils; #[inline] fn token_value(it: &mut std::str::Chars) -> String { // Use .find() to skip until the first double quote (") character. - if it.find(|&ch| ch == '"').is_some() { + if it.any(|ch| ch == '"') { // If a quote is found, call get_value to extract the value. return get_value(it); } @@ -93,7 +93,7 @@ macro_rules! process_header_fields { }; } -fn process_header(input: &ConfigInput, video_suffixes: &Vec<&str>, content: &str, url: &str) -> PlaylistItemHeader { +fn process_header(input: &ConfigInput, video_suffixes: &[&str], content: &str, url: &str) -> PlaylistItemHeader { let mut plih = create_empty_playlistitem_header(input.id, url); let mut it = content.chars(); let line_token = token_till(&mut it, ':', false); @@ -186,7 +186,7 @@ where continue; } if let Some(header_value) = header { - let item = PlaylistItem { header: RefCell::new(process_header(input, &video_suffixes, &header_value, &line)) }; + let item = PlaylistItem { header: RefCell::new(process_header(input, &video_suffixes, &header_value, line)) }; let mut header = item.header.borrow_mut(); if header.group.is_empty() { if let Some(group_value) = group { diff --git a/src/processing/playlist_processor.rs b/src/processing/playlist_processor.rs index 7244f7ce1..1442c2c03 100644 --- a/src/processing/playlist_processor.rs +++ b/src/processing/playlist_processor.rs @@ -68,41 +68,38 @@ fn playlistitem_comparator(a: &PlaylistItem, b: &PlaylistItem, channel_sort: &Co let raw_value_b = get_field_value(b, &channel_sort.field); let value_a = if match_as_ascii { Rc::new(unidecode(&raw_value_a)) } else { raw_value_a }; let value_b = if match_as_ascii { Rc::new(unidecode(&raw_value_b)) } else { raw_value_b }; - match &channel_sort.sequence { - Some(custom_order) => { - // Check indices in the custom order vector - let index_a = custom_order.iter().position(|s| s == value_a.as_ref()); - let index_b = custom_order.iter().position(|s| s == value_b.as_ref()); + if let Some(custom_order) = &channel_sort.sequence { + // Check indices in the custom order vector + let index_a = custom_order.iter().position(|s| s == value_a.as_ref()); + let index_b = custom_order.iter().position(|s| s == value_b.as_ref()); - match (index_a, index_b) { - (Some(idx_a), Some(idx_b)) => { - // Both items found in custom order, compare indices - idx_a.cmp(&idx_b) - } - (Some(_), None) => { - // Only 'a' found in custom order, it comes first - Ordering::Less - } - (None, Some(_)) => { - // Only 'b' found in custom order, it comes first - Ordering::Greater - } - (None, None) => { - // Neither found, fall back to default ordering - let ordering = value_a.partial_cmp(&value_b).unwrap(); - match channel_sort.order { - Asc => ordering, - Desc => ordering.reverse(), - } + match (index_a, index_b) { + (Some(idx_a), Some(idx_b)) => { + // Both items found in custom order, compare indices + idx_a.cmp(&idx_b) + } + (Some(_), None) => { + // Only 'a' found in custom order, it comes first + Ordering::Less + } + (None, Some(_)) => { + // Only 'b' found in custom order, it comes first + Ordering::Greater + } + (None, None) => { + // Neither found, fall back to default ordering + let ordering = value_a.partial_cmp(&value_b).unwrap(); + match channel_sort.order { + Asc => ordering, + Desc => ordering.reverse(), } } } - None => { - let ordering = value_a.partial_cmp(&value_b).unwrap(); - match channel_sort.order { - Asc => ordering, - Desc => ordering.reverse() - } + } else { + let ordering = value_a.partial_cmp(&value_b).unwrap(); + match channel_sort.order { + Asc => ordering, + Desc => ordering.reverse() } } } @@ -251,7 +248,7 @@ fn map_playlist_counter(target: &ConfigTarget, playlist: &[PlaylistGroup]) { if let Some(counter_list) = &mapping.t_counter { for counter in counter_list { let mut cntval = counter.value.lock().unwrap(); - for plg in &*playlist { + for plg in playlist { for channel in &plg.channels { let provider = ValueProvider { pli: RefCell::new(channel) }; if counter.filter.filter(&provider, &mut mock_processor) { @@ -260,12 +257,12 @@ fn map_playlist_counter(target: &ConfigTarget, playlist: &[PlaylistGroup]) { } else { let value = match channel.header.borrow_mut().get_field(&counter.field) { Some(field_value) => field_value.to_string(), - None => "".to_string(), + None => String::new(), }; if counter.modifier == CounterModifier::Suffix { - format!("{}{}{}", value, counter.concat, cntval.to_string()) + format!("{value}{}{cntval}", counter.concat) } else { - format!("{}{}{}", cntval.to_string(), counter.concat, value) + format!("{cntval}{}{value}", counter.concat) } }; channel.header.borrow_mut().set_field(&counter.field, new_value.as_str()); @@ -352,7 +349,7 @@ async fn process_source(cfg: Arc, source_idx: usize, user_targets: Arc

{} - Err(mut err) => errors.extend(err.drain(..)) + Err(mut err) => errors.append(&mut err) } } } @@ -363,8 +360,8 @@ async fn process_source(cfg: Arc, source_idx: usize, user_targets: Arc

InputStats { InputStats { name: input_name.to_string(), - input_type: input_type, - error_count: error_count, + input_type, + error_count, raw_stats: PlaylistStats { group_count, channel_count, @@ -500,7 +497,7 @@ async fn process_playlist<'a>(playlists: &mut [FetchedPlaylist<'a>], let epg_channel_ids: HashSet<_> = fp.playlistgroups.iter().flat_map(|g| &g.channels) .filter_map(|c| c.header.borrow().epg_channel_id.clone()).collect(); - new_playlist.extend(fp.playlistgroups.drain(..)); + new_playlist.append(&mut fp.playlistgroups); if !epg_channel_ids.is_empty() { if let Some(tv_guide) = fp.epg { debug!("found epg information for {}", &target.name); diff --git a/src/repository/bplustree.rs b/src/repository/bplustree.rs index f7cac14bc..dd4e6fdb7 100644 --- a/src/repository/bplustree.rs +++ b/src/repository/bplustree.rs @@ -1,6 +1,6 @@ use std::array::TryFromSliceError; use std::fs::{File, OpenOptions}; -use std::io::{self, Read, Seek, SeekFrom, Write}; +use std::io::{self, Error, ErrorKind, Read, Seek, SeekFrom, Write}; use std::marker::PhantomData; use std::path::Path; @@ -35,15 +35,15 @@ fn u32_from_bytes(bytes: &[u8]) -> io::Result { } #[inline] -fn bincode_serialize(value: &T) -> io::Result> +fn bincode_serialize(value: &T) -> io::Result> where - T: serde::Serialize, + T: ?Sized + serde::Serialize, { bincode::serialize(value).map_err(|e| io::Error::new(io::ErrorKind::Other, e.to_string())) } #[inline] -fn bincode_deserialize(value: &[u8]) -> io::Result +fn bincode_deserialize(value: &[u8]) -> io::Result where T: for<'a> serde::Deserialize<'a>, { @@ -51,10 +51,9 @@ where } -fn get_entry_index_upper_bound(keys: &Vec, key: &K) -> usize +fn get_entry_index_upper_bound(keys: &[K], key: &K) -> usize where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone, - V: Serialize + for<'de> Deserialize<'de> + Clone, { let mut left = 0; let mut right = keys.len(); @@ -103,19 +102,19 @@ where order >> 1 } - fn find_leaf_entry<'a>(&self, node: &'a BPlusTreeNode) -> &'a K { + fn find_leaf_entry(node: &BPlusTreeNode) -> &K { if node.is_leaf { - node.keys.get(0).unwrap() + node.keys.first().unwrap() } else { - let child = node.children.get(0).unwrap(); - self.find_leaf_entry(child) + let child = node.children.first().unwrap(); + Self::find_leaf_entry(child) } } #[allow(dead_code)] fn query(&self, key: &K) -> Option<&V> { if self.is_leaf { - return match self.keys.binary_search(&key) { + return match self.keys.binary_search(key) { Ok(idx) => self.values.get(idx), Err(_) => None, }; @@ -133,19 +132,18 @@ where while left <= right { let mid = left + ((right - left) >> 1); let mid_key = &self.keys[mid]; - if mid_key == key { - return Some(mid); - } else if mid_key > key { - right = mid.checked_sub(1)?; - } else { - left = mid + 1; + + match mid_key.cmp(key) { + std::cmp::Ordering::Equal => return Some(mid), + std::cmp::Ordering::Greater => right = mid.checked_sub(1)?, + std::cmp::Ordering::Less => left = mid + 1, } } None } fn get_entry_index_upper_bound(&self, key: &K) -> usize { - get_entry_index_upper_bound::(&self.keys, key) + get_entry_index_upper_bound::(&self.keys, key) } fn insert(&mut self, key: K, v: V, inner_order: usize, leaf_order: usize) -> Option> { @@ -169,9 +167,9 @@ where let child = self.children.get_mut(pos).unwrap(); let node = child.insert(key.clone(), v, inner_order, leaf_order); if node.is_some() { - let leaf_key = self.find_leaf_entry(node.as_ref().unwrap()); + let leaf_key = Self::find_leaf_entry(node.as_ref().unwrap()); let idx = self.get_entry_index_upper_bound(leaf_key); - if let Err(_) = self.keys.binary_search(&key) { + if self.keys.binary_search(&key).is_err() { self.keys.insert(idx, leaf_key.clone()); self.children.insert(idx + 1, node.unwrap()); if self.is_overflow(inner_order) { @@ -194,7 +192,7 @@ where let mut node = BPlusTreeNode::new(false); node.keys = self.keys.split_off(median + 1); node.children = self.children.split_off(median + 1); - self.children.push(node.children.get(0).unwrap().clone()); + self.children.push(node.children.first().unwrap().clone()); node } } @@ -214,13 +212,13 @@ where let buffer_slice = &mut buffer[..]; // Write node type (leaf or internal) - buffer_slice[0] = if self.is_leaf { 1u8 } else { 0u8 }; + buffer_slice[0] = u8::from(self.is_leaf); let mut write_pos = 1; // Serialize and write keys let keys_encoded = bincode_serialize(&self.keys)?; let keys_bytes_len = keys_encoded.len(); - buffer_slice[write_pos..write_pos + 4].copy_from_slice(&(keys_bytes_len as u32).to_le_bytes()); + buffer_slice[write_pos..write_pos + 4].copy_from_slice(&(u32::try_from(keys_bytes_len).map_err(|err| Error::new(ErrorKind::Other, err))?).to_le_bytes()); write_pos += 4; buffer_slice[write_pos..write_pos + keys_bytes_len].copy_from_slice(&keys_encoded); write_pos += keys_bytes_len; @@ -228,13 +226,21 @@ where // If leaf, serialize and write values if self.is_leaf { let values_encoded = bincode_serialize(&self.values)?; - let mut encoder = flate2::write::ZlibEncoder::new(Vec::new(), Compression::fast()); - encoder.write_all(&values_encoded)?; - let compressed_bytes = encoder.finish()?; - let values_bytes_len = compressed_bytes.len(); - buffer_slice[write_pos..write_pos + 4].copy_from_slice(&(values_bytes_len as u32).to_le_bytes()); + let use_compression = values_encoded.len() + 5 > BLOCK_SIZE; + let compression_byte = if use_compression {1u8.to_le_bytes()} else {0u8.to_le_bytes()}; + buffer_slice[write_pos..=write_pos].copy_from_slice(&compression_byte); + write_pos += 1; + let content_bytes = if use_compression { + let mut encoder = flate2::write::ZlibEncoder::new(Vec::new(), Compression::fast()); + encoder.write_all(&values_encoded)?; + encoder.finish()? + } else { + values_encoded + }; + let values_bytes_len = content_bytes.len(); + buffer_slice[write_pos..write_pos + 4].copy_from_slice(&(u32::try_from(values_bytes_len).map_err(|err| Error::new(ErrorKind::Other, err))?).to_le_bytes()); write_pos += 4; - buffer_slice[write_pos..write_pos + values_bytes_len].copy_from_slice(&compressed_bytes); + buffer_slice[write_pos..write_pos + values_bytes_len].copy_from_slice(&content_bytes); write_pos += values_bytes_len; } @@ -252,7 +258,7 @@ where } let pointer_encoded = bincode_serialize(&pointer)?; - let pointer_bytes_len = pointer_encoded.len() as u32; + let pointer_bytes_len = u32::try_from(pointer_encoded.len()).map_err(|err| Error::new(ErrorKind::Other, err))?; file.seek(SeekFrom::Start(pointer_offset))?; file.write_all(&pointer_bytes_len.to_le_bytes())?; @@ -278,13 +284,19 @@ where // Deserialize values if leaf node let values = if is_leaf { + let use_compression = u8::from_le_bytes(buffer[read_pos..=read_pos].try_into().unwrap()) == 1; + read_pos += 1; let values_length = u32_from_bytes(&buffer[read_pos..read_pos + 4])? as usize; read_pos += 4; - let compressed_bytes = &buffer[read_pos..read_pos + values_length]; - let mut decoder = flate2::write::ZlibDecoder::new(Vec::new()); - decoder.write_all(&compressed_bytes)?; - let values_bytes = decoder.finish()?; - let values: Vec = bincode_deserialize(&values_bytes)?; + let content_bytes = &buffer[read_pos..read_pos + values_length]; + let values_bytes = if use_compression { + let mut decoder = flate2::write::ZlibDecoder::new(Vec::new()); + decoder.write_all(content_bytes)?; + &decoder.finish()? + } else { + content_bytes + }; + let values: Vec = bincode_deserialize(values_bytes)?; read_pos += values_length; values } else { @@ -292,7 +304,9 @@ where }; // Deserialize children indices if internal node - let (children, children_pointer) = if !is_leaf { + let (children, children_pointer) = if is_leaf { + (vec![], None) + } else { let pointers_length = u32_from_bytes(&buffer[read_pos..read_pos + 4])? as usize; read_pos += 4; let pointers: Vec = bincode_deserialize(&buffer[read_pos..read_pos + pointers_length])?; @@ -310,16 +324,9 @@ where } else { (vec![], Some(pointers)) } - } else { - (vec![], None) }; - Ok((BPlusTreeNode { - is_leaf, - keys, - values, - children, - }, children_pointer)) + Ok((BPlusTreeNode { keys, children, is_leaf, values }, children_pointer)) } } @@ -366,9 +373,9 @@ where if let Some(node) = self.root.insert(key, value, self.inner_order, self.leaf_order) { let child_key = if node.is_leaf { - node.keys.get(0).as_ref().unwrap() + node.keys.first().as_ref().unwrap() } else { - node.find_leaf_entry(&node) + BPlusTreeNode::::find_leaf_entry(&node) }; let mut new_root = BPlusTreeNode::::new(false); @@ -386,7 +393,7 @@ where } pub(crate) fn serialize(&self, filepath: &Path) -> io::Result { - let mut file = OpenOptions::new().write(true).create(true).open(filepath)?; + let mut file = OpenOptions::new().write(true).create(true).truncate(true).open(filepath)?; let mut buffer = vec![0u8; BLOCK_SIZE]; let result = self.root.serialize_to_block(&mut file, &mut buffer, 0u64); file.flush()?; @@ -427,7 +434,7 @@ where Err(_) => None, }; } - let child_idx = get_entry_index_upper_bound::(&node.keys, key); + let child_idx = get_entry_index_upper_bound::(&node.keys, key); offset = *pointers.unwrap().get(child_idx).unwrap(); } Err(err) => { @@ -453,8 +460,8 @@ where let file = is_file_valid(File::open(filepath)?)?; Ok(BPlusTreeQuery { file, - _marker_k: Default::default(), - _marker_v: Default::default(), + _marker_k: PhantomData, + _marker_v: PhantomData, }) } @@ -485,8 +492,8 @@ where let file = is_file_valid(file)?; Ok(BPlusTreeUpdate { file, - _marker_k: Default::default(), - _marker_v: Default::default(), + _marker_k: PhantomData, + _marker_v: PhantomData, }) } @@ -494,7 +501,7 @@ where query_tree(&mut self.file, key) } - fn serialize_node(&mut self, offset: u64, node: BPlusTreeNode) -> io::Result { + fn serialize_node(&mut self, offset: u64, node: &BPlusTreeNode) -> io::Result { let mut buffer = vec![0u8; BLOCK_SIZE]; let result = node.serialize_to_block(&mut self.file, &mut buffer, offset); self.file.flush()?; @@ -511,16 +518,16 @@ where return match node.keys.binary_search(key) { Ok(idx) => { let old_value = node.values.get(idx); - if let Some(_) = old_value { + if old_value.is_some() { node.values[idx] = value; - return self.serialize_node(offset, node); + return self.serialize_node(offset, &node); } Err(io::Error::new(io::ErrorKind::NotFound, "Entry not found")) } Err(_) => Err(io::Error::new(io::ErrorKind::NotFound, "Entry not found")), }; } - let child_idx = get_entry_index_upper_bound::(&node.keys, key); + let child_idx = get_entry_index_upper_bound::(&node.keys, key); offset = *pointers.unwrap().get(child_idx).unwrap(); } Err(err) => { @@ -632,13 +639,13 @@ mod tests { for i in 0u32..=500 { tree.insert(i, Record { id: i, - data: format!("Entry {}", i+1), + data: format!("Entry {}", i + 1), }); } tree.traverse(|keys, values| { keys.iter().zip(values.iter()).for_each(|(k, v)| { - assert!(format!("Entry {}", k+1).eq(&v.data), "Wrong entry") + assert!(format!("Entry {}", k + 1).eq(&v.data), "Wrong entry") }); }); diff --git a/src/repository/epg_repository.rs b/src/repository/epg_repository.rs index 528497d71..806e902d1 100644 --- a/src/repository/epg_repository.rs +++ b/src/repository/epg_repository.rs @@ -51,12 +51,11 @@ pub(crate) fn epg_write(target: &ConfigTarget, cfg: &Config, target_path: &Path, M3uFilterErrorKind::Notify, format!("write epg for target {} failed: No filename set", target.name))); } - if let Some(path) = m3u_get_epg_file_path(target_path) { - if log_enabled!(Level::Debug) { - debug!("writing m3u epg to {}", path.to_str().unwrap_or("?")); - } - epg_write_file(target, epg_data, &path)?; + let path = m3u_get_epg_file_path(target_path); + if log_enabled!(Level::Debug) { + debug!("writing m3u epg to {}", path.to_str().unwrap_or("?")); } + epg_write_file(target, epg_data, &path)?; } TargetType::Xtream => { match xtream_get_storage_path(cfg, &target.name) { diff --git a/src/repository/indexed_document_reader.rs b/src/repository/indexed_document_reader.rs index 8c6d4a1c0..ae0e48be2 100644 --- a/src/repository/indexed_document_reader.rs +++ b/src/repository/indexed_document_reader.rs @@ -3,22 +3,23 @@ use std::fs::File; use std::io::{Error, ErrorKind, Read, Seek, SeekFrom}; use std::marker::PhantomData; use std::path::Path; + use crate::repository::bplustree::BPlusTreeQuery; use crate::repository::indexed_document_writer::OffsetPointer; fn get_offset(index_path: &Path, doc_id: u32) -> Result { - match BPlusTreeQuery::::try_new(index_path) { + match BPlusTreeQuery::::try_new(index_path) { Ok(mut tree) => { match tree.query(&doc_id) { - Some(offset) => Ok(offset as u64), + Some(offset) => Ok(u64::from(offset)), None => Err(Error::new(ErrorKind::NotFound, format!("doc_id not found {doc_id}"))), } - }, + } Err(err) => Err(err) } } -fn read_content_size(main_file: &mut File) -> Result { +fn read_content_size(main_file: &mut File) -> Result { let mut size_bytes = [0u8; 4]; main_file.read_exact(&mut size_bytes)?; let buf_size = u32::from_le_bytes(size_bytes) as usize; @@ -34,7 +35,7 @@ pub(in crate::repository) struct IndexedDocumentReader { t_type: PhantomData, } -impl IndexedDocumentReader { +impl IndexedDocumentReader { pub fn new(main_path: &Path) -> Result, Error> { if main_path.exists() { match File::open(main_path) { @@ -82,26 +83,24 @@ impl IndexedDocumentReader { self.t_buffer.reserve(buf_size - self.t_buffer.capacity()); } self.t_buffer.resize(buf_size, 0u8); - // read content self.main_file.read_exact(&mut self.t_buffer[0..buf_size])?; - self.cursor += buf_size as u32; - + self.cursor += u32::try_from(buf_size).map_err(|err| Error::new(ErrorKind::Other, err))?; // deserialize buffer - return match bincode::deserialize::(&self.t_buffer[0..buf_size]) { + match bincode::deserialize::(&self.t_buffer[0..buf_size]) { Ok(value) => Ok(Some(value)), Err(err) => { self.failed = true; Err(Error::new(ErrorKind::Other, format!("Failed to deserialize document {err}"))) } - }; + } } pub(in crate::repository) fn read_indexed_item(main_path: &Path, index_path: &Path, doc_id: u32) -> Result { if main_path.exists() && index_path.exists() { // get the offset from index - let offset = crate::repository::indexed_document_reader::get_offset(index_path, doc_id)?; + let offset = get_offset(index_path, doc_id)?; let mut main_file = File::open(main_path)?; main_file.seek(SeekFrom::Start(offset))?; let buf_size = read_content_size(&mut main_file)?; @@ -115,7 +114,7 @@ impl IndexedDocumentReader { } } -impl Iterator for IndexedDocumentReader { +impl Iterator for IndexedDocumentReader { type Item = T; // Implement the next() method diff --git a/src/repository/indexed_document_writer.rs b/src/repository/indexed_document_writer.rs index a98acae7f..e80e82a5e 100644 --- a/src/repository/indexed_document_writer.rs +++ b/src/repository/indexed_document_writer.rs @@ -16,7 +16,7 @@ pub(in crate::repository) type OffsetPointer = u32; * - index * * Layout of content file record is: -* - content-size (u32) + content +* - content-size (u32) + content (deflate) * * index file is a bplustree */ @@ -30,32 +30,32 @@ pub(in crate::repository) struct IndexedDocumentWriter { impl IndexedDocumentWriter { fn new_with_mode(main_path: PathBuf, index_path: PathBuf, append: bool) -> Result { - match open_file_append(&main_path, append) { - Ok(main_file) => { - let main_offset = match &main_file.metadata() { - Ok(meta) => u32::try_from(meta.len()).map_err(|err| Error::new(ErrorKind::Other, err))?, - Err(_) => 0 - }; + // Attempt to open the main file in the specified mode (append or not) + let main_file = open_file_append(&main_path, append)?; - let index_tree = if append && index_path.exists() { - BPlusTree::::deserialize(&index_path).unwrap_or_else(|err| { - error!("Failed to load index {:?} {err}", index_path); - BPlusTree::::new() - }) - } else { - BPlusTree::::new() - }; + // Retrieve file size and convert to `u32` for `main_offset`, if possible + let main_offset = main_file + .metadata() + .and_then(|meta| u32::try_from(meta.len()).map_err(|err| Error::new(ErrorKind::Other, err))) + .unwrap_or(0); - Ok(Self { - main_path, - index_path, - main_file, - main_offset, - index_tree - }) - } - Err(e) => Err(e) - } + // Initialize the index tree (BPlusTree) - either by deserializing an existing one or creating a new one + let index_tree = if append && index_path.exists() { + BPlusTree::::deserialize(&index_path).unwrap_or_else(|err| { + error!("Failed to load index {:?}: {}", index_path, err); + BPlusTree::::new() + }) + } else { + BPlusTree::::new() + }; + + Ok(Self { + main_path, + index_path, + main_file, + main_offset, + index_tree, + }) } pub fn new(main_path: PathBuf, index_path: PathBuf) -> Result { @@ -74,19 +74,17 @@ impl IndexedDocumentWriter { pub fn write_doc(&mut self, doc_id: u32, doc: &T) -> Result<(), Error> where T: ?Sized + serde::Serialize { - if let Ok(encoded) = bincode::serialize(doc) { - let content_bytes_len = u32::try_from(encoded.len()).map_err(|err| Error::new(ErrorKind::Other, err))?; - let mut data: Vec = content_bytes_len.to_le_bytes().to_vec(); - data.extend(&encoded); - match file_utils::check_write(&self.main_file.write_all(&data)) { - Ok(()) => { - self.index_tree.insert(doc_id, self.main_offset); - let written_bytes = u32::try_from(data.len()).map_err(|err| Error::new(ErrorKind::Other, err))?; - self.main_offset += written_bytes; - } - Err(err) => { - return Err(Error::new(ErrorKind::Other, format!("failed to write document: {} - {}", self.main_path.to_str().unwrap(), err))); - } + let encoded_bytes = bincode::serialize(doc).map_err(|_| Error::new(ErrorKind::InvalidData, "Failed to deserialize document"))?; + let encoded_bytes_len = u32::try_from(encoded_bytes.len()).map_err(|err| Error::new(ErrorKind::Other, err))?; + self.main_file.write_all(&encoded_bytes_len.to_le_bytes())?; + match file_utils::check_write(&self.main_file.write_all(&encoded_bytes)) { + Ok(()) => { + self.index_tree.insert(doc_id, self.main_offset); + let written_bytes = u32::try_from(encoded_bytes.len() + 4).map_err(|err| Error::new(ErrorKind::Other, err))?; + self.main_offset += written_bytes; + } + Err(err) => { + return Err(Error::new(ErrorKind::Other, format!("failed to write document: {} - {}", self.main_path.to_str().unwrap(), err))); } } Ok(()) diff --git a/src/repository/m3u_repository.rs b/src/repository/m3u_repository.rs index b6028de9d..e50f0cc16 100644 --- a/src/repository/m3u_repository.rs +++ b/src/repository/m3u_repository.rs @@ -21,24 +21,20 @@ macro_rules! cant_write_result { } } -fn m3u_get_base_file_path(target_path: &Path) -> Option { - Some(target_path.join(PathBuf::from("m3u.db"))) +fn m3u_get_base_file_path(target_path: &Path) -> PathBuf { + target_path.join(PathBuf::from("m3u.db")) } -pub(crate) fn m3u_get_file_paths(target_path: &Path) -> Option<(PathBuf, PathBuf)> { - match m3u_get_base_file_path(target_path) { - Some(m3u_path) => { - let extension = m3u_path.extension().map(|ext| format!("{}_", ext.to_str().unwrap_or(""))); - let index_path = m3u_path.with_extension(format!("{}idx", &extension.unwrap_or(String::new()))); - Some((m3u_path, index_path)) - } - None => None - } +pub(crate) fn m3u_get_file_paths(target_path: &Path) -> (PathBuf, PathBuf) { + let m3u_path = m3u_get_base_file_path(target_path); + let extension = m3u_path.extension().map(|ext| format!("{}_", ext.to_str().unwrap_or(""))); + let index_path = m3u_path.with_extension(format!("{}idx", &extension.unwrap_or_default())); + (m3u_path, index_path) } -pub(crate) fn m3u_get_epg_file_path(target_path: &Path) -> Option { - m3u_get_base_file_path(target_path) - .map(|path| file_utils::add_prefix_to_filename(&path, "epg_", Some("xml"))) +pub(crate) fn m3u_get_epg_file_path(target_path: &Path) -> PathBuf { + let path = m3u_get_base_file_path(target_path); + file_utils::add_prefix_to_filename(&path, "epg_", Some("xml")) } fn persist_m3u_playlist_as_text(target: &ConfigTarget, cfg: &Config, m3u_playlist: &Vec) { @@ -63,27 +59,26 @@ fn persist_m3u_playlist_as_text(target: &ConfigTarget, cfg: &Config, m3u_playlis pub(crate) fn m3u_write_playlist(target: &ConfigTarget, cfg: &Config, target_path: &Path, new_playlist: &[PlaylistGroup]) -> Result<(), M3uFilterError> { if !new_playlist.is_empty() { - if let Some((m3u_path, idx_path)) = m3u_get_file_paths(target_path) { - let m3u_playlist = new_playlist.iter() - .flat_map(|pg| &pg.channels) - .filter(|&pli| pli.header.borrow().item_type != PlaylistItemType::SeriesInfo) - .map(PlaylistItem::to_m3u).collect::>(); + let (m3u_path, idx_path) = m3u_get_file_paths(target_path); + let m3u_playlist = new_playlist.iter() + .flat_map(|pg| &pg.channels) + .filter(|&pli| pli.header.borrow().item_type != PlaylistItemType::SeriesInfo) + .map(PlaylistItem::to_m3u).collect::>(); - persist_m3u_playlist_as_text(target, cfg, &m3u_playlist); - { - let _file_lock = cfg.file_locks.write_lock(&m3u_path).map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Info, format!("{}", err)))?; - match IndexedDocumentWriter::new(m3u_path.clone(), idx_path) { - Ok(mut writer) => { - for m3u in m3u_playlist { - match writer.write_doc(m3u.virtual_id, &m3u) { - Ok(_) => {} - Err(err) => return Err(cant_write_result!(&m3u_path, err)) - } + persist_m3u_playlist_as_text(target, cfg, &m3u_playlist); + { + let _file_lock = cfg.file_locks.write_lock(&m3u_path).map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Info, format!("{err}")))?; + match IndexedDocumentWriter::new(m3u_path.clone(), idx_path) { + Ok(mut writer) => { + for m3u in m3u_playlist { + match writer.write_doc(m3u.virtual_id, &m3u) { + Ok(()) => {} + Err(err) => return Err(cant_write_result!(&m3u_path, err)) } - writer.flush().map_err(|err| cant_write_result!(&m3u_path, err))?; } - Err(err) => return Err(cant_write_result!(&m3u_path, err)) + writer.flush().map_err(|err| cant_write_result!(&m3u_path, err))?; } + Err(err) => return Err(cant_write_result!(&m3u_path, err)) } } } @@ -93,42 +88,39 @@ pub(crate) fn m3u_write_playlist(target: &ConfigTarget, cfg: &Config, target_pat pub(crate) fn m3u_load_rewrite_playlist(cfg: &Config, target: &ConfigTarget, user: &ProxyUserCredentials) -> Option { match ensure_target_storage_path(cfg, target.name.as_str()) { Ok(target_path) => { - if let Some((m3u_path, _)) = m3u_get_file_paths(&target_path) { - { - let _file_lock = cfg.file_locks.read_lock(&m3u_path).map_err(|err| { - error!("Could not lock document {:?}: {}", m3u_path, err); - Error::new(ErrorKind::Other, format!("Document Reader error for target {}", &target.name)) - }).ok()?; + let (m3u_path, _) = m3u_get_file_paths(&target_path); + { + let _file_lock = cfg.file_locks.read_lock(&m3u_path).map_err(|err| { + error!("Could not lock document {:?}: {}", m3u_path, err); + Error::new(ErrorKind::Other, format!("Document Reader error for target {}", &target.name)) + }).ok()?; - match IndexedDocumentReader::::new(&m3u_path) { - Ok(mut reader) => { - let server_info = get_user_server_info(cfg, user); - let url = format!("{}/m3u-stream/{}/{}", server_info.get_base_url(), user.username, user.password); - let mut result = vec![]; - result.push("#EXTM3U".to_string()); - for m3u_pli in reader.by_ref() { - match user.proxy { - ProxyType::Reverse => { - result.push(m3u_pli.to_m3u(target, Some(format!("{url}/{}", m3u_pli.virtual_id).as_str()))); - } - ProxyType::Redirect => { - result.push(m3u_pli.to_m3u(target, None)); - } + match IndexedDocumentReader::::new(&m3u_path) { + Ok(mut reader) => { + let server_info = get_user_server_info(cfg, user); + let url = format!("{}/m3u-stream/{}/{}", server_info.get_base_url(), user.username, user.password); + let mut result = vec![]; + result.push("#EXTM3U".to_string()); + for m3u_pli in reader.by_ref() { + match user.proxy { + ProxyType::Reverse => { + result.push(m3u_pli.to_m3u(target, Some(format!("{url}/{}", m3u_pli.virtual_id).as_str()))); + } + ProxyType::Redirect => { + result.push(m3u_pli.to_m3u(target, None)); } - }; - if reader.by_ref().has_error() { - error!("Could not deserialize m3u item {}", &m3u_path.to_str().unwrap()); - } else { - return Some(result.join("\n")); } - } - Err(err) => { - error!("Could not deserialize file {} - {}", &m3u_path.to_str().unwrap(), err); + }; + if reader.by_ref().has_error() { + error!("Could not deserialize m3u item {}", &m3u_path.to_str().unwrap()); + } else { + return Some(result.join("\n")); } } + Err(err) => { + error!("Could not deserialize file {} - {}", &m3u_path.to_str().unwrap(), err); + } } - } else { - error!("Could not open files for target {}", &target.name); } } Err(err) => { @@ -143,7 +135,7 @@ pub(crate) fn m3u_get_item_for_stream_id(cfg: &Config, stream_id: u32, m3u_path: return Err(Error::new(ErrorKind::Other, "id should start with 1")); } { - let _file_lock = cfg.file_locks.read_lock(&m3u_path)?; + let _file_lock = cfg.file_locks.read_lock(m3u_path)?; IndexedDocumentReader::::read_indexed_item(m3u_path, idx_path, stream_id) } } \ No newline at end of file diff --git a/src/repository/playlist_repository.rs b/src/repository/playlist_repository.rs index 298369eed..bd926f9a3 100644 --- a/src/repository/playlist_repository.rs +++ b/src/repository/playlist_repository.rs @@ -28,7 +28,7 @@ pub(crate) fn persist_playlist(playlist: &mut [PlaylistGroup], epg: Option<&Epg> Some(provider_id) => { let uuid = header.get_uuid(); let item_type = header.item_type; - header.virtual_id = target_id_mapping.insert_entry(**uuid, provider_id, &item_type, 0); + header.virtual_id = target_id_mapping.insert_entry(**uuid, provider_id, item_type, 0); } None => { errors.push(M3uFilterError::new(M3uFilterErrorKind::Info, format!("Playlistitem has no provider id: {}", &header.title))); diff --git a/src/repository/target_id_mapping.rs b/src/repository/target_id_mapping.rs index 5851462f4..f93125f05 100644 --- a/src/repository/target_id_mapping.rs +++ b/src/repository/target_id_mapping.rs @@ -26,7 +26,7 @@ pub(crate) struct VirtualIdRecord { impl VirtualIdRecord { fn new(provider_id: u32, virtual_id: u32, item_type: PlaylistItemType, parent_virtual_id: u32, uuid: [u8; 32]) -> Self { let last_updated = Local::now().timestamp(); - VirtualIdRecord { provider_id, virtual_id, uuid, item_type, parent_virtual_id, last_updated } + VirtualIdRecord { virtual_id, provider_id, uuid, item_type, parent_virtual_id, last_updated } } pub(crate) fn is_expired(&self) -> bool { @@ -48,7 +48,7 @@ pub(crate) struct TargetIdMapping { impl TargetIdMapping { pub(crate) fn new(path: &Path) -> Self { - let tree_virtual_id: BPlusTree = match BPlusTree::::deserialize(&path) { + let tree_virtual_id: BPlusTree = match BPlusTree::::deserialize(path) { Ok(tree) => tree, _ => BPlusTree::::new() }; @@ -61,9 +61,9 @@ impl TargetIdMapping { virtual_id_counter = max(virtual_id_counter, *max_value); } } - values.iter().for_each(|v| { + for v in values { tree_uuid.insert(v.uuid, v.virtual_id); - }); + } }); TargetIdMapping { dirty: false, @@ -74,12 +74,12 @@ impl TargetIdMapping { } } - pub(crate) fn insert_entry(&mut self, uuid: [u8; 32], provider_id: u32, item_type: &PlaylistItemType, parent_virtual_id: u32) -> u32 { + pub(crate) fn insert_entry(&mut self, uuid: [u8; 32], provider_id: u32, item_type: PlaylistItemType, parent_virtual_id: u32) -> u32 { match self.by_uuid.get(&uuid) { None => { self.dirty = true; self.virtual_id_counter += 1; - let record = VirtualIdRecord::new(provider_id, self.virtual_id_counter, *item_type, parent_virtual_id, uuid); + let record = VirtualIdRecord::new(provider_id, self.virtual_id_counter, item_type, parent_virtual_id, uuid); self.by_virtual_id.insert(self.virtual_id_counter, record); self.virtual_id_counter } @@ -98,11 +98,8 @@ impl TargetIdMapping { impl Drop for TargetIdMapping { fn drop(&mut self) { - match self.persist() { - Ok(_) => {} - Err(err) => { - error!("Failed to persist target id mapping {:?} err:{}", &self.path, err.to_string()) - } + if let Err(err) = self.persist() { + error!("Failed to persist target id mapping {:?} err:{err}", &self.path); } } } \ No newline at end of file diff --git a/src/repository/xtream_repository.rs b/src/repository/xtream_repository.rs index 387866ae8..076195c7c 100644 --- a/src/repository/xtream_repository.rs +++ b/src/repository/xtream_repository.rs @@ -57,16 +57,16 @@ fn xtream_get_info_file_paths(storage_path: &Path, cluster: XtreamCluster) -> Op fn write_playlists_to_file(cfg: &Config, storage_path: &Path, collections: Vec<(XtreamCluster, &mut [PlaylistItem])>) -> Result<(), M3uFilterError> { for (cluster, playlist) in collections { - let (xtream_path, idx_path) = xtream_get_file_paths(storage_path, &cluster); + let (xtream_path, idx_path) = xtream_get_file_paths(storage_path, cluster); { - let _file_lock = cfg.file_locks.write_lock(&xtream_path).map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Info, format!("{}", err)))?; + let _file_lock = cfg.file_locks.write_lock(&xtream_path).map_err(|err| M3uFilterError::new(M3uFilterErrorKind::Info, format!("{err}")))?; match IndexedDocumentWriter::new(xtream_path.clone(), idx_path) { Ok(mut writer) => { for item in playlist { match item.to_xtream() { Ok(xtream) => { match writer.write_doc(item.header.borrow().virtual_id, &xtream) { - Ok(_) => {} + Ok(()) => {} Err(err) => return Err(cant_write_result!(&xtream_path, err)) } } @@ -120,10 +120,7 @@ fn load_old_category_ids(path: &Path) -> (u32, HashMap) { } pub(crate) fn xtream_get_storage_path(cfg: &Config, target_name: &str) -> Option { - match get_target_storage_path(cfg, target_name) { - Some(target_path) => Some(target_path.join(std::path::PathBuf::from("xtream"))), - None => None, - } + get_target_storage_path(cfg, target_name).map(|target_path| target_path.join(PathBuf::from("xtream"))) } pub(crate) fn xtream_get_epg_file_path(path: &Path) -> PathBuf { @@ -131,13 +128,13 @@ pub(crate) fn xtream_get_epg_file_path(path: &Path) -> PathBuf { } fn xtream_get_file_paths_for_name(storage_path: &Path, name: &str) -> (PathBuf, PathBuf) { - let xtream_path = storage_path.join(format!("{}.db", name)); + let xtream_path = storage_path.join(format!("{name}.db")); let extension = xtream_path.extension().map(|ext| format!("{}_", ext.to_str().unwrap_or(""))); let index_path = xtream_path.with_extension(format!("{}idx", &extension.unwrap_or_default())); (xtream_path, index_path) } -pub(crate) fn xtream_get_file_paths(storage_path: &Path, cluster: &XtreamCluster) -> (PathBuf, PathBuf) { +pub(crate) fn xtream_get_file_paths(storage_path: &Path, cluster: XtreamCluster) -> (PathBuf, PathBuf) { xtream_get_file_paths_for_name(storage_path, &cluster.as_str().to_lowercase()) } @@ -182,22 +179,17 @@ pub(crate) fn xtream_write_playlist(target: &ConfigTarget, cfg: &Config, playlis // we skip resolved series, because this is only necessary when writing m3u files let col = if header.item_type == PlaylistItemType::Series { None + } else if header.get_provider_id().is_some() { + header.category_id = *cat_id; + Some(match header.xtream_cluster { + XtreamCluster::Live => &mut live_col, + XtreamCluster::Series => &mut series_col, + XtreamCluster::Video => &mut vod_col, + }) } else { - match header.get_provider_id() { - Some(_) => { - header.category_id = *cat_id; - Some(match header.xtream_cluster { - XtreamCluster::Live => &mut live_col, - XtreamCluster::Series => &mut series_col, - XtreamCluster::Video => &mut vod_col, - }) - } - None => { - let title = header.title.as_str(); - errors.push(format!("Channel does not have an id: {title}")); - None - } - } + let title = header.title.as_str(); + errors.push(format!("Channel does not have an id: {title}")); + None }; drop(header); if let Some(pl) = col { @@ -248,7 +240,7 @@ pub(crate) fn xtream_get_collection_path(cfg: &Config, target_name: &str, collec Err(Error::new(ErrorKind::Other, format!("Cant find collection: {target_name}/{collection_name}"))) } -fn xtream_read_item_for_stream_id(cfg: &Config, stream_id: u32, storage_path: &Path, cluster: &XtreamCluster) -> Result { +fn xtream_read_item_for_stream_id(cfg: &Config, stream_id: u32, storage_path: &Path, cluster: XtreamCluster) -> Result { let (xtream_path, idx_path) = xtream_get_file_paths(storage_path, cluster); { let _file_lock = cfg.file_locks.read_lock(&xtream_path)?; @@ -260,7 +252,7 @@ fn xtream_read_series_item_for_stream_id(cfg: &Config, stream_id: u32, storage_p let (xtream_path, idx_path) = xtream_get_file_paths_for_series(storage_path); { let _file_lock = cfg.file_locks.read_lock(&xtream_path)?; - return IndexedDocumentReader::::read_indexed_item(&xtream_path, &idx_path, stream_id); + IndexedDocumentReader::::read_indexed_item(&xtream_path, &idx_path, stream_id) } } @@ -284,10 +276,10 @@ pub(crate) 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) - .map_err(|err| Error::new(ErrorKind::Other, format!("Could not get lock for id mapping for target {} err:{}", target.name, err.to_string())))?; + .map_err(|err| Error::new(ErrorKind::Other, format!("Could not get lock for id mapping for target {} err:{err}", target.name)))?; let mut target_id_mapping = BPlusTreeQuery::::try_new(&target_id_mapping_file) - .map_err(|err| Error::new(ErrorKind::Other, format!("Could not load id mapping for target {} err:{}", target.name, err.to_string())))?; + .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) @@ -302,20 +294,20 @@ pub(crate) fn xtream_get_item_for_stream_id( } PlaylistItemType::Catchup => { let cluster = try_cluster!(xtream_cluster, mapping.item_type, virtual_id)?; - let mut item = xtream_read_item_for_stream_id(config, mapping.parent_virtual_id, &storage_path, &cluster)?; + let mut item = xtream_read_item_for_stream_id(config, mapping.parent_virtual_id, &storage_path, cluster)?; item.provider_id = mapping.provider_id; Ok(item) } _ => { let cluster = try_cluster!(xtream_cluster, mapping.item_type, virtual_id)?; - xtream_read_item_for_stream_id(config, virtual_id, &storage_path, &cluster) + xtream_read_item_for_stream_id(config, virtual_id, &storage_path, cluster) } } } } -pub(crate) fn xtream_load_rewrite_playlist(cluster: &XtreamCluster, config: &Config, target: &ConfigTarget, category_id: u32) -> Result { +pub(crate) fn xtream_load_rewrite_playlist(cluster: XtreamCluster, config: &Config, target: &ConfigTarget, category_id: u32) -> Result { if let Some(storage_path) = xtream_get_storage_path(config, target.name.as_str()) { let (xtream_path, _) = xtream_get_file_paths(&storage_path, cluster); { @@ -445,7 +437,7 @@ pub(crate) fn write_and_get_xtream_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) - .map_err(|err| Error::new(ErrorKind::Other, format!("Could not load id mapping for target {} err:{}", target.name, err.to_string())))?; + .map_err(|err| Error::new(ErrorKind::Other, format!("Could not load id mapping for target {} err:{err}", target.name)))?; let mut target_id_mapping = TargetIdMapping::new(&target_id_mapping_file); let options = XtreamMappingOptions::from_target_options(target.options.as_ref()); @@ -453,7 +445,7 @@ pub(crate) fn write_and_get_xtream_series_info( for episode in episode_list.iter_mut().filter_map(Value::as_object_mut) { if let Some(provider_id) = episode.get("id").and_then(Value::as_str).and_then(|id| id.parse::().ok()) { let uuid = hash_string(&format!("{}/{}", pli_series_info.url, provider_id)); - let virtual_id = target_id_mapping.insert_entry(uuid, provider_id, &PlaylistItemType::SeriesEpisode, pli_series_info.virtual_id); + let virtual_id = target_id_mapping.insert_entry(uuid, provider_id, PlaylistItemType::SeriesEpisode, pli_series_info.virtual_id); episode.insert("id".to_string(), Value::String(virtual_id.to_string())); } if options.skip_series_direct_source { diff --git a/src/utils/config_reader.rs b/src/utils/config_reader.rs index 3bef9b1e1..f1e27e696 100644 --- a/src/utils/config_reader.rs +++ b/src/utils/config_reader.rs @@ -21,7 +21,7 @@ pub(crate) fn read_mappings(args_mapping: Option, cfg: &mut Config) -> R Ok(mappings) => { match mappings { None => {debug!("no mapping loaded");} - Some(mappings_cfg) => {cfg.set_mappings(mappings_cfg);} + Some(mappings_cfg) => {cfg.set_mappings(&mappings_cfg);} } Ok(()) } diff --git a/src/utils/download.rs b/src/utils/download.rs index a5ddff2de..7df8d46a7 100644 --- a/src/utils/download.rs +++ b/src/utils/download.rs @@ -56,7 +56,7 @@ pub(crate) async fn get_xtream_playlist_series<'a>(fpl: &mut FetchedPlaylist<'a> 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.extend(series.drain(..)); + group_series.append(&mut series); } } Err(err) => errors.push(err), @@ -133,7 +133,7 @@ pub(crate) async fn get_xtream_playlist(input: &ConfigInput, working_dir: &Strin &stream_content) { Ok(sub_playlist_parsed) => { if let Some(mut xtream_sub_playlist) = sub_playlist_parsed { - playlist_groups.extend(xtream_sub_playlist.drain(..)); + playlist_groups.append(&mut xtream_sub_playlist); } } Err(err) => errors.push(err) @@ -142,8 +142,7 @@ pub(crate) async fn get_xtream_playlist(input: &ConfigInput, working_dir: &Strin (Err(err1), Err(err2)) => { errors.extend([err1, err2]); }, - (Err(err), _) => errors.push(err), - (_, Err(err)) => errors.push(err), + (_, Err(err)) | (Err(err), _) => errors.push(err), } } } diff --git a/src/utils/request_utils.rs b/src/utils/request_utils.rs index 9d13c5c1a..ee12ed3d7 100644 --- a/src/utils/request_utils.rs +++ b/src/utils/request_utils.rs @@ -110,7 +110,7 @@ pub(crate) fn get_request_headers(defined_headers: &HashMap, cus fn get_local_file_content(file_path: &PathBuf) -> Result { // Check if the file is accessible if file_path.exists() && file_path.is_file() { - if let Ok(content) = fs::read(&file_path) { + if let Ok(content) = fs::read(file_path) { if content.len() >= 2 && is_gzip(&content[0..2]) { let mut decoder = GzDecoder::new(&content[..]); let mut decode_buffer = String::new(); diff --git a/test/rest-api.http b/test/rest-api.http index e7a59df00..dc3e6eaf6 100644 --- a/test/rest-api.http +++ b/test/rest-api.http @@ -29,7 +29,7 @@ GET http://localhost:8901/player_api.php?action=get_series_categories&username=x GET http://localhost:8901/player_api.php?action=get_series&username=xt&password=xt&category_id=120 ### xtream series info -GET http://localhost:8901/player_api.php?action=get_series_info&username=xt&password=xt&series_id=16441 +GET http://localhost:8901/player_api.php?action=get_series_info&username=xt&password=xt&series_id=16467 ### xtream series stream - GET http://localhost:8901/series/xt/xt/18130 \ No newline at end of file + GET http://localhost:8901/series/xt/xt/18137 \ No newline at end of file