diff --git a/src/api/xtream_api.rs b/src/api/xtream_api.rs index 309f78f3f..c52aac87c 100644 --- a/src/api/xtream_api.rs +++ b/src/api/xtream_api.rs @@ -17,7 +17,7 @@ use crate::api::api_utils::{get_user_server_info, get_user_target, get_user_targ use crate::model::api_proxy::{ProxyType, ProxyUserCredentials}; use crate::model::config::{Config, ConfigInput, ConfigTarget}; use crate::model::config::TargetType; -use crate::model::playlist::{PlaylistItem, XtreamCluster}; +use crate::model::playlist::{XtreamCluster, XtreamPlaylistItem}; use crate::model::xtream::XtreamMappingOptions; use crate::repository::xtream_repository; use crate::repository::xtream_repository::get_xtream_item_for_stream_id; @@ -168,13 +168,10 @@ async fn xtream_player_api_stream( match get_xtream_item_for_stream_id(req_stream_id, &app_state.config, target) { Ok(pli) => { - let input_id: u16 = match FromStr::from_str(pli.header.borrow().source.as_str()) { - Ok(id) => id, - Err(_) => return HttpResponse::BadRequest().finish() - }; + let input_id: u16 = pli.input_id; if let Some(input) = app_state.config.get_input_by_id(&input_id) { let mut query_path = if stream_req.action_path.is_empty() { "".to_string() } else { format!("{}/", stream_req.action_path) }; - query_path = format!("{}{}{}", query_path, pli.header.borrow().id, stream_ext); + query_path = format!("{}{}{}", query_path, pli.id, stream_ext); if let Some(stream_url) = get_xtream_player_api_stream_url(input, stream_req.context.to_string().as_str(), query_path.as_str()) { if user.proxy == ProxyType::Redirect { debug!("Redirecting stream request to {}", stream_url); @@ -250,7 +247,7 @@ async fn xtream_player_api_timeshift_stream( xtream_player_api_stream(&req, &api_req, &app_state, XtreamApiStreamRequest::from(XtreamApiStreamContext::Timeshift, &username, &password, &stream_id, &action_path)).await } -async fn xtream_get_stream_info(input: &ConfigInput, target: &ConfigTarget, pli: &PlaylistItem, +async fn xtream_get_stream_info(input: &ConfigInput, target: &ConfigTarget, pli: &XtreamPlaylistItem, info_url: &str, cluster: &XtreamCluster) -> Result { if let Ok(url) = Url::parse(info_url) { let client = request_utils::get_client_request(Some(input), url, None); @@ -264,15 +261,15 @@ async fn xtream_get_stream_info(input: &ConfigInput, target: &ConfigTarget, pli: XtreamCluster::Video => { if let Ok(mut doc) = serde_json::from_str::>(content.as_str()) { if let Some(Value::Object(movie_data) ) = doc.get_mut("movie_data") { - let stream_id = pli.header.borrow().id.parse::().unwrap(); - let category_id = pli.header.borrow().category_id; + let stream_id = pli.id; + let category_id = pli.category_id; movie_data.insert("stream_id".to_string(), Value::Number(serde_json::value::Number::from(stream_id))); movie_data.insert("category_id".to_string(), Value::Number(serde_json::value::Number::from(category_id))); let options = XtreamMappingOptions::from_target_options(target.options.as_ref()); if options.skip_video_direct_source { movie_data.insert("direct_source".to_string(), Value::String("".to_string())); } else { - movie_data.insert("direct_source".to_string(), Value::String(pli.header.borrow().url.to_string())); + movie_data.insert("direct_source".to_string(), Value::String(pli.url.to_string())); } if let Ok(result) = serde_json::to_string(&doc) { return Ok(result); @@ -296,7 +293,7 @@ async fn xtream_get_stream_info(input: &ConfigInput, target: &ConfigTarget, pli: } } Err(Error::new(std::io::ErrorKind::Other, format!("Cant find stream with id: {}/{}/{}", - target.name.as_str(), &cluster, pli.header.borrow().stream_id.as_ref()))) + target.name.as_str(), &cluster, pli.stream_id))) } async fn xtream_get_stream_info_response(app_state: &AppState, user: &ProxyUserCredentials, @@ -308,9 +305,9 @@ async fn xtream_get_stream_info_response(app_state: &AppState, user: &ProxyUserC }; if let Ok(pli) = get_xtream_item_for_stream_id(req_stream_id, &app_state.config, target) { - let input_id = pli.header.borrow().source.parse::().unwrap(); + let input_id = pli.input_id; if let Some(input) = app_state.config.get_input_by_id(&input_id) { - let stream_id = pli.header.borrow().id.parse::().unwrap(); + let stream_id = pli.id; if let Some(info_url) = get_xtream_player_api_info_url(input, cluster, stream_id) { if user.proxy == ProxyType::Redirect { return HttpResponse::Found().insert_header(("Location", info_url)).finish(); @@ -332,13 +329,10 @@ async fn xtream_get_short_epg(app_state: &AppState, user: &ProxyUserCredentials, }; if let Ok(pli) = get_xtream_item_for_stream_id(req_stream_id, &app_state.config, target) { - let input_id: u16 = match FromStr::from_str(pli.header.borrow().source.as_str()) { - Ok(id) => id, - Err(_) => return HttpResponse::BadRequest().finish() - }; + let input_id: u16 = pli.input_id; if let Some(input) = app_state.config.get_input_by_id(&input_id) { if let Some(action_url) = get_xtream_player_api_action_url(input, "get_short_epg") { - let mut info_url = format!("{}&stream_id={}", action_url, pli.header.borrow().id.as_str()); + let mut info_url = format!("{}&stream_id={}", action_url, pli.id); if !(limit.is_empty() || limit.eq("0")) { info_url = format!("{}&limit={}", info_url, limit); } @@ -425,13 +419,15 @@ async fn xtream_player_api( api_req.limit.trim()).await } _ => { - match xtream_player_api_handle_content_action(&_app_state.config, target_name, action, api_req.category_id.as_str(), req).await { + let category_id = api_req.category_id.as_str().trim(); + match xtream_player_api_handle_content_action(&_app_state.config, target_name, action, category_id, req).await { Some(response) => response, _ => { + let cat_id = if category_id.is_empty() { 0 } else { category_id.parse::().unwrap_or(0) }; match match action { - "get_live_streams" => xtream_repository::load_rewrite_xtream_playlist(&XtreamCluster::Live, &_app_state.config, target), - "get_vod_streams" => xtream_repository::load_rewrite_xtream_playlist(&XtreamCluster::Video, &_app_state.config, target), - "get_series" => xtream_repository::load_rewrite_xtream_playlist(&XtreamCluster::Series, &_app_state.config, target), + "get_live_streams" => xtream_repository::load_rewrite_xtream_playlist(&XtreamCluster::Live, &_app_state.config, target, cat_id), + "get_vod_streams" => xtream_repository::load_rewrite_xtream_playlist(&XtreamCluster::Video, &_app_state.config, target, cat_id), + "get_series" => xtream_repository::load_rewrite_xtream_playlist(&XtreamCluster::Series, &_app_state.config, target, cat_id), _ => Err(Error::new(ErrorKind::Unsupported, format!("Cant find action: {}/{}", target_name, action))), } { Ok(payload) => HttpResponse::Ok().body(payload), diff --git a/src/model/playlist.rs b/src/model/playlist.rs index 101a9ae14..ff980bc13 100644 --- a/src/model/playlist.rs +++ b/src/model/playlist.rs @@ -50,7 +50,7 @@ impl Display for XtreamCluster { pub(crate) fn default_stream_cluster() -> XtreamCluster { XtreamCluster::Live } -#[derive(Debug, Clone, PartialEq)] +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] pub(crate) enum PlaylistItemType { Live = 1, Movie = 2, @@ -60,6 +60,7 @@ pub(crate) enum PlaylistItemType { pub(crate) fn default_playlist_item_type() -> PlaylistItemType { PlaylistItemType::Live } fn default_as_zero_u32() -> u32 { 0 } +fn default_as_zero_u16() -> u16 { 0 } pub(crate) trait FieldAccessor { fn get_field(&self, field: &str) -> Option>; @@ -80,19 +81,19 @@ pub(crate) struct PlaylistItemHeader { pub audio_track: Rc, pub time_shift: Rc, pub rec: Rc, - pub source: Rc, pub url: Rc, pub epg_channel_id: Option>, #[serde(default = "default_stream_cluster")] pub xtream_cluster: XtreamCluster, - pub additional_properties: Option>, + pub additional_properties: Option, #[serde(default = "default_playlist_item_type", skip_serializing, skip_deserializing)] pub item_type: PlaylistItemType, #[serde(default = "default_as_false", skip_serializing, skip_deserializing)] pub series_fetched: bool, // only used for series_info #[serde(default = "default_as_zero_u32")] pub category_id: u32, - + #[serde(default = "default_as_zero_u16")] + pub input_id: u16, } macro_rules! update_fields { @@ -133,7 +134,7 @@ macro_rules! to_m3u_non_empty_fields { impl FieldAccessor for PlaylistItemHeader { fn get_field(&self, field: &str) -> Option> { - let val = get_fields!(self, field, id, stream_id, name, logo, logo_small, group, title, parent_code, audio_track, time_shift, rec, source, url;); + let val = get_fields!(self, field, id, stream_id, name, logo, logo_small, group, title, parent_code, audio_track, time_shift, rec, url;); if val.is_some() { return val; } @@ -145,7 +146,7 @@ impl FieldAccessor for PlaylistItemHeader { fn set_field(&mut self, field: &str, value: &str) -> bool { let val = String::from(value); - let updated = update_fields!(self, field, id, stream_id, name, logo, logo_small, group, title, parent_code, audio_track, time_shift, rec, source, url; val); + let updated = update_fields!(self, field, id, stream_id, name, logo, logo_small, group, title, parent_code, audio_track, time_shift, rec, url; val); if updated { return updated; } @@ -174,6 +175,7 @@ pub(crate) struct M3uPlaylistItem { pub rec: Rc, pub url: Rc, pub epg_channel_id: Option>, + pub input_id: u16, } @@ -201,6 +203,33 @@ impl M3uPlaylistItem { } } +#[derive(Debug, Clone, Serialize, Deserialize)] +pub(crate) struct XtreamPlaylistItem { + pub stream_id: u32, + pub id: u32, + pub name: Rc, + pub logo: Rc, + pub logo_small: Rc, + pub group: Rc, + pub title: Rc, + pub parent_code: Rc, + pub rec: Rc, + pub url: Rc, + pub epg_channel_id: Option>, + pub xtream_cluster: XtreamCluster, + pub additional_properties: Option, + pub item_type: PlaylistItemType, + pub series_fetched: bool, // only used for series_info + pub category_id: u32, + pub input_id: u16, +} + +impl XtreamPlaylistItem { + pub fn to_doc(&self, options: &XtreamMappingOptions) -> Value { + xtream_playlistitem_to_document(self, options) + } +} + #[derive(Debug, Clone, Serialize, Deserialize)] pub(crate) struct PlaylistItem { pub header: RefCell, @@ -223,11 +252,37 @@ impl PlaylistItem { rec: Rc::clone(&header.rec), url: Rc::clone(&header.url), epg_channel_id: header.epg_channel_id.clone(), + input_id: header.input_id, } } - pub fn to_xtream(&self, options: &XtreamMappingOptions) -> Value { - xtream_playlistitem_to_document(self, options) + pub fn to_xtream(&self) -> XtreamPlaylistItem { + let header = self.header.borrow(); + XtreamPlaylistItem { + stream_id: header.stream_id.parse::().unwrap(), + id: header.id.parse::().unwrap(), + name: Rc::clone(&header.name), + logo: Rc::clone(&header.logo), + logo_small: Rc::clone(&header.logo_small), + group: Rc::clone(&header.group), + title: Rc::clone(&header.title), + parent_code: Rc::clone(&header.parent_code), + rec: Rc::clone(&header.rec), + url: Rc::clone(&header.url), + epg_channel_id: header.epg_channel_id.clone(), + xtream_cluster: header.xtream_cluster.clone(), + additional_properties: match &header.additional_properties { + None => None, + Some(props) => match serde_json::to_string(props) { + Ok(val) => Some(val), + Err(_) => None + } + }, + item_type: header.item_type.clone(), + series_fetched: header.series_fetched, + category_id: header.category_id, + input_id: header.input_id + } } } diff --git a/src/model/xtream.rs b/src/model/xtream.rs index 7e77cc57f..b7059e507 100644 --- a/src/model/xtream.rs +++ b/src/model/xtream.rs @@ -1,14 +1,13 @@ -use std::cell::Ref; use std::collections::HashMap; use std::iter::{FromIterator}; use std::rc::Rc; use serde::{Deserialize, Deserializer, Serialize}; use serde::de::DeserializeOwned; -use serde_json::Value; +use serde_json::{Map, Value}; use crate::model::config::{ConfigTargetOptions, default_as_empty_rc_str}; -use crate::model::playlist::{PlaylistItem, PlaylistItemHeader, XtreamCluster}; +use crate::model::playlist::{PlaylistItem, XtreamCluster, XtreamPlaylistItem}; const LIVE_STREAM_FIELDS: &[&str] = &[]; @@ -186,36 +185,36 @@ pub(crate) struct XtreamStream { macro_rules! add_str_property_if_exists { ($vec:expr, $prop:expr, $prop_name:expr) => { - $vec.push((String::from($prop_name), Value::String($prop.to_string()))); + $vec.insert(String::from($prop_name), Value::String($prop.to_string())); } } macro_rules! add_rc_str_property_if_exists { ($vec:expr, $prop:expr, $prop_name:expr) => { - $prop.as_ref().map(|v| $vec.push((String::from($prop_name), Value::String(v.to_string())))); + $prop.as_ref().map(|v| $vec.insert(String::from($prop_name), Value::String(v.to_string()))); } } macro_rules! add_opt_i64_property_if_exists { ($vec:expr, $prop:expr, $prop_name:expr) => { - $prop.as_ref().map(|v| $vec.push((String::from($prop_name), Value::Number(serde_json::value::Number::from(i64::from(*v)))))); + $prop.as_ref().map(|v| $vec.insert(String::from($prop_name), Value::Number(serde_json::value::Number::from(i64::from(*v))))); } } macro_rules! add_opt_f64_property_if_exists { ($vec:expr, $prop:expr, $prop_name:expr) => { - $prop.as_ref().map(|v| $vec.push((String::from($prop_name), Value::Number(serde_json::value::Number::from_f64(f64::from(*v)).unwrap())))); + $prop.as_ref().map(|v| $vec.insert(String::from($prop_name), Value::Number(serde_json::value::Number::from_f64(f64::from(*v)).unwrap()))); } } macro_rules! add_f64_property_if_exists { ($vec:expr, $prop:expr, $prop_name:expr) => { - $vec.push((String::from($prop_name), Value::Number(serde_json::value::Number::from_f64(f64::from($prop)).unwrap()))); + $vec.insert(String::from($prop_name), Value::Number(serde_json::value::Number::from_f64(f64::from($prop)).unwrap())); } } macro_rules! add_i64_property_if_exists { ($vec:expr, $prop:expr, $prop_name:expr) => { - $vec.push((String::from($prop_name), Value::Number(serde_json::value::Number::from(i64::from($prop))))); + $vec.insert(String::from($prop_name), Value::Number(serde_json::value::Number::from(i64::from($prop)))); } } @@ -225,11 +224,11 @@ impl XtreamStream { self.stream_id.map_or_else(|| self.series_id.map_or_else(|| String::from(""), |seid| format!("{}", seid)), |sid| format!("{}", sid)) } - pub(crate) fn get_additional_properties(&self) -> Option> { - let mut result = vec![]; + pub(crate) fn get_additional_properties(&self) -> Option { + let mut result = Map::new(); if let Some(bdpath) = self.backdrop_path.as_ref() { if !bdpath.is_empty() { - result.push((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().unwrap()))]))); } } add_rc_str_property_if_exists!(result, self.added, "added"); @@ -251,7 +250,7 @@ impl XtreamStream { add_rc_str_property_if_exists!(result, self.epg_channel_id, "epg_channel_id"); add_opt_i64_property_if_exists!(result, self.tv_archive, "tv_archive"); add_opt_i64_property_if_exists!(result, self.tv_archive_duration, "tv_archive_duration"); - if result.is_empty() { None } else { Some(result) } + if result.is_empty() { None } else { Some(Value::Object(result)) } } } @@ -326,11 +325,11 @@ pub(crate) struct XtreamSeriesInfo { impl XtreamSeriesInfoEpisode { - pub(crate) fn get_additional_properties(&self, series_info: &XtreamSeriesInfo) -> Option> { - let mut result = vec![]; + pub(crate) fn get_additional_properties(&self, series_info: &XtreamSeriesInfo) -> Option { + let mut result = Map::new(); let bdpath = &series_info.info.backdrop_path; if !bdpath.is_empty() { - result.push((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().unwrap()))]))); } add_str_property_if_exists!(result, self.added.as_str(), "added"); add_str_property_if_exists!(result, series_info.info.cast.as_str(), "cast"); @@ -346,7 +345,7 @@ impl XtreamSeriesInfoEpisode { add_str_property_if_exists!(result, self.title, "title"); add_i64_property_if_exists!(result, self.season, "season"); add_str_property_if_exists!(result, series_info.info.youtube_trailer, "youtube_trailer"); - if result.is_empty() { None } else { Some(result) } + if result.is_empty() { None } else { Some(Value::Object(result)) } } } @@ -359,8 +358,9 @@ impl XtreamMappingOptions { pub fn from_target_options(options: Option<&ConfigTargetOptions>) -> Self { let (skip_live_direct_source, skip_video_direct_source) = options .map_or((false, false), |o| (o.xtream_skip_live_direct_source, o.xtream_skip_video_direct_source)); - XtreamMappingOptions{ - skip_live_direct_source, skip_video_direct_source + XtreamMappingOptions { + skip_live_direct_source, + skip_video_direct_source, } } } @@ -394,46 +394,48 @@ fn append_mandatory_fields(document: &mut serde_json::Map, fields } } -fn append_prepared_series_properties(header: &Ref, document: &mut serde_json::Map) { - if let Some(add_props) = &header.additional_properties { - match add_props.iter().find(|(key, _)| key.eq("rating")) { - Some((_, value)) => { - document.insert("rating".to_string(), match value { - Value::Number(val) => Value::String(format!("{:.0}", val.as_f64().unwrap())), - Value::String(val) => Value::String(val.to_string()), - _ => Value::String("0".to_string()), - }); - } - None => { - document.insert("rating".to_string(), Value::String("0".to_string())); +fn append_prepared_series_properties(add_props: Option<&Map>, document: &mut Map) { + match add_props { + Some(props) => { + match props.get("rating") { + Some(value) => { + document.insert("rating".to_string(), match value { + Value::Number(val) => Value::String(format!("{:.0}", val.as_f64().unwrap())), + Value::String(val) => Value::String(val.to_string()), + _ => Value::String("0".to_string()), + }); + } + None => { + document.insert("rating".to_string(), Value::String("0".to_string())); + } } } + None => {} } } -pub(crate) fn xtream_playlistitem_to_document(pli: &PlaylistItem, options: &XtreamMappingOptions) -> serde_json::Value { - let header = &pli.header.borrow(); - let stream_id_value = Value::Number(serde_json::Number::from(header.stream_id.parse::().unwrap())); +pub(crate) fn xtream_playlistitem_to_document(pli: &XtreamPlaylistItem, options: &XtreamMappingOptions) -> serde_json::Value { + let stream_id_value = Value::Number(serde_json::Number::from(pli.stream_id)); let mut document = serde_json::Map::from_iter([ - ("category_id".to_string(), Value::String(format!("{}", &header.category_id))), - ("category_ids".to_string(), Value::Array(Vec::from([Value::Number(serde_json::Number::from(header.category_id))]))), - ("name".to_string(), Value::String(header.name.as_ref().clone())), + ("category_id".to_string(), Value::String(format!("{}", &pli.category_id))), + ("category_ids".to_string(), Value::Array(Vec::from([Value::Number(serde_json::Number::from(pli.category_id))]))), + ("name".to_string(), Value::String(pli.name.as_ref().clone())), ("num".to_string(), stream_id_value.clone()), - ("title".to_string(), Value::String(header.title.as_ref().clone())), - ("stream_icon".to_string(), Value::String(header.logo.as_ref().clone())), + ("title".to_string(), Value::String(pli.title.as_ref().clone())), + ("stream_icon".to_string(), Value::String(pli.logo.as_ref().clone())), ]); - match header.xtream_cluster { + match pli.xtream_cluster { XtreamCluster::Live => { document.insert("stream_id".to_string(), stream_id_value); if options.skip_live_direct_source { document.insert("direct_source".to_string(), Value::String("".to_string())); } else { - document.insert("direct_source".to_string(), Value::String(header.url.as_ref().clone())); + document.insert("direct_source".to_string(), Value::String(pli.url.as_ref().clone())); } - document.insert("thumbnail".to_string(), Value::String(header.logo_small.as_ref().clone())); + document.insert("thumbnail".to_string(), Value::String(pli.logo_small.as_ref().clone())); document.insert("custom_sid".to_string(), Value::String("".to_string())); - document.insert("epg_channel_id".to_string(), match &header.epg_channel_id { + document.insert("epg_channel_id".to_string(), match &pli.epg_channel_id { None => Value::Null, Some(epg_id) => Value::String(epg_id.as_ref().clone()) }); @@ -443,7 +445,7 @@ pub(crate) fn xtream_playlistitem_to_document(pli: &PlaylistItem, options: &Xtre if options.skip_video_direct_source { document.insert("direct_source".to_string(), Value::String("".to_string())); } else { - document.insert("direct_source".to_string(), Value::String(header.url.as_ref().clone())); + document.insert("direct_source".to_string(), Value::String(pli.url.as_ref().clone())); } document.insert("custom_sid".to_string(), Value::String("".to_string())); } @@ -452,13 +454,23 @@ pub(crate) fn xtream_playlistitem_to_document(pli: &PlaylistItem, options: &Xtre } }; - if let Some(add_props) = &header.additional_properties { + let props = match &pli.additional_properties { + None => None, + Some(add_props) => { + match serde_json::from_str::>(add_props) { + Ok(p) => Some(p), + Err(_) => None, + } + } + }; + + if let Some(ref add_props) = props { for (field_name, field_value) in add_props { document.insert(field_name.to_string(), field_value.to_owned()); } } - match header.xtream_cluster { + match pli.xtream_cluster { XtreamCluster::Live => { append_mandatory_fields(&mut document, LIVE_STREAM_FIELDS); } @@ -466,10 +478,10 @@ pub(crate) fn xtream_playlistitem_to_document(pli: &PlaylistItem, options: &Xtre append_mandatory_fields(&mut document, VIDEO_STREAM_FIELDS); } XtreamCluster::Series => { - append_prepared_series_properties(header, &mut document); + append_prepared_series_properties(props.as_ref(), &mut document); append_mandatory_fields(&mut document, SERIES_STREAM_FIELDS); append_release_date(&mut document); } }; Value::Object(document) -} \ No newline at end of file +} diff --git a/src/processing/m3u_parser.rs b/src/processing/m3u_parser.rs index 5ed254a6d..34ae22306 100644 --- a/src/processing/m3u_parser.rs +++ b/src/processing/m3u_parser.rs @@ -66,7 +66,6 @@ fn create_empty_playlistitem_header(input_id: u16, url: String) -> PlaylistItemH audio_track: default_as_empty_rc_str(), time_shift: default_as_empty_rc_str(), rec: default_as_empty_rc_str(), - source: Rc::new(input_id.to_string()),//Rc::new(content.to_owned()), url: Rc::new(url), epg_channel_id: None, item_type: default_playlist_item_type(), @@ -74,6 +73,7 @@ fn create_empty_playlistitem_header(input_id: u16, url: String) -> PlaylistItemH additional_properties: None, series_fetched: false, category_id: 0, + input_id } } diff --git a/src/processing/playlist_processor.rs b/src/processing/playlist_processor.rs index d34b0d974..b85d958cd 100644 --- a/src/processing/playlist_processor.rs +++ b/src/processing/playlist_processor.rs @@ -281,7 +281,7 @@ async fn process_source(cfg: Arc, source_idx: usize, user_targets: Arc

Result>, M3uFilterError> { match map_to_xtream_category(category) { Ok(mut categories) => { - let input_id = input.id.to_string(); + let input_id = input.id; let url = input.url.as_str(); let username = input.username.as_ref().map_or("", |v| v); let password = input.password.as_ref().map_or("", |v| v); @@ -109,7 +109,6 @@ pub(crate) fn parse_xtream(input: &ConfigInput, audio_track: default_as_empty_rc_str(), time_shift: default_as_empty_rc_str(), rec: default_as_empty_rc_str(), - source: Rc::new(input_id.to_owned()), url: if stream.direct_source.is_empty() { let stream_base_url = match xtream_cluster { XtreamCluster::Live => format!("{}/live/{}/{}/{}.ts", url, username, password, &stream.get_stream_id()), @@ -135,6 +134,7 @@ pub(crate) fn parse_xtream(input: &ConfigInput, additional_properties: stream.get_additional_properties(), series_fetched: false, category_id: 0, + input_id }), }; grp.add(item); diff --git a/src/repository/m3u_repository.rs b/src/repository/m3u_repository.rs index 3169e375d..35fb12aac 100644 --- a/src/repository/m3u_repository.rs +++ b/src/repository/m3u_repository.rs @@ -99,9 +99,9 @@ pub(crate) fn load_rewrite_m3u_playlist(cfg: &Config, target: &ConfigTarget, use result.push("#EXTM3U".to_string()); while cursor < size { - let tuple = IndexRecord::from_bytes(&encoded_idx, &mut cursor); - let start_offset = tuple.index as usize; - let end_offset = start_offset + tuple.size as usize; + let index_record = IndexRecord::from_bytes(&encoded_idx, &mut cursor); + let start_offset = index_record.index as usize; + let end_offset = start_offset + index_record.size as usize; if let Ok(m3u_pli) = bincode::deserialize::(&encoded_m3u[start_offset..end_offset]) { match user.proxy { ProxyType::Reverse => { diff --git a/src/repository/repository_utils.rs b/src/repository/repository_utils.rs index ab4f8b979..5c5b0e384 100644 --- a/src/repository/repository_utils.rs +++ b/src/repository/repository_utils.rs @@ -35,7 +35,7 @@ impl IndexRecord { let index_bytes: [u8; 4] = bytes[*cursor..*cursor + 4].try_into().unwrap(); *cursor += 4; let size_bytes: [u8; 2] = bytes[*cursor..*cursor + 2].try_into().unwrap(); - *cursor += 4; + *cursor += 2; let index = u32::from_le_bytes(index_bytes); let size = u16::from_le_bytes(size_bytes); IndexRecord { index, size } diff --git a/src/repository/xtream_repository.rs b/src/repository/xtream_repository.rs index 182752b73..d58ebe1dc 100644 --- a/src/repository/xtream_repository.rs +++ b/src/repository/xtream_repository.rs @@ -8,7 +8,7 @@ use log::{error}; use serde::Serialize; use serde_json::{json, Value}; use crate::model::config::{Config, ConfigTarget}; -use crate::model::playlist::{PlaylistGroup, PlaylistItem, PlaylistItemType, XtreamCluster}; +use crate::model::playlist::{PlaylistGroup, PlaylistItem, PlaylistItemType, XtreamCluster, XtreamPlaylistItem}; use crate::{create_m3u_filter_error_result}; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::xtream::{XtreamMappingOptions}; @@ -94,7 +94,7 @@ fn write_playlist_to_file(storage_path: &Path, stream_id: &mut u32, cluster: &Xt let mut idx_offset: u32 = 0; for pli in playlist.iter_mut() { pli.header.borrow_mut().stream_id = Rc::new(stream_id.to_string()); - if let Ok(encoded) = bincode::serialize(&pli) { + if let Ok(encoded) = bincode::serialize(&pli.to_xtream()) { match file_utils::check_write(main_file.write_all(&encoded)) { Ok(_) => { let bytes_written = encoded.len() as u16; @@ -144,7 +144,7 @@ fn write_playlists_to_file(storage_path: &Path, collections: Vec<(XtreamCluster, } match save_stream_id_cluster_mapping(storage_path, &mut id_list) { Ok(_) => Ok(()), - Err(err) => Err(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("failed to write xtream playlist: {} - {}", storage_path.to_str().unwrap() , err))) + Err(err) => Err(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("failed to write xtream playlist: {} - {}", storage_path.to_str().unwrap(), err))) } } @@ -186,7 +186,7 @@ pub(crate) fn write_xtream_playlist(target: &ConfigTarget, cfg: &Config, playlis let col = if header.item_type != PlaylistItemType::Series { if header.id.parse::().is_ok() { header.category_id = *cat_id; - Some(match pli.header.borrow().xtream_cluster { + Some(match header.xtream_cluster { XtreamCluster::Live => &mut live_col, XtreamCluster::Series => &mut series_col, XtreamCluster::Video => &mut vod_col, @@ -197,7 +197,7 @@ pub(crate) fn write_xtream_playlist(target: &ConfigTarget, cfg: &Config, playlis } } else { None }; drop(header); - if let Some(pl) = col { + if let Some(pl) = col { pl.push(pli); } } @@ -277,10 +277,10 @@ 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))) } -pub(crate) fn get_xtream_item_for_stream_id(stream_id: u32, config: &Config, target: &ConfigTarget) -> Result { +pub(crate) fn get_xtream_item_for_stream_id(stream_id: u32, config: &Config, target: &ConfigTarget) -> Result { if let Some(storage_path) = get_xtream_storage_path(config, target.name.as_str()) { if let Some(mapping) = load_stream_id_cluster_mapping(&storage_path) { - if let Some((cluster, cluster_index)) = mapping.iter().find(|(_, c)| stream_id > *c) { + if let Some((cluster, cluster_index)) = mapping.iter().find(|(_, c)| stream_id >= *c) { let (xtream_path, idx_path) = get_xtream_file_paths(&storage_path, cluster); if xtream_path.exists() && idx_path.exists() { let offset: u64 = IndexRecord::get_index_offset(stream_id - cluster_index) as u64; @@ -290,7 +290,7 @@ pub(crate) fn get_xtream_item_for_stream_id(stream_id: u32, config: &Config, tar xtream_file.seek(SeekFrom::Start(index_record.index as u64))?; let mut buffer: Vec = vec![0; index_record.size as usize]; xtream_file.read_exact(&mut buffer)?; - if let Ok(pli) = bincode::deserialize::(&buffer[..]) { + if let Ok(pli) = bincode::deserialize::(&buffer[..]) { return Ok(pli); } } @@ -300,30 +300,37 @@ pub(crate) fn get_xtream_item_for_stream_id(stream_id: u32, config: &Config, tar Err(Error::new(ErrorKind::Other, format!("Failed to read xtream item for stream-id {}", stream_id))) } -pub(crate) fn load_rewrite_xtream_playlist(cluster: &XtreamCluster, config: &Config, target: &ConfigTarget) -> Result { +pub(crate) fn load_rewrite_xtream_playlist(cluster: &XtreamCluster, config: &Config, target: &ConfigTarget, category_id: u32) -> Result { if let Some(storage_path) = get_xtream_storage_path(config, target.name.as_str()) { let (xtream_path, idx_path) = get_xtream_file_paths(&storage_path, cluster); if xtream_path.exists() && idx_path.exists() { match std::fs::read(&xtream_path) { Ok(encoded_xtream) => { match std::fs::read(&idx_path) { - Ok(encoded_idx) => { - let mut cursor = 0; - let size = encoded_idx.len(); - let options = XtreamMappingOptions::from_target_options(target.options.as_ref()); - let mut result = vec![]; - while cursor < size { - let tuple = IndexRecord::from_bytes(&encoded_idx, &mut cursor); - let start_offset = tuple.index as usize; - let end_offset = start_offset + tuple.size as usize; - if let Ok(pli) = bincode::deserialize::(&encoded_xtream[start_offset..end_offset]) { - result.push(pli.to_xtream(&options)); - } else { - error!("Could not deserialize item {}", &xtream_path.to_str().unwrap()); - } - } - return Ok(serde_json::to_string(&result).unwrap()); - } + Ok(encoded_idx) => { + let mut cursor = 0; + let size = encoded_idx.len(); + let options = XtreamMappingOptions::from_target_options(target.options.as_ref()); + let mut result = vec![]; + let mut deserialize_error = false; + while cursor < size { + let index_record = IndexRecord::from_bytes(&encoded_idx, &mut cursor); + let start_offset = index_record.index as usize; + let end_offset = start_offset + index_record.size as usize; + match bincode::deserialize::(&encoded_xtream[start_offset..end_offset]) { + Ok(pli) => { + if category_id == 0 || pli.category_id == category_id { + result.push(pli.to_doc(&options)); + } + }, + Err(_) => deserialize_error = true, + }; + } + if deserialize_error { + error!("Could not deserialize item {}", &xtream_path.to_str().unwrap()); + } + return Ok(serde_json::to_string(&result).unwrap()); + } Err(err) => error!("Could not open file {}: {}", &idx_path.to_str().unwrap(), err), } }