diff --git a/src/api/api_model.rs b/src/api/api_model.rs index e5e0f6e76..393f495ec 100644 --- a/src/api/api_model.rs +++ b/src/api/api_model.rs @@ -284,14 +284,12 @@ pub(crate) struct ServerTargetConfig { pub watch: Option>, } - #[derive(Deserialize, Serialize, Debug, Clone)] pub(crate) struct ServerSourceConfig { pub inputs: Vec, pub targets: Vec, } - #[derive(Deserialize, Serialize, Debug, Clone)] pub(crate) struct ServerConfig { pub api: ConfigApi, diff --git a/src/api/download_api.rs b/src/api/download_api.rs index 0298434c5..75781dbc0 100644 --- a/src/api/download_api.rs +++ b/src/api/download_api.rs @@ -113,18 +113,18 @@ macro_rules! download_info { pub(crate) async fn queue_download_file( req: web::Json, - _app_state: web::Data, + app_state: web::Data, ) -> HttpResponse { - if let Some(download_cfg) = &_app_state.config.video.as_ref().unwrap().download { + if let Some(download_cfg) = &app_state.config.video.as_ref().unwrap().download { if download_cfg.directory.is_none() { return HttpResponse::BadRequest().json(json!({"error": "Server config missing video.download.directory configuration"})); } match FileDownload::new(req.url.as_str(), req.filename.as_str(), download_cfg) { Some(file_download) => { let response = HttpResponse::Ok().json(download_info!(file_download)); - _app_state.downloads.queue.lock().unwrap().push_back(file_download); - if _app_state.downloads.active.read().unwrap().is_none() { - match run_download_queue(download_cfg, Arc::clone(&_app_state.downloads)) { + app_state.downloads.queue.lock().unwrap().push_back(file_download); + if app_state.downloads.active.read().unwrap().is_none() { + match run_download_queue(download_cfg, Arc::clone(&app_state.downloads)) { Ok(_) => {} Err(err) => return HttpResponse::InternalServerError().json(json!({"error": err})), } @@ -139,12 +139,12 @@ pub(crate) async fn queue_download_file( } pub(crate) async fn download_file_info( - _app_state: web::Data, + app_state: web::Data, ) -> HttpResponse { - let finished_list: &[Value] = &_app_state.downloads.finished.write().unwrap().drain(..) + let finished_list: &[Value] = &app_state.downloads.finished.write().unwrap().drain(..) .map(|fd| download_info!(fd)).collect::>(); - match &*_app_state.downloads.active.read().unwrap() { + match &*app_state.downloads.active.read().unwrap() { None => HttpResponse::Ok().json(json!({ "completed": true, "downloads": finished_list })), diff --git a/src/api/m3u_api.rs b/src/api/m3u_api.rs index cf663f7e6..02f8928ac 100644 --- a/src/api/m3u_api.rs +++ b/src/api/m3u_api.rs @@ -4,7 +4,7 @@ use log::error; use crate::api::api_utils::{get_user_target, get_user_target_by_credentials, stream_response}; use crate::api::api_model::{AppState, UserApiRequest}; use crate::model::config::TargetType; -use crate::repository::m3u_repository::{get_m3u_file_paths, get_m3u_item_for_stream_id, rewrite_m3u_playlist}; +use crate::repository::m3u_repository::{get_m3u_file_paths, get_m3u_item_for_stream_id, load_rewrite_m3u_playlist}; async fn m3u_api( api_req: web::Query, @@ -16,7 +16,7 @@ async fn m3u_api( match get_user_target(&api_req, &app_state) { Some((user, target)) => { // let filename = target.get_m3u_filename(); - if let Some(content) = rewrite_m3u_playlist(&app_state.config, target, &user) { + if let Some(content) = load_rewrite_m3u_playlist(&app_state.config, target, &user) { HttpResponse::Ok().content_type(mime::TEXT_PLAIN_UTF_8).body(content) } else { HttpResponse::NoContent().finish() diff --git a/src/api/v1_api.rs b/src/api/v1_api.rs index 018caebe1..6d97af9a8 100644 --- a/src/api/v1_api.rs +++ b/src/api/v1_api.rs @@ -56,12 +56,12 @@ pub(crate) async fn save_config_api_proxy_user( pub(crate) async fn save_config_main( req: web::Json, - mut _app_state: web::Data, + app_state: web::Data, ) -> HttpResponse { let cfg = req.0; if cfg.is_valid() { - let file_path = _app_state.config._config_file_path.as_str(); - let backup_dir = _app_state.config.backup_dir.as_ref().unwrap().as_str(); + let file_path = app_state.config._config_file_path.as_str(); + let backup_dir = app_state.config.backup_dir.as_ref().unwrap().as_str(); if let Some(err) = _save_config_main(file_path, backup_dir, &cfg) { return HttpResponse::InternalServerError().json(json!({"error": err.to_string()})); } @@ -93,14 +93,14 @@ pub(crate) async fn save_config_api_proxy_config( pub(crate) async fn playlist_update( req: web::Json>, - _app_state: web::Data, + app_state: web::Data, ) -> HttpResponse { let targets = req.0; let user_targets = if targets.is_empty() { None } else { Some(targets) }; - let process_targets = validate_targets(&user_targets, &_app_state.config.sources); + let process_targets = validate_targets(&user_targets, &app_state.config.sources); match process_targets { Ok(valid_targets) => { - actix_rt::spawn(playlist_processor::exec_processing(Arc::clone(&_app_state.config), Arc::new(valid_targets))); + actix_rt::spawn(playlist_processor::exec_processing(Arc::clone(&app_state.config), Arc::new(valid_targets))); HttpResponse::Ok().finish() } Err(err) => { @@ -135,11 +135,11 @@ fn create_config_input_for_url(url: &str) -> ConfigInput { pub(crate) async fn playlist( req: web::Json, - _app_state: web::Data, + app_state: web::Data, ) -> HttpResponse { match match &req.input_id { Some(input_id) => { - _app_state.config.get_input_by_id(input_id) + app_state.config.get_input_by_id(input_id) } None => { let url = req.url.as_deref().unwrap_or(""); @@ -150,8 +150,8 @@ pub(crate) async fn playlist( Some(input) => { let (result, errors) = match input.input_type { - InputType::M3u => download::get_m3u_playlist(&_app_state.config, &input, &_app_state.config.working_dir).await, - InputType::Xtream => download::get_xtream_playlist(&input, &_app_state.config.working_dir).await, + InputType::M3u => download::get_m3u_playlist(&app_state.config, &input, &app_state.config.working_dir).await, + InputType::Xtream => download::get_xtream_playlist(&input, &app_state.config.working_dir).await, }; if result.is_empty() { let error_strings: Vec = errors.iter().map(|err| err.to_string()).collect(); diff --git a/src/api/xtream_api.rs b/src/api/xtream_api.rs index a835a03e2..c6bcfeb28 100644 --- a/src/api/xtream_api.rs +++ b/src/api/xtream_api.rs @@ -1,6 +1,7 @@ // https://github.com/tellytv/go.xtream-codes/blob/master/structs.go use std::collections::HashMap; +use std::fmt::{Display, Formatter}; use std::io::{Error}; use std::path::Path; use std::str::FromStr; @@ -14,8 +15,9 @@ use crate::api::api_model::{AppState, UserApiRequest, XtreamAuthorizationRespons use crate::model::api_proxy::{ProxyType, ProxyUserCredentials}; use crate::model::config::{Config, ConfigInput, ConfigTarget}; use crate::model::config::{TargetType}; -use crate::model::playlist::XtreamCluster; +use crate::model::playlist::{XtreamCluster}; use crate::repository::xtream_repository; +use crate::repository::xtream_repository::get_xtream_item_for_stream_id; use crate::utils::{json_utils, request_utils}; pub(crate) async fn serve_query(file_path: &Path, filter: &HashMap<&str, &str>) -> HttpResponse { @@ -93,7 +95,7 @@ fn get_user_info(user: &ProxyUserCredentials, cfg: &Config) -> XtreamAuthorizati } } -fn separate_number_and_rest(input: &str) -> (String, String) { +fn xtream_api_request_separate_number_and_rest(input: &str) -> (String, String) { if let Some(dot_index) = input.find('.') { let number_part = input[..dot_index].to_string(); let rest = input[dot_index..].to_string(); @@ -103,46 +105,68 @@ fn separate_number_and_rest(input: &str) -> (String, String) { } } +enum XtreamApiStreamContext { + LiveAlt, + Live, + Movie, + Series, + Timeshift +} + +impl Display for XtreamApiStreamContext { + fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { + write!(f, "{}", match self { + XtreamApiStreamContext::LiveAlt => "", + XtreamApiStreamContext::Live => "live", + XtreamApiStreamContext::Movie => "movie", + XtreamApiStreamContext::Series => "series", + XtreamApiStreamContext::Timeshift=> "timeshift", + }) + } +} + async fn xtream_player_api_stream( req: &HttpRequest, api_req: &web::Query, - _app_state: &web::Data, - context: &str, + app_state: &web::Data, + context: XtreamApiStreamContext, username: &str, password: &str, - action_path: &str, + stream_id: &str, + action_path: &str ) -> HttpResponse { - if let Some((user, target)) = get_user_target_by_credentials(username, password, api_req, _app_state) { + if let Some((user, target)) = get_user_target_by_credentials(username, password, api_req, app_state) { let target_name = &target.name; if target.has_output(&TargetType::Xtream) { - let mut stream_id = action_path.to_owned(); - let mut input: Option<&ConfigInput> = None; - if target.is_multi_input() { - let (action_stream_id, action_ext) = separate_number_and_rest(action_path); - if let Ok(num) = action_stream_id.trim().parse() { - let (xtream_id, cfg_input) = get_xtream_mapped_id_and_input_for_stream_id(_app_state, target_name, num); - if cfg_input.is_some() { - input = cfg_input; - stream_id = format!("{}{}", xtream_id, action_ext); - } - } - } else if let Some(inputs) = _app_state.config.get_inputs_for_target(target_name) { - input = inputs.first().copied(); - } + let (action_stream_id, stream_ext) = xtream_api_request_separate_number_and_rest(stream_id); + let req_stream_id: u32 = match FromStr::from_str(action_stream_id.trim()) { + Ok(id) => id, + Err(_) => return HttpResponse::BadRequest().finish() + }; - if let Some(target_input) = input { - if let Some(stream_url) = get_xtream_player_api_stream_url(target_input, context, stream_id.as_str()) { - if user.proxy == ProxyType::Redirect { - debug!("Redirecting stream request to {}", stream_url); - return HttpResponse::Found().insert_header(("Location", stream_url)).finish(); + 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() + }; + if let Some(input) = &app_state.config.get_input_by_id(&input_id) { + let mut query_path = if action_path.is_empty() { "".to_string() } else { format!("{}/", action_path) }; + query_path = format!("{}{}{}", query_path, pli.header.borrow().id, stream_ext); + if let Some(stream_url) = get_xtream_player_api_stream_url(input, context.to_string().as_str(), query_path.as_str()) { + if user.proxy == ProxyType::Redirect { + debug!("Redirecting stream request to {}", stream_url); + return HttpResponse::Found().insert_header(("Location", stream_url)).finish(); + } + return stream_response(&stream_url, req, Some(input)).await + } else { + error!("Cant find stream url for target {}, context {}, stream_id {}", target_name, context, req_stream_id); + } + } else { + error!("Cant find input for target {}, context {}, stream_id {}", target_name, context, req_stream_id); } - return stream_response(&stream_url, req, Some(target_input)).await - } else { - debug!("Cant figure out stream url for target {}, context {}, action {}", - target_name, context, action_path); - } - } else { - debug!("Cant find input definition for target {}", target_name); + }, + Err(_) => error!("Failed to read xtream item for stream id {}", req_stream_id), } } else { debug!("Target has no xtream output {}", target_name); @@ -157,92 +181,82 @@ async fn xtream_player_api_live_stream( req: HttpRequest, api_req: web::Query, path: web::Path<(String, String, String)>, - _app_state: web::Data, + app_state: web::Data, ) -> HttpResponse { let (username, password, stream_id) = path.into_inner(); - xtream_player_api_stream(&req, &api_req, &_app_state, "live", &username, &password, &stream_id).await + xtream_player_api_stream(&req, &api_req, &app_state, XtreamApiStreamContext::Live, &username, &password, &stream_id, "").await } async fn xtream_player_api_live_stream_alt( req: HttpRequest, api_req: web::Query, path: web::Path<(String, String, String)>, - _app_state: web::Data, + app_state: web::Data, ) -> HttpResponse { let (username, password, stream_id) = path.into_inner(); - xtream_player_api_stream(&req, &api_req, &_app_state, "", &username, &password, &stream_id).await + xtream_player_api_stream(&req, &api_req, &app_state, XtreamApiStreamContext::LiveAlt, &username, &password, &stream_id, "").await } async fn xtream_player_api_series_stream( req: HttpRequest, api_req: web::Query, path: web::Path<(String, String, String)>, - _app_state: web::Data, + app_state: web::Data, ) -> HttpResponse { let (username, password, stream_id) = path.into_inner(); - xtream_player_api_stream(&req, &api_req, &_app_state, "series", &username, &password, &stream_id).await + xtream_player_api_stream(&req, &api_req, &app_state, XtreamApiStreamContext::Series, &username, &password, &stream_id, "").await } async fn xtream_player_api_movie_stream( req: HttpRequest, api_req: web::Query, path: web::Path<(String, String, String)>, - _app_state: web::Data, + app_state: web::Data, ) -> HttpResponse { let (username, password, stream_id) = path.into_inner(); - xtream_player_api_stream(&req, &api_req, &_app_state, "movie", &username, &password, &stream_id).await + xtream_player_api_stream(&req, &api_req, &app_state, XtreamApiStreamContext::Movie, &username, &password, &stream_id, "").await } async fn xtream_player_api_timeshift_stream( req: HttpRequest, api_req: web::Query, path: web::Path<(String, String, String, String, String)>, - _app_state: web::Data, + app_state: web::Data, ) -> HttpResponse { let (username, password, duration, start, stream_id) = path.into_inner(); - let action_path = format!("{}/{}/{}", duration, start, stream_id); - xtream_player_api_stream(&req, &api_req, &_app_state, "timeshift", &username, &password, &action_path).await -} - -fn get_xtream_mapped_id_and_input_for_stream_id<'a>(app_state: &'a AppState, target_name: &str, stream_id: i32) -> (i32, Option<&'a ConfigInput>) { - if let Some(inputs) = app_state.config.get_inputs_for_target(target_name) { - if let Ok(Some(mapping)) = xtream_repository::read_xtream_mapping(stream_id as u32, app_state.config.as_ref(), target_name) { - if let Some(cfg_input) = inputs.iter().find(|&&inp| inp.id == mapping.input_id).cloned() { - return (mapping.stream_id as i32, Some(cfg_input)); - } - } - } - (stream_id, None) + let action_path = format!("{}/{}/", duration, start); + xtream_player_api_stream(&req, &api_req, &app_state, XtreamApiStreamContext::Timeshift, &username, &password, &stream_id, &action_path).await } async fn xtream_get_stream_info(app_state: &AppState, target_name: &str, stream_id: i32, cluster: &XtreamCluster) -> Result { - let (xtream_id, input) = get_xtream_mapped_id_and_input_for_stream_id(app_state, target_name, stream_id); - if let Some(target_input) = input { - if let Ok(content) = xtream_repository::xtream_get_stored_stream_info(app_state, target_name, stream_id, cluster, target_input).await { - return Ok(content); - } - if let Some(info_url) = get_xtream_player_api_info_url(target_input, cluster, xtream_id) { - if let Ok(url) = Url::parse(&info_url) { - let client = request_utils::get_client_request(Some(target_input), url, None); - if let Ok(response) = client.send().await { - debug!("{}", response.status()); - if response.status().is_success() { - match response.text().await { - Ok(content) => { - // TODO we are not replacing direct_source, we should add an option to do this. - xtream_repository::xtream_persist_stream_info(app_state, target_name, stream_id, cluster, - target_input, content.as_str()).await; - return Ok(content); - } - Err(err) => { error!("Failed to download info {}", err.to_string()); } - } - } - } - } - } - } + // let (xtream_id, input) = get_xtream_mapped_id_and_input_for_stream_id(app_state, target_name, stream_id); + // if let Some(target_input) = input { + // if let Ok(content) = xtream_repository::xtream_get_stored_stream_info(app_state, target_name, stream_id, cluster, target_input).await { + // return Ok(content); + // } + // + // if let Some(info_url) = get_xtream_player_api_info_url(target_input, cluster, xtream_id) { + // if let Ok(url) = Url::parse(&info_url) { + // let client = request_utils::get_client_request(Some(target_input), url, None); + // if let Ok(response) = client.send().await { + // debug!("{}", response.status()); + // if response.status().is_success() { + // match response.text().await { + // Ok(content) => { + // // TODO we are not replacing direct_source, we should add an option to do this. + // xtream_repository::xtream_persist_stream_info(app_state, target_name, stream_id, cluster, + // target_input, content.as_str()).await; + // return Ok(content); + // } + // Err(err) => { error!("Failed to download info {}", err.to_string()); } + // } + // } + // } + // } + // } + // } Err(Error::new(std::io::ErrorKind::Other, format!("Cant find stream with id: {}/{}/{}", target_name, &cluster, stream_id))) } @@ -277,35 +291,35 @@ async fn xtream_get_short_epg(app_state: &AppState, user: &ProxyUserCredentials, Err(_) => return HttpResponse::BadRequest().finish() }; - let (xtream_id, input) = get_xtream_mapped_id_and_input_for_stream_id(app_state, target_name, xtream_stream_id); - if let Some(target_input) = input { - if let Some(action_url) = get_xtream_player_api_action_url(target_input, "get_short_epg") { - let mut info_url = format!("{}&stream_id={}", action_url, xtream_id); - if !(limit.is_empty() || limit.eq("0")) { - info_url = format!("{}&limit={}", info_url, limit); - } - if let Ok(url) = Url::parse(&info_url) { - if user.proxy == ProxyType::Redirect { - return HttpResponse::Found().insert_header(("Location", info_url)).finish(); - } - - let client = request_utils::get_client_request(Some(target_input), url, None); - if let Ok(response) = client.send().await { - if response.status().is_success() { - return match response.text().await { - Ok(content) => { - HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body(content) - } - Err(err) => { - error!("Failed to download epg {}", err.to_string()); - HttpResponse::NoContent().finish() - } - }; - } - } - } - } - } + // let (xtream_id, input) = get_xtream_mapped_id_and_input_for_stream_id(app_state, target_name, xtream_stream_id); + // if let Some(target_input) = input { + // if let Some(action_url) = get_xtream_player_api_action_url(target_input, "get_short_epg") { + // let mut info_url = format!("{}&stream_id={}", action_url, xtream_id); + // if !(limit.is_empty() || limit.eq("0")) { + // info_url = format!("{}&limit={}", info_url, limit); + // } + // if let Ok(url) = Url::parse(&info_url) { + // if user.proxy == ProxyType::Redirect { + // return HttpResponse::Found().insert_header(("Location", info_url)).finish(); + // } + // + // let client = request_utils::get_client_request(Some(target_input), url, None); + // if let Ok(response) = client.send().await { + // if response.status().is_success() { + // return match response.text().await { + // Ok(content) => { + // HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body(content) + // } + // Err(err) => { + // error!("Failed to download epg {}", err.to_string()); + // HttpResponse::NoContent().finish() + // } + // }; + // } + // } + // } + // } + // } error!("Cant find short epg with id: {}/{}", target_name, stream_id); HttpResponse::NoContent().finish() } @@ -391,16 +405,16 @@ async fn xtream_player_api( async fn xtream_player_api_get(req: HttpRequest, api_req: web::Query, - _app_state: web::Data, + app_state: web::Data, ) -> HttpResponse { - xtream_player_api(&req, api_req.into_inner(), &_app_state).await + xtream_player_api(&req, api_req.into_inner(), &app_state).await } async fn xtream_player_api_post(req: HttpRequest, api_req: web::Form, - _app_state: web::Data, + app_state: web::Data, ) -> HttpResponse { - xtream_player_api(&req, api_req.into_inner(), &_app_state).await + xtream_player_api(&req, api_req.into_inner(), &app_state).await } pub(crate) fn xtream_api_register(cfg: &mut web::ServiceConfig) { diff --git a/src/model/api_proxy.rs b/src/model/api_proxy.rs index dd96c0f2b..1cb9a18a3 100644 --- a/src/model/api_proxy.rs +++ b/src/model/api_proxy.rs @@ -24,11 +24,10 @@ impl ProxyType { impl Display for ProxyType { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - let str = match self { - ProxyType::Reverse => "reverse".to_string(), - ProxyType::Redirect => "redirect".to_string() - }; - write!(f, "{}", str) + write!(f, "{}", match self { + ProxyType::Reverse => "reverse", + ProxyType::Redirect => "redirect" + }) } } diff --git a/src/model/config.rs b/src/model/config.rs index 719998375..7a8ae1b14 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -47,7 +47,7 @@ pub(crate) fn default_as_empty_str() -> String { String::from("") } pub(crate) fn default_as_empty_rc_str() -> Rc { Rc::new(String::from("")) } -pub(crate) fn default_as_zero() -> u8 { 0 } +fn default_as_zero_u8() -> u8 { 0 } fn default_as_frm() -> ProcessingOrder { ProcessingOrder::Frm } @@ -470,11 +470,10 @@ pub(crate) enum InputType { impl Display for InputType { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - let str = match self { - InputType::M3u => "m3u".to_string(), - InputType::Xtream => "xtream".to_string() - }; - write!(f, "{}", str) + write!(f, "{}", match self { + InputType::M3u => "m3u", + InputType::Xtream => "xtream" + }) } } @@ -701,7 +700,7 @@ impl VideoConfig { #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub(crate) struct ConfigDto { - #[serde(default = "default_as_zero")] + #[serde(default = "default_as_zero_u8")] pub threads: u8, pub api: ConfigApi, pub working_dir: String, @@ -801,7 +800,7 @@ impl WebAuthConfig { #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub(crate) struct Config { - #[serde(default = "default_as_zero")] + #[serde(default = "default_as_zero_u8")] pub threads: u8, pub api: ConfigApi, pub sources: Vec, diff --git a/src/model/playlist.rs b/src/model/playlist.rs index d4ee92556..9f46dc428 100644 --- a/src/model/playlist.rs +++ b/src/model/playlist.rs @@ -8,6 +8,7 @@ use serde_json::Value; use crate::model::config::{ConfigInput, ConfigTarget}; use crate::model::config::{default_as_false}; use crate::model::xmltv::TVGuide; +use crate::model::xtream::{xtream_playlistitem_to_document, XtreamMappingOptions}; // https://de.wikipedia.org/wiki/M3U // https://siptv.eu/howto/playlist.html @@ -30,7 +31,7 @@ impl FetchedPlaylist<'_> { } } -#[derive(Debug, Clone, PartialEq)] +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] pub(crate) enum XtreamCluster { Live = 1, Video = 2, @@ -58,7 +59,7 @@ pub(crate) enum PlaylistItemType { } pub(crate) fn default_playlist_item_type() -> PlaylistItemType { PlaylistItemType::Live } - +fn default_as_zero_u32() -> u32 { 0 } pub(crate) trait FieldAccessor { fn get_field(&self, field: &str) -> Option>; @@ -79,18 +80,19 @@ pub(crate) struct PlaylistItemHeader { pub audio_track: Rc, pub time_shift: Rc, pub rec: Rc, - // this is the source content not the url pub source: Rc, pub url: Rc, pub epg_channel_id: Option>, - #[serde(default = "default_stream_cluster", skip_serializing, skip_deserializing)] + #[serde(default = "default_stream_cluster")] pub xtream_cluster: XtreamCluster, - #[serde(skip_serializing, skip_deserializing)] 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, + } macro_rules! update_fields { @@ -223,6 +225,10 @@ impl PlaylistItem { epg_channel_id: header.epg_channel_id.clone(), } } + + pub fn to_xtream(&self, options: &XtreamMappingOptions) -> Value { + xtream_playlistitem_to_document(&self, options) + } } #[derive(Debug, Clone, Serialize, Deserialize)] diff --git a/src/model/xtream.rs b/src/model/xtream.rs index 06c263bb3..7e77cc57f 100644 --- a/src/model/xtream.rs +++ b/src/model/xtream.rs @@ -1,12 +1,29 @@ +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 crate::model::config::{default_as_empty_rc_str}; -use crate::model::playlist::{PlaylistItem}; +use crate::model::config::{ConfigTargetOptions, default_as_empty_rc_str}; +use crate::model::playlist::{PlaylistItem, PlaylistItemHeader, XtreamCluster}; + +const LIVE_STREAM_FIELDS: &[&str] = &[]; + +const VIDEO_STREAM_FIELDS: &[&str] = &[ + "release_date", "cast", + "director", "episode_run_time", "genre", + "stream_type", "title", "year", "youtube_trailer", + "plot", "rating_5based", "stream_icon", "container_extension" +]; + +const SERIES_STREAM_FIELDS: &[&str] = &[ + "backdrop_path", "cast", "cover", "director", "episode_run_time", "genre", + "last_modified", "name", "plot", "rating_5based", + "stream_type", "title", "year", "youtube_trailer", +]; fn default_as_empty_list() -> Vec { vec![] } @@ -333,8 +350,126 @@ impl XtreamSeriesInfoEpisode { } } -#[derive(Debug, Clone, Serialize, Deserialize)] -pub(crate) struct MultiXtreamMapping { - pub stream_id: u32, - pub input_id: u16, +pub(crate) struct XtreamMappingOptions { + pub skip_live_direct_source: bool, + pub skip_video_direct_source: bool, } + +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 + } + } +} + +fn append_release_date(document: &mut serde_json::Map) { + // Do we really need releaseDate ? + let has_release_date_1 = document.contains_key("release_date"); + let has_release_date_2 = document.contains_key("releaseDate"); + if !(has_release_date_1 && has_release_date_2) { + let release_date = if has_release_date_1 { + document.get("release_date") + } else if has_release_date_2 { + document.get("releaseDate") + } else { + None + }.map_or_else(|| Value::Null, |v| v.clone()); + if !&has_release_date_1 { + document.insert("release_date".to_string(), release_date.clone()); + } + if !&has_release_date_2 { + document.insert("releaseDate".to_string(), release_date.clone()); + } + } +} + +fn append_mandatory_fields(document: &mut serde_json::Map, fields: &[&str]) { + for &field in fields { + if !document.contains_key(field) { + document.insert(field.to_string(), Value::Null); + } + } +} + +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())); + } + } + } +} + +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())); + 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())), + ("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())), + ]); + + match header.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("thumbnail".to_string(), Value::String(header.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 { + None => Value::Null, + Some(epg_id) => Value::String(epg_id.as_ref().clone()) + }); + } + XtreamCluster::Video => { + document.insert("stream_id".to_string(), stream_id_value); + 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("custom_sid".to_string(), Value::String("".to_string())); + } + XtreamCluster::Series => { + document.insert("series_id".to_string(), stream_id_value); + } + }; + + if let Some(add_props) = &header.additional_properties { + for (field_name, field_value) in add_props { + document.insert(field_name.to_string(), field_value.to_owned()); + } + } + + match header.xtream_cluster { + XtreamCluster::Live => { + append_mandatory_fields(&mut document, LIVE_STREAM_FIELDS); + } + XtreamCluster::Video => { + append_mandatory_fields(&mut document, VIDEO_STREAM_FIELDS); + } + XtreamCluster::Series => { + append_prepared_series_properties(header, &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 1c4f877e8..c57980ca5 100644 --- a/src/processing/m3u_parser.rs +++ b/src/processing/m3u_parser.rs @@ -1,7 +1,7 @@ use std::borrow::{BorrowMut}; use std::cell::RefCell; use std::rc::Rc; -use crate::model::config::Config; +use crate::model::config::{Config, ConfigInput}; use crate::model::config::default_as_empty_rc_str; use crate::model::playlist::{default_playlist_item_type, default_stream_cluster, PlaylistGroup, PlaylistItem, PlaylistItemHeader, PlaylistItemType, XtreamCluster}; use crate::utils::string_utils; @@ -53,7 +53,7 @@ fn skip_digit(it: &mut std::str::Chars) -> Option { } } -fn create_empty_playlistitem_header(content: &String, url: String) -> PlaylistItemHeader { +fn create_empty_playlistitem_header(input_id: u16, url: String) -> PlaylistItemHeader { PlaylistItemHeader { id: default_as_empty_rc_str(), stream_id: default_as_empty_rc_str(), @@ -66,13 +66,14 @@ fn create_empty_playlistitem_header(content: &String, url: String) -> PlaylistIt audio_track: default_as_empty_rc_str(), time_shift: default_as_empty_rc_str(), rec: default_as_empty_rc_str(), - source: Rc::new(content.to_owned()), + 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(), xtream_cluster: default_stream_cluster(), additional_properties: None, series_fetched: false, + category_id: 0, } } @@ -87,8 +88,8 @@ macro_rules! process_header_fields { }; } -fn process_header(video_suffixes: &Vec<&str>, content: &String, url: String) -> PlaylistItemHeader { - let mut plih = create_empty_playlistitem_header(content, url.clone()); +fn process_header(input: &ConfigInput, video_suffixes: &Vec<&str>, content: &String, url: String) -> PlaylistItemHeader { + let mut plih = create_empty_playlistitem_header(input.id, url.clone()); let mut it = content.chars(); let line_token = token_till(&mut it, ':'); if line_token == Some(String::from("#EXTINF")) { @@ -154,7 +155,7 @@ fn extract_id_from_url(url: &str) -> Option { None } -pub(crate) fn consume_m3u(cfg: &Config, lines: impl Iterator, mut visit: F) { +pub(crate) fn consume_m3u(cfg: &Config, input: &ConfigInput, lines: impl Iterator, mut visit: F) { let mut header: Option = None; let mut group: Option = None; @@ -172,15 +173,17 @@ pub(crate) fn consume_m3u(cfg: &Config, lines: impl Iter continue; } if let Some(header_value) = header { - let item = PlaylistItem { header: RefCell::new(process_header(&video_suffixes, &header_value,line)) }; - if item.header.borrow().group.is_empty() { + 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 { - item.header.borrow_mut().group = Rc::new(group_value); + header.group = Rc::new(group_value); } else { - let current_title = item.header.borrow().title.to_owned(); - item.header.borrow_mut().group = Rc::new(string_utils::get_title_group(current_title.as_str())); + let current_title = header.title.to_owned(); + header.group = Rc::new(string_utils::get_title_group(current_title.as_str())); } } + drop(header); visit(item); } header = None; @@ -188,11 +191,11 @@ pub(crate) fn consume_m3u(cfg: &Config, lines: impl Iter } } -pub(crate) fn parse_m3u(cfg: &Config, lines: &[String]) -> Vec { +pub(crate) fn parse_m3u(cfg: &Config, input: &ConfigInput, lines: &[String]) -> Vec { let mut groups: std::collections::HashMap, Vec> = std::collections::HashMap::new(); let mut sort_order: Vec> = vec![]; let mut playlist = Vec::new(); - consume_m3u(cfg, lines.iter().cloned(), |item| playlist.push(item)); + consume_m3u(cfg, input, lines.iter().cloned(), |item| playlist.push(item)); playlist.drain(..).for_each(|item| { let key = Rc::clone(&item.header.borrow().group); // let key2 = String::from(&item.header.group); diff --git a/src/processing/playlist_processor.rs b/src/processing/playlist_processor.rs index bc0ae75f5..d34b0d974 100644 --- a/src/processing/playlist_processor.rs +++ b/src/processing/playlist_processor.rs @@ -1,5 +1,6 @@ extern crate unidecode; +use core::cmp::Ordering; use std::cell::RefCell; use std::collections::{HashMap, HashSet}; use std::rc::Rc; @@ -14,20 +15,17 @@ use crate::{Config, get_errors_notify_message, model::config}; use crate::filter::{get_field_value, MockValueProcessor, set_field_value, ValueProvider}; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::messaging::{MsgKind, send_message}; -use crate::model::config::{ConfigTarget, default_as_default, InputType, ProcessTargets}; -use crate::model::config::{ItemField, ProcessingOrder, SortOrder::{Asc, Desc}, TargetType}; +use crate::model::config::{ConfigSortChannel, ConfigSortGroup, ConfigTarget, default_as_default, InputType, + ItemField, ProcessingOrder, ProcessTargets, SortOrder::{Asc, Desc}}; use crate::model::mapping::{Mapping, MappingValueProcessor}; - use crate::model::playlist::{FetchedPlaylist, PlaylistGroup, PlaylistItem}; use crate::model::stats::{InputStats, PlaylistStats}; -use crate::model::xtream::MultiXtreamMapping; +use crate::processing::affix_processor::apply_affixes; use crate::processing::playlist_watch::process_group_watch; use crate::processing::xmltv_parser::flatten_tvguide; use crate::processing::xtream_processor::playlist_resolve_series; use crate::repository::playlist_repository::persist_playlist; -use crate::repository::xtream_repository::write_xtream_mapping; use crate::utils::download; -use crate::processing::affix_processor::apply_affixes; fn is_valid(pli: &PlaylistItem, target: &ConfigTarget) -> bool { let provider = ValueProvider { pli: RefCell::new(pli) }; @@ -53,37 +51,41 @@ fn filter_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Opt Some(new_playlist) } +fn playlistgroup_comparator(a: &PlaylistGroup, b: &PlaylistGroup, group_sort: &ConfigSortGroup, match_as_ascii: bool) -> Ordering { + let value_a = if match_as_ascii { Rc::new(unidecode(&a.title)) } else { Rc::clone(&a.title) }; + let value_b = if match_as_ascii { Rc::new(unidecode(&b.title)) } else { Rc::clone(&b.title) }; + let ordering = value_a.partial_cmp(&value_b).unwrap(); + match group_sort.order { + Asc => ordering, + Desc => ordering.reverse() + } +} + +fn playlistitem_comparator(a: &PlaylistItem, b: &PlaylistItem, channel_sort: &ConfigSortChannel, match_as_ascii: bool) -> Ordering { + let raw_value_a = get_field_value(a, &channel_sort.field); + 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 }; + let ordering = value_a.partial_cmp(&value_b).unwrap(); + match channel_sort.order { + Asc => ordering, + Desc => ordering.reverse() + } +} + fn sort_playlist(target: &ConfigTarget, new_playlist: &mut [PlaylistGroup]) { if let Some(sort) = &target.sort { - let match_as_ascii = &sort.match_as_ascii; + let match_as_ascii = sort.match_as_ascii; if let Some(group_sort) = &sort.groups { - new_playlist.sort_by(|a, b| { - let value_a = if *match_as_ascii { Rc::new(unidecode(&a.title)) } else { Rc::clone(&a.title) }; - let value_b = if *match_as_ascii { Rc::new(unidecode(&b.title)) } else { Rc::clone(&b.title) }; - let ordering = value_a.partial_cmp(&value_b).unwrap(); - match group_sort.order { - Asc => ordering, - Desc => ordering.reverse() - } - }); + new_playlist.sort_by(|a, b| playlistgroup_comparator(a, b, group_sort, match_as_ascii)); } if let Some(channel_sorts) = &sort.channels { channel_sorts.iter().for_each(|channel_sort| { let regexp = channel_sort.re.as_ref().unwrap(); new_playlist.iter_mut().for_each(|group| { - let group_title = if *match_as_ascii { Rc::new(unidecode(&group.title)) } else { Rc::clone(&group.title) }; + let group_title = if match_as_ascii { Rc::new(unidecode(&group.title)) } else { Rc::clone(&group.title) }; if regexp.is_match(group_title.as_str()) { - group.channels.sort_by(|a, b| { - let raw_value_a = get_field_value(a, &channel_sort.field); - 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 }; - let ordering = value_a.partial_cmp(&value_b).unwrap(); - match channel_sort.order { - Asc => ordering, - Desc => ordering.reverse() - } - }); + group.channels.sort_by(|chan1, chan2| playlistitem_comparator(chan1, chan2, channel_sort, match_as_ascii)); } }); }); @@ -392,68 +394,22 @@ async fn process_playlist<'a>(playlists: &mut [FetchedPlaylist<'a>], apply_affixes(&mut new_fetched_playlists); - if target.is_multi_input() && target.has_output(&TargetType::Xtream) { - let mut stream_id_mappings: Vec = Vec::new(); - let mut counter: u32 = 0; - new_fetched_playlists.iter() - .flat_map(|pl| { - let input_id = &pl.input.id; - pl.playlist.iter().map(move |plg| (input_id, &plg.channels)) - }) - .flat_map(|(input_id, channels)| channels.iter().map(move |chan| (input_id, chan))) - .for_each(|(input_id, chan)| { - let mut header = chan.header.borrow_mut(); - if header.stream_id.is_empty() { - header.stream_id = Rc::clone(&header.id); - } - match header.stream_id.parse::() { - Ok(stream_id) => { - let xtream_mapping = MultiXtreamMapping { - stream_id, - input_id: *input_id, - }; - stream_id_mappings.push(xtream_mapping); - } - Err(_) => { - error!("Failed to parse stream_id: {}", &header.id) - } - } - counter += 1; - header.id = Rc::new(counter.to_string()); - }); - - match write_xtream_mapping(&stream_id_mappings, cfg, &target.name) { - Ok(_) => { - debug!("wrote multi xtream input mapping for {}", &target.name); - } - Err(err) => { - return Err(vec![M3uFilterError::new( - M3uFilterErrorKind::Notify, - format!("Write multi xtream input mapping {} failed: {}", target.name, err))]); - } - } - } - let mut new_playlist = vec![]; let mut new_epg = vec![]; - let mut tv_guides = vec![]; + new_fetched_playlists.drain(..).for_each(|mut fp| { + let epg_channel_ids: HashSet<_> = fp.playlist.iter().flat_map(|g| &g.channels) + .filter_map(|c| c.header.borrow().epg_channel_id.clone()).collect(); fp.playlist.drain(..).for_each(|group| new_playlist.push(group)); - if let Some(tv_guide) = fp.epg { - tv_guides.push(tv_guide); - let guide = tv_guides.last().unwrap(); - if log_enabled!(Level::Debug) { + if !epg_channel_ids.is_empty() { + if let Some(tv_guide) = fp.epg { debug!("found epg information for {}", &target.name); - } - let channel_ids: HashSet<_> = new_playlist.iter().flat_map(|g| &g.channels) - .filter_map(|c| c.header.borrow().epg_channel_id.clone()).collect(); - if !channel_ids.is_empty() { - if let Some(epg) = guide.filter(&channel_ids) { + if let Some(epg) = tv_guide.filter(&epg_channel_ids) { new_epg.push(epg); } - } else if log_enabled!(Level::Debug) { - debug!("channel ids are empty"); } + } else if log_enabled!(Level::Debug) { + debug!("channel ids are empty"); } }); diff --git a/src/processing/xtream_parser.rs b/src/processing/xtream_parser.rs index 2e49fe2d5..73f1f3e61 100644 --- a/src/processing/xtream_parser.rs +++ b/src/processing/xtream_parser.rs @@ -50,7 +50,6 @@ pub(crate) fn parse_xtream_series_info(info: &Value, group_title: &str, input: & audio_track: default_as_empty_rc_str(), time_shift: default_as_empty_rc_str(), rec: default_as_empty_rc_str(), - // source is meant to hold the original provider data source: default_as_empty_rc_str(), url: if episode.direct_source.is_empty() { let ext = episode.container_extension.to_owned(); @@ -64,6 +63,7 @@ pub(crate) fn parse_xtream_series_info(info: &Value, group_title: &str, input: & xtream_cluster: XtreamCluster::Series, additional_properties: episode.get_additional_properties(&series_info), series_fetched: false, + category_id: 0, }) }).collect(); if result.is_empty() { Ok(None) } else { Ok(Some(result)) } @@ -80,6 +80,7 @@ pub(crate) fn parse_xtream(input: &ConfigInput, streams: &Value) -> Result>, M3uFilterError> { match map_to_xtream_category(category) { Ok(mut categories) => { + let input_id = input.id.to_string(); 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); @@ -108,8 +109,7 @@ 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 is meant to hold the original provider data - source: 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()), @@ -134,6 +134,7 @@ pub(crate) fn parse_xtream(input: &ConfigInput, xtream_cluster: xtream_cluster.clone(), additional_properties: stream.get_additional_properties(), series_fetched: false, + category_id: 0, }), }; grp.add(item); @@ -146,7 +147,7 @@ pub(crate) fn parse_xtream(input: &ConfigInput, id: cat.category_id.parse::().unwrap_or(0), xtream_cluster: xtream_cluster.clone(), title: Rc::clone(&cat.category_name), - channels: cat.channels.clone() + channels: cat.channels.clone(), } }).collect())) } diff --git a/src/repository/m3u_repository.rs b/src/repository/m3u_repository.rs index 64c3cafde..3169e375d 100644 --- a/src/repository/m3u_repository.rs +++ b/src/repository/m3u_repository.rs @@ -13,6 +13,7 @@ use crate::model::config::{Config, ConfigTarget}; use crate::model::playlist::{M3uPlaylistItem, PlaylistGroup, PlaylistItemType}; use crate::repository::repository_utils::IndexRecord; use crate::utils::file_utils; +use crate::utils::file_utils::create_file_tuple; macro_rules! cant_write_result { ($path:expr, $err:expr) => { @@ -36,18 +37,6 @@ pub(crate) fn get_m3u_epg_file_path(cfg: &Config, filename: &Option) -> .map(|path| file_utils::add_prefix_to_filename(&path, "epg_", Some("xml"))) } -fn create_m3u_files(m3u_path: &Path, idx_path: &Path) -> Result<(File, File), M3uFilterError> { - match File::create(m3u_path) { - Ok(m3u_file) => { - match File::create(idx_path) { - Ok(idx_file) => Ok((m3u_file, idx_file)), - Err(e) => cant_write_result!(&idx_path, e), - } - } - Err(e) => cant_write_result!(&m3u_path, e), - } -} - pub(crate) fn write_m3u_playlist(target: &ConfigTarget, cfg: &Config, new_playlist: &[PlaylistGroup], filename: &Option) -> Result<(), M3uFilterError> { if !new_playlist.is_empty() { if filename.is_none() { @@ -57,7 +46,7 @@ pub(crate) fn write_m3u_playlist(target: &ConfigTarget, cfg: &Config, new_playli } if let Some((m3u_path, idx_path)) = get_m3u_file_paths(cfg, filename) { - match create_m3u_files(&m3u_path, &idx_path) { + match create_file_tuple(&m3u_path, &idx_path) { Ok((mut m3u_file, mut m3u_idx_file)) => { let mut idx_offset: u32 = 0; let m3u_playlist = new_playlist.iter() @@ -70,12 +59,12 @@ pub(crate) fn write_m3u_playlist(target: &ConfigTarget, cfg: &Config, new_playli if let Ok(encoded) = bincode::serialize(&m3u) { match file_utils::check_write(m3u_file.write_all(&encoded)) { Ok(_) => { - let bytes_written = encoded.len() as u32; - let combined_bytes: [u8; 8] = IndexRecord::new(idx_offset, bytes_written).to_bytes(); + let bytes_written = encoded.len() as u16; + let combined_bytes = IndexRecord::new(idx_offset, bytes_written).to_bytes(); if let Err(err) = file_utils::check_write(m3u_idx_file.write_all(&combined_bytes)) { - return cant_write_result!(&m3u_path, err); + return cant_write_result!(&idx_path, err); } - idx_offset += bytes_written; + idx_offset += bytes_written as u32; stream_id += 1; } Err(err) => { @@ -92,7 +81,7 @@ pub(crate) fn write_m3u_playlist(target: &ConfigTarget, cfg: &Config, new_playli Ok(()) } -pub(crate) fn rewrite_m3u_playlist(cfg: &Config, target: &ConfigTarget, user: &ProxyUserCredentials) -> Option { +pub(crate) fn load_rewrite_m3u_playlist(cfg: &Config, target: &ConfigTarget, user: &ProxyUserCredentials) -> Option { let filename = target.get_m3u_filename(); if filename.is_some() { if let Some((m3u_path, idx_path)) = get_m3u_file_paths(cfg, &filename) { diff --git a/src/repository/playlist_repository.rs b/src/repository/playlist_repository.rs index 2169fa869..0b31556e0 100644 --- a/src/repository/playlist_repository.rs +++ b/src/repository/playlist_repository.rs @@ -4,7 +4,7 @@ use crate::model::playlist::PlaylistGroup; use crate::model::xmltv::Epg; use crate::repository::epg_repository::write_epg; use crate::repository::kodi_repository::write_strm_playlist; -use crate::repository::m3u_repository::{write_m3u_playlist}; +use crate::repository::m3u_repository::write_m3u_playlist; use crate::repository::xtream_repository::write_xtream_playlist; pub(crate) fn persist_playlist(playlist: &mut [PlaylistGroup], epg: Option, diff --git a/src/repository/repository_utils.rs b/src/repository/repository_utils.rs index d657c00de..ab4f8b979 100644 --- a/src/repository/repository_utils.rs +++ b/src/repository/repository_utils.rs @@ -12,43 +12,43 @@ of the encoded file, and size is the size of the encoded struct. */ pub struct IndexRecord { pub index: u32, - pub size: u32, + pub size: u16, } impl IndexRecord { - pub fn new(index: u32, size: u32) -> IndexRecord { + pub fn new(index: u32, size: u16) -> Self { IndexRecord { index, size } } - pub fn from_file(file: &mut File, offset: u64) -> Result { + pub fn from_file(file: &mut File, offset: u64) -> Result { file.seek(SeekFrom::Start(offset))?; let mut index_bytes = [0u8; 4]; - let mut size_bytes = [0u8; 4]; + let mut size_bytes = [0u8; 2]; file.read_exact(&mut index_bytes)?; file.read_exact(&mut size_bytes)?; let index = u32::from_le_bytes(index_bytes); - let size = u32::from_le_bytes(size_bytes); + let size = u16::from_le_bytes(size_bytes); Ok(IndexRecord { index, size }) } - pub fn from_bytes(bytes: &[u8], cursor: &mut usize) -> IndexRecord { + pub fn from_bytes(bytes: &[u8], cursor: &mut usize) -> Self { let index_bytes: [u8; 4] = bytes[*cursor..*cursor + 4].try_into().unwrap(); *cursor += 4; - let size_bytes: [u8; 4] = bytes[*cursor..*cursor + 4].try_into().unwrap(); + let size_bytes: [u8; 2] = bytes[*cursor..*cursor + 2].try_into().unwrap(); *cursor += 4; let index = u32::from_le_bytes(index_bytes); - let size = u32::from_le_bytes(size_bytes); + let size = u16::from_le_bytes(size_bytes); IndexRecord { index, size } } - pub fn to_bytes(&self) -> [u8; 8] { + pub fn to_bytes(&self) -> [u8; 6] { let index_bytes: [u8; 4] = self.index.to_le_bytes(); - let size_bytes: [u8; 4] = self.size.to_le_bytes(); - let mut combined_bytes: [u8; 8] = [0; 8]; + let size_bytes: [u8; 2] = self.size.to_le_bytes(); + let mut combined_bytes: [u8; 6] = [0; 6]; combined_bytes[..4].copy_from_slice(&index_bytes); combined_bytes[4..].copy_from_slice(&size_bytes); combined_bytes } - pub fn get_index_offset(index: u32) -> u32 { index * 8 } + pub fn get_index_offset(index: u32) -> u32 { index * 6 } } \ No newline at end of file diff --git a/src/repository/xtream_repository.rs b/src/repository/xtream_repository.rs index b59ae72cc..4c4baca33 100644 --- a/src/repository/xtream_repository.rs +++ b/src/repository/xtream_repository.rs @@ -1,20 +1,21 @@ -use std::cell::Ref; use std::collections::{BTreeMap, HashMap}; use std::{fs, io}; use std::fs::{File, OpenOptions}; use std::io::{BufReader, BufWriter, Error, ErrorKind, Read, Seek, SeekFrom, Write}; -use std::iter::FromIterator; use std::path::{Path, PathBuf}; +use std::rc::Rc; use log::{error}; use serde::Serialize; -use serde_json::{json, Map, Value}; +use serde_json::{json, Value}; use crate::model::config::{Config, ConfigInput, ConfigTarget}; -use crate::model::playlist::{PlaylistGroup, PlaylistItemHeader, PlaylistItemType, XtreamCluster}; +use crate::model::playlist::{PlaylistGroup, PlaylistItem, PlaylistItemType, XtreamCluster}; use crate::{create_m3u_filter_error_result}; use crate::api::api_model::AppState; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; -use crate::model::xtream::MultiXtreamMapping; +use crate::model::xtream::{XtreamMappingOptions}; +use crate::repository::repository_utils::IndexRecord; use crate::utils::file_utils; +use crate::utils::file_utils::create_file_tuple; use crate::utils::json_utils::iter_json_array; type IndexTree = BTreeMap; @@ -26,21 +27,6 @@ pub(crate) static COL_LIVE: &str = "live"; pub(crate) static COL_SERIES: &str = "series"; pub(crate) static COL_VOD: &str = "vod"; -const LIVE_STREAM_FIELDS: &[&str] = &[]; - -const VIDEO_STREAM_FIELDS: &[&str] = &[ - "release_date", "cast", - "director", "episode_run_time", "genre", - "stream_type", "title", "year", "youtube_trailer", - "plot", "rating_5based", "stream_icon", "container_extension" -]; - -const SERIES_STREAM_FIELDS: &[&str] = &[ - "backdrop_path", "cast", "cover", "director", "episode_run_time", "genre", - "last_modified", "name", "plot", "rating_5based", - "stream_type", "title", "year", "youtube_trailer", -]; - pub(crate) fn get_xtream_storage_path(cfg: &Config, target_name: &str) -> Option { file_utils::get_file_path(&cfg.working_dir, Some(std::path::PathBuf::from(target_name.replace(' ', "_")))) } @@ -49,6 +35,21 @@ pub(crate) fn get_xtream_epg_file_path(path: &Path) -> PathBuf { path.join("epg.xml") } +pub(crate) fn get_xtream_file_paths(storage_path: &Path, cluster: &XtreamCluster) -> (PathBuf, PathBuf) { + let xtream_path = storage_path.join(format!("{}.db", match cluster { + XtreamCluster::Live => COL_LIVE, + XtreamCluster::Video => COL_VOD, + XtreamCluster::Series => COL_SERIES + })); + let extension = xtream_path.extension().map(|ext| format!("{}_", ext.to_str().unwrap_or(""))).unwrap_or("".to_owned()); + let index_path = xtream_path.with_extension(format!("{}idx", &extension)); + (xtream_path, index_path) +} + +pub(crate) fn get_xtream_id_index_file_path(storage_path: &Path) -> PathBuf { + storage_path.join("stream_id.db") +} + fn get_collection_path(path: &Path, collection: &str) -> PathBuf { path.join(format!("{}.json", collection)) } @@ -61,22 +62,6 @@ fn get_info_idx_path(path: &Path, collection: &str) -> PathBuf { path.join(format!("{}_info.idx", collection)) } -fn write_to_file(file: &Path, value: &T) -> Result<(), Error> - where - T: ?Sized + Serialize { - match File::create(file) { - Ok(file) => { - let mut writer = BufWriter::new(file); - serde_json::to_writer(&mut writer, value)?; - match writer.flush() { - Ok(_) => Ok(()), - Err(e) => Err(e) - } - } - Err(e) => Err(e) - } -} - fn get_info_collection_and_idx_path(path: &Path, cluster: &XtreamCluster) -> (PathBuf, PathBuf) { let collection = match cluster { XtreamCluster::Live => COL_LIVE, @@ -126,23 +111,98 @@ fn ensure_xtream_storage_path(cfg: &Config, target_name: &str) -> Result(file: &Path, value: &T) -> Result<(), Error> + where + T: ?Sized + Serialize { + match File::create(file) { + Ok(file) => { + let mut writer = BufWriter::new(file); + serde_json::to_writer(&mut writer, value)?; + match writer.flush() { + Ok(_) => Ok(()), + Err(e) => Err(e) + } + } + Err(e) => Err(e) + } +} + +macro_rules! cant_write_result { + ($path:expr, $err:expr) => { + create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "failed to write xtream playlist: {} - {}", $path.to_str().unwrap() ,$err) + } +} + +fn write_playlist_to_file(storage_path: &Path, stream_id: &mut u32, cluster: &XtreamCluster, playlist: &mut [PlaylistItem]) -> Result<(), M3uFilterError> { + let (xtream_path, idx_path) = get_xtream_file_paths(storage_path, cluster); + match create_file_tuple(&xtream_path, &idx_path) { + Ok((mut main_file, mut idx_file)) => { + 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) { + match file_utils::check_write(main_file.write_all(&encoded)) { + Ok(_) => { + let bytes_written = encoded.len() as u16; + let combined_bytes = IndexRecord::new(idx_offset, bytes_written).to_bytes(); + if let Err(err) = file_utils::check_write(idx_file.write_all(&combined_bytes)) { + return cant_write_result!(&idx_path, err); + } + idx_offset += bytes_written as u32; + *stream_id += 1; + } + Err(err) => { + return cant_write_result!(&xtream_path, err); + } + } + } + } + } + Err(err) => return cant_write_result!(&xtream_path, err), + } + Ok(()) +} + +fn save_stream_id_cluster_mapping(storage_path: &Path, id_data: &mut Vec<(XtreamCluster, u32)>) -> Result<(), Error> { + let stream_id_path = get_xtream_id_index_file_path(storage_path); + id_data.sort_by(|(_, a), (_, b)| b.cmp(a)); + let encoded: Vec = bincode::serialize(id_data).unwrap(); + fs::write(stream_id_path, encoded) +} + +fn load_stream_id_cluster_mapping(storage_path: &Path) -> Option> { + let path = get_xtream_id_index_file_path(storage_path); + match fs::read(path) { + Ok(encoded) => { + let decoded: Vec<(XtreamCluster, u32)> = bincode::deserialize(&encoded[..]).unwrap(); + Some(decoded) + } + Err(_) => None, + } +} + +fn write_playlists_to_file(storage_path: &Path, collections: Vec<(XtreamCluster, &mut [PlaylistItem])>) -> Result<(), M3uFilterError> { + let mut id_list: Vec<(XtreamCluster, u32)> = vec![]; + let mut stream_id: u32 = 1; + for (cluster, playlist) in collections { + id_list.push((cluster.clone(), stream_id)); + write_playlist_to_file(storage_path, &mut stream_id, &cluster, playlist)?; + } + 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))) + } +} + pub(crate) fn write_xtream_playlist(target: &ConfigTarget, cfg: &Config, playlist: &mut [PlaylistGroup]) -> Result<(), M3uFilterError> { match ensure_xtream_storage_path(cfg, target.name.as_str()) { Ok(path) => { - let (skip_live_direct_source, skip_video_direct_source) = target.options.as_ref() - .map_or((false, false), |o| (o.xtream_skip_live_direct_source, o.xtream_skip_video_direct_source)); - let mut cat_live_col = vec![]; let mut cat_series_col = vec![]; let mut cat_vod_col = vec![]; let mut live_col = vec![]; let mut series_col = vec![]; let mut vod_col = vec![]; - - let mut vod_map = HashMap::::new(); - let mut series_map = HashMap::::new(); - - let mut channel_num: i32 = 0; let mut errors = Vec::new(); // preserve category_ids @@ -160,97 +220,31 @@ pub(crate) fn write_xtream_playlist(target: &ConfigTarget, cfg: &Config, playlis XtreamCluster::Live => &mut cat_live_col, XtreamCluster::Series => &mut cat_series_col, XtreamCluster::Video => &mut cat_vod_col, - }.push( - json!({ - "category_id": format!("{}", &cat_id), - "category_name": plg.title.clone(), - "parent_id": 0 - })); + }.push(json!({ + "category_id": format!("{}", &cat_id), + "category_name": plg.title.clone(), + "parent_id": 0 + })); - for pli in &plg.channels { - let header = &pli.header.borrow(); - if let Ok(stream_id) = header.id.parse::() { - if header.item_type == PlaylistItemType::Series { - // we skip resolved series, because this is only necessary when writing m3u files - continue; + for pli in plg.channels.drain(..) { + let mut header = pli.header.borrow_mut(); + // we skip resolved series, because this is only necessary when writing m3u files + let col = if header.item_type != PlaylistItemType::Series { + if let Ok(_) = header.id.parse::() { + header.category_id = *cat_id; + Some(match pli.header.borrow().xtream_cluster { + XtreamCluster::Live => &mut live_col, + XtreamCluster::Series => &mut series_col, + XtreamCluster::Video => &mut vod_col, + }) + } else { + errors.push(format!("Channel does not have an id: {}", pli.header.borrow().title.as_str())); + None } - channel_num += 1; - let mut document = serde_json::Map::from_iter([ - ("category_id".to_string(), Value::String(format!("{}", &cat_id))), - ("category_ids".to_string(), Value::Array(Vec::from([Value::Number(serde_json::Number::from(cat_id.to_owned()))]))), - ("name".to_string(), Value::String(header.name.as_ref().clone())), - ("num".to_string(), Value::Number(serde_json::Number::from(channel_num))), - ("title".to_string(), Value::String(header.title.as_ref().clone())), - ("stream_icon".to_string(), Value::String(header.logo.as_ref().clone())), - ]); - - let stream_id_value = Value::Number(serde_json::Number::from(stream_id)); - match header.xtream_cluster { - XtreamCluster::Live => { - document.insert("stream_id".to_string(), stream_id_value); - if 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("thumbnail".to_string(), Value::String(header.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 { - None => Value::Null, - Some(epg_id) => Value::String(epg_id.as_ref().clone()) - }); - } - XtreamCluster::Video => { - document.insert("stream_id".to_string(), stream_id_value); - if 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("custom_sid".to_string(), Value::String("".to_string())); - } - XtreamCluster::Series => { - document.insert("series_id".to_string(), stream_id_value); - } - }; - - if let Some(add_props) = &header.additional_properties { - for (field_name, field_value) in add_props { - document.insert(field_name.to_string(), field_value.to_owned()); - } - } - - match header.xtream_cluster { - XtreamCluster::Live => { - append_mandatory_fields(&mut document, LIVE_STREAM_FIELDS); - } - XtreamCluster::Video => { - append_mandatory_fields(&mut document, VIDEO_STREAM_FIELDS); - } - XtreamCluster::Series => { - append_prepared_series_properties(header, &mut document); - append_mandatory_fields(&mut document, SERIES_STREAM_FIELDS); - append_release_date(&mut document); - } - }; - - match header.xtream_cluster { - XtreamCluster::Live => {} - XtreamCluster::Series => { - series_map.insert(stream_id, serde_json::to_string(&document).unwrap()); - } - XtreamCluster::Video => { - vod_map.insert(stream_id, serde_json::to_string(&document).unwrap()); - } - } - - match header.xtream_cluster { - XtreamCluster::Live => &mut live_col, - XtreamCluster::Series => &mut series_col, - XtreamCluster::Video => &mut vod_col, - }.push(Value::Object(document)); - } else { - errors.push(format!("Channel does not have an id: {}", &header.title)); + } else { None }; + drop(header); + if let Some(pl) = col { + pl.push(pli); } } } @@ -259,17 +253,25 @@ pub(crate) fn write_xtream_playlist(target: &ConfigTarget, cfg: &Config, playlis for (col_path, data) in [ (get_collection_path(&path, COL_CAT_LIVE), &cat_live_col), (get_collection_path(&path, COL_CAT_VOD), &cat_vod_col), - (get_collection_path(&path, COL_CAT_SERIES), &cat_series_col), - (get_collection_path(&path, COL_LIVE), &live_col), - (get_collection_path(&path, COL_VOD), &vod_col), - (get_collection_path(&path, COL_SERIES), &series_col)] { - match write_to_file(&col_path, data) { + (get_collection_path(&path, COL_CAT_SERIES), &cat_series_col)] { + match write_documents_to_file(&col_path, data) { Ok(()) => {} Err(err) => { errors.push(format!("Persisting collection failed: {}: {}", &col_path.to_str().unwrap(), err)); } } } + + match write_playlists_to_file(&path, vec![ + (XtreamCluster::Live, &mut live_col), + (XtreamCluster::Video, &mut vod_col), + (XtreamCluster::Series, &mut series_col)]) { + Ok(()) => {} + Err(err) => { + errors.push(format!("Persisting collection failed:{}", err)); + } + } + if !errors.is_empty() { return create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "{}", errors.join("\n")); } @@ -311,52 +313,6 @@ fn load_old_category_ids(path: &Path) -> (u32, HashMap) { (max_id, result) } -fn append_prepared_series_properties(header: &Ref, document: &mut 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_release_date(document: &mut Map) { - // Do we really need releaseDate ? - let has_release_date_1 = document.contains_key("release_date"); - let has_release_date_2 = document.contains_key("releaseDate"); - if !(has_release_date_1 && has_release_date_2) { - let release_date = if has_release_date_1 { - document.get("release_date") - } else if has_release_date_2 { - document.get("releaseDate") - } else { - None - }.map_or_else(|| Value::Null, |v| v.clone()); - if !&has_release_date_1 { - document.insert("release_date".to_string(), release_date.clone()); - } - if !&has_release_date_2 { - document.insert("releaseDate".to_string(), release_date.clone()); - } - } -} - -fn append_mandatory_fields(document: &mut Map, fields: &[&str]) { - for &field in fields { - if !document.contains_key(field) { - document.insert(field.to_string(), Value::Null); - } - } -} - pub(crate) fn xtream_get_collection_path(cfg: &Config, target_name: &str, collection_name: &str) -> Result<(Option, Option), Error> { if let Some(path) = get_xtream_storage_path(cfg, target_name) { let col_path = get_collection_path(&path, collection_name); @@ -452,53 +408,60 @@ pub(crate) async fn xtream_persist_stream_info( } } -fn get_id_mapping_path(path: &Path) -> PathBuf { - path.join("id_mapping.db") -} - -pub(crate) fn write_xtream_mapping(mappings: &[MultiXtreamMapping], config: &Config, target_name: &str) -> Result<(), M3uFilterError> { - if let Some(path) = get_xtream_storage_path(config, target_name) { - if fs::create_dir_all(&path).is_err() { - let msg = format!("Failed to save, can't create directory {}", &path.to_str().unwrap()); - return Err(M3uFilterError::new(M3uFilterErrorKind::Notify, msg)); - } - - let err_map = |e: Error| M3uFilterError::new(M3uFilterErrorKind::Notify, e.to_string()); - let mut file = File::create(get_id_mapping_path(&path)).map_err(err_map)?; - // We assume the mappings list is created with a counter as id - // and id 1 means the 0 index. We write all the data and can calculate the offset inside the - // file by (u32 size + u16 size) * index. - for mapping in mappings { - file.write_all(&mapping.stream_id.to_le_bytes()).map_err(err_map)?; - file.write_all(&mapping.input_id.to_le_bytes()).map_err(err_map)?; - }; - return Ok(()); - } - Err(M3uFilterError::new(M3uFilterErrorKind::Notify, format!("Failed to find the xtream storage path for {}", target_name))) -} - -pub(crate) fn read_xtream_mapping(id: u32, config: &Config, target_name: &str) -> io::Result> { - if id < 1 { - return Err(Error::new(ErrorKind::Other, "id should start with 1")); - } - - if let Some(path) = get_xtream_storage_path(config, target_name) { - let mapping_file_path = get_id_mapping_path(&path); - if mapping_file_path.exists() { - let mut file = File::open(&mapping_file_path)?; - let index = (id - 1) as u64; - let mapping_size = 4 + 2; // u32 + u16 - let offset = mapping_size * index; - - file.seek(SeekFrom::Start(offset))?; - let mut stream_id_bytes = [0u8; 4]; - file.read_exact(&mut stream_id_bytes)?; - let stream_id = u32::from_le_bytes(stream_id_bytes); - let mut input_id_bytes = [0u8; 2]; - file.read_exact(&mut input_id_bytes)?; - let input_id = u16::from_le_bytes(input_id_bytes); - return Ok(Some(MultiXtreamMapping { stream_id, input_id })); +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) { + 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; + let mut idx_file = File::open(idx_path)?; + let mut xtream_file = File::open(xtream_path)?; + let index_record = IndexRecord::from_file(&mut idx_file, offset)?; + 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[..]) { + return Ok(pli); + } + } + } } } - Ok(None) -} \ No newline at end of file + 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 { + 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()); + } + Err(err) => error!("Could not open file {}: {}", &idx_path.to_str().unwrap(), err), + } + } + Err(err) => error!("Could not open file {}: {}", &xtream_path.to_str().unwrap(), err), + } + } + } + Err(Error::new(ErrorKind::Other, format!("Failed to find xtream storage for target {}", &target.name))) +} + diff --git a/src/utils/download.rs b/src/utils/download.rs index 67138839f..263f27801 100644 --- a/src/utils/download.rs +++ b/src/utils/download.rs @@ -34,7 +34,7 @@ pub(crate) async fn get_m3u_playlist(cfg: &Config, input: &ConfigInput, working_ match request_utils::get_input_text_content(input, working_dir, &url, persist_file_path).await { Ok(text) => { let lines = text.lines().map(String::from).collect::>(); - (m3u_parser::parse_m3u(cfg, &lines), vec![]) + (m3u_parser::parse_m3u(cfg, input, &lines), vec![]) } Err(err) => (vec![], vec![err]) } diff --git a/src/utils/file_utils.rs b/src/utils/file_utils.rs index ce7aeb910..ab089c590 100644 --- a/src/utils/file_utils.rs +++ b/src/utils/file_utils.rs @@ -1,8 +1,11 @@ use std::fs; +use std::fs::File; use std::io::{Write}; use std::path::{Path, PathBuf}; use log::{debug, error}; use path_absolutize::*; +use crate::create_m3u_filter_error_result; +use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; #[macro_export] macro_rules! exit { @@ -172,3 +175,17 @@ pub(crate) fn check_write(res: std::io::Result<()>) -> Result<(), std::io::Error Err(_) => Err(std::io::Error::new(std::io::ErrorKind::Other, "Unable to write file")), } } + +pub(crate) fn create_file_tuple(path1: &Path, path2: &Path) -> Result<(File, File), M3uFilterError> { + match File::create(path1) { + Ok(file1) => { + match File::create(path2) { + Ok(file2) => Ok((file1, file2)), + Err(err) => + create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "failed to create file: {} - {}", path2.to_str().unwrap(), err) + } + } + Err(err) => + create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "failed to create file: {} - {}", path1.to_str().unwrap(), err) + } +}