diff --git a/CHANGELOG.md b/CHANGELOG.md index 12f13e941..2ae254c8b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -21,6 +21,7 @@ - `userfile` optional userfile with generated userfile in format "username: password" per file, default name is user.txt in config path * Password generation argument --genpwd to generate passwords for userfile. + # v1.1.8(2024-03-06) * Fixed WebUI Option-Select * WebUI: added gallery view as second view for playlist diff --git a/src/api/api_model.rs b/src/api/api_model.rs index f73d3e843..e5e0f6e76 100644 --- a/src/api/api_model.rs +++ b/src/api/api_model.rs @@ -89,7 +89,7 @@ impl FileDownload { let mut x: usize = 1; while file_path.is_file() { filename = format!("{}_{}.{}", file_stem, x, file_ext); - file_path = file_dir.clone(); + file_path.clone_from(&file_dir); file_path.push(&filename); x += 1; } diff --git a/src/api/api_utils.rs b/src/api/api_utils.rs index 986e367bf..16a44e945 100644 --- a/src/api/api_utils.rs +++ b/src/api/api_utils.rs @@ -79,4 +79,4 @@ pub(crate) async fn stream_response(stream_url: &str, req: &HttpRequest, input: error!("Url is malformed {}", &stream_url) } HttpResponse::BadRequest().finish() -} \ No newline at end of file +} diff --git a/src/api/m3u_api.rs b/src/api/m3u_api.rs index 1cdc8f008..cf663f7e6 100644 --- a/src/api/m3u_api.rs +++ b/src/api/m3u_api.rs @@ -1,35 +1,26 @@ use actix_web::{HttpRequest, HttpResponse, web}; use log::error; -use crate::api::api_utils::{get_user_target, get_user_target_by_credentials, serve_file, stream_response}; +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::api_proxy::ProxyType; use crate::model::config::TargetType; -use crate::repository::m3u_repository::{get_m3u_file_paths, get_m3u_url_for_stream_id, rewrite_m3u_playlist}; +use crate::repository::m3u_repository::{get_m3u_file_paths, get_m3u_item_for_stream_id, rewrite_m3u_playlist}; async fn m3u_api( api_req: web::Query, //_api_req: web::Query>, - req: HttpRequest, + _req: HttpRequest, app_state: web::Data, ) -> HttpResponse { //let api_req = UserApiRequest::from_map(&_api_req); match get_user_target(&api_req, &app_state) { Some((user, target)) => { - let filename = target.get_m3u_filename(); - if filename.is_some() { - if let Some((m3u_path, _url_path, _idx_path)) = get_m3u_file_paths(&app_state.config, &filename) { - if user.proxy == ProxyType::Reverse { - if let Some(content) = rewrite_m3u_playlist(&app_state.config, target, &user) { - return HttpResponse::Ok().content_type(mime::TEXT_PLAIN_UTF_8).body(content); - } - HttpResponse::NoContent().finish(); - } else { - return serve_file(&m3u_path, &req, mime::TEXT_PLAIN_UTF_8).await; - } - } + // let filename = target.get_m3u_filename(); + if let Some(content) = rewrite_m3u_playlist(&app_state.config, target, &user) { + HttpResponse::Ok().content_type(mime::TEXT_PLAIN_UTF_8).body(content) + } else { + HttpResponse::NoContent().finish() } - HttpResponse::NoContent().finish() } None => { HttpResponse::BadRequest().finish() @@ -49,10 +40,10 @@ async fn m3u_api_stream( if target.has_output(&TargetType::M3u) { let filename = target.get_m3u_filename(); if filename.is_some() { - if let Some((_m3u_path, url_path, idx_path)) = get_m3u_file_paths(&app_state.config, &filename) { - match get_m3u_url_for_stream_id(m3u_stream_id, &url_path, &idx_path) { - Ok(stream_url) => { - return stream_response(&stream_url, &req, None).await + if let Some((m3u_path, idx_path)) = get_m3u_file_paths(&app_state.config, &filename) { + match get_m3u_item_for_stream_id(m3u_stream_id, &m3u_path, &idx_path) { + Ok(m3u_item) => { + return stream_response(m3u_item.url.as_str(), &req, None).await } Err(err) => { error!("Failed to get m3u url: {}", err); @@ -67,9 +58,9 @@ async fn m3u_api_stream( } pub(crate) fn m3u_api_register(cfg: &mut web::ServiceConfig) { - cfg.service(web::resource("/get.php").route(web::get().to(m3u_api))); - cfg.service(web::resource("/get.php").route(web::post().to(m3u_api))); - cfg.service(web::resource("/apiget").route(web::get().to(m3u_api))); - cfg.service(web::resource("/m3u").route(web::get().to(m3u_api))); - cfg.service(web::resource("/m3u-stream/{username}/{password}/{stream_id}").route(web::get().to(m3u_api_stream))); + cfg.service(web::resource("/get.php").route(web::get().to(m3u_api))) + .service(web::resource("/get.php").route(web::post().to(m3u_api))) + .service(web::resource("/apiget").route(web::get().to(m3u_api))) + .service(web::resource("/m3u").route(web::get().to(m3u_api))) + .service(web::resource("/m3u-stream/{username}/{password}/{stream_id}").route(web::get().to(m3u_api_stream))); } \ No newline at end of file diff --git a/src/api/main_api.rs b/src/api/main_api.rs index d123f9bec..1db53d379 100644 --- a/src/api/main_api.rs +++ b/src/api/main_api.rs @@ -21,11 +21,9 @@ use crate::processing::playlist_processor; fn get_web_dir_path(web_ui_enabled: bool, web_root: &str) -> Result { let web_dir = web_root.to_string(); let web_dir_path = PathBuf::from(&web_dir); - if web_ui_enabled { - if !&web_dir_path.exists() || !&web_dir_path.is_dir() { - return Err(std::io::Error::new(ErrorKind::NotFound, - format!("web_root does not exists or is not an directory: {:?}", &web_dir_path))); - } + if web_ui_enabled && (!&web_dir_path.exists() || !&web_dir_path.is_dir()) { + return Err(std::io::Error::new(ErrorKind::NotFound, + format!("web_root does not exists or is not an directory: {:?}", &web_dir_path))); }; Ok(web_dir_path) } diff --git a/src/api/v1_api.rs b/src/api/v1_api.rs index 919d14215..018caebe1 100644 --- a/src/api/v1_api.rs +++ b/src/api/v1_api.rs @@ -220,13 +220,13 @@ pub(crate) async fn config( // if we didn't read it from file then we should use it from app_state if result.api_proxy.is_none() { - result.api_proxy = app_state.config._api_proxy.read().unwrap().clone(); + result.api_proxy.clone_from(&app_state.config._api_proxy.read().unwrap()); } HttpResponse::Ok().json(result) } -pub(crate) fn v1_api_register(web_auth_enabled: bool) -> impl Fn(&mut web::ServiceConfig) -> () { +pub(crate) fn v1_api_register(web_auth_enabled: bool) -> impl Fn(&mut web::ServiceConfig) { return move |cfg: &mut web::ServiceConfig| { cfg.service(web::scope("/api/v1") .wrap(Condition::new(web_auth_enabled, HttpAuthentication::with_fn(validator))) diff --git a/src/api/web_index.rs b/src/api/web_index.rs index 4440b5259..5d3313309 100644 --- a/src/api/web_index.rs +++ b/src/api/web_index.rs @@ -1,5 +1,5 @@ use std::collections::HashMap; -use std::path::PathBuf; +use std::path::{Path, PathBuf}; use actix_files::NamedFile; use actix_web::{HttpRequest, HttpResponse, web}; @@ -19,16 +19,14 @@ async fn token( app_state: web::Data, ) -> HttpResponse { match &app_state.config.web_auth { - None => { - return no_web_auth_token(); - }, + None => no_web_auth_token(), Some(web_auth) => { let username = req.username.as_str(); let password = req.password.as_str(); - if username.len() > 0 && password.len() > 0 { + if !(username.is_empty() || password.is_empty()) { if let Some(hash) = web_auth.get_user_password(username) { - if verify_password(hash, &password.as_bytes()) { + if verify_password(hash, password.as_bytes()) { req.zeroize(); if let Ok(token) = create_jwt(web_auth) { return HttpResponse::Ok().json(HashMap::from([("token", token)])); @@ -49,7 +47,7 @@ async fn token_refresh( ) -> HttpResponse { match &app_state.config.web_auth { None => { - return no_web_auth_token(); + no_web_auth_token() }, Some(web_auth) => { let secret_key = web_auth.secret.as_ref(); @@ -71,14 +69,13 @@ async fn index( NamedFile::open(path) } -pub(crate) fn index_register(web_dir_path: &PathBuf) -> impl Fn(&mut web::ServiceConfig) -> () { - let wdp = web_dir_path.clone(); +pub(crate) fn index_register(web_dir_path: &Path) -> impl Fn(&mut web::ServiceConfig) + '_ { return move |cfg: &mut web::ServiceConfig| { cfg.service(web::scope("/auth") .route("/token", web::post().to(token)) .route("/refresh", web::post().to(token_refresh))); cfg.service(web::scope("/") .route("", web::get().to(index)) - .service(actix_files::Files::new("/", &wdp))); + .service(actix_files::Files::new("/", web_dir_path))); }; } \ No newline at end of file diff --git a/src/api/xmltv_api.rs b/src/api/xmltv_api.rs index 4067be7d0..fd0958e59 100644 --- a/src/api/xmltv_api.rs +++ b/src/api/xmltv_api.rs @@ -2,37 +2,37 @@ use std::path::PathBuf; use actix_web::{HttpRequest, HttpResponse, web}; use log::{debug, info}; -use url::Url; +use url::{ParseError, Url}; use crate::api::api_model::{AppState, UserApiRequest}; use crate::api::api_utils::{get_user_target, serve_file}; -use crate::model::api_proxy::ProxyType; -use crate::model::config::{Config, ConfigTarget, InputType}; +use crate::model::api_proxy::{ProxyType, ProxyUserCredentials}; +use crate::model::config::{Config, ConfigInput, ConfigTarget}; use crate::model::config::TargetType; use crate::repository::m3u_repository::get_m3u_epg_file_path; use crate::repository::xtream_repository::{get_xtream_epg_file_path, get_xtream_storage_path}; use crate::utils::{file_utils, request_utils}; +fn get_epg_path_for_target_of_type(target_name: &str, file_path: Option) -> Option { + if let Some(epg_path) = file_path { + if file_utils::path_exists(&epg_path) { + return Some(epg_path); + } else { + info!("Cant find epg file for {target_name} target: {}", epg_path.to_str().unwrap_or("?")) + } + } + None +} + fn get_epg_path_for_target(config: &Config, target: &ConfigTarget) -> Option { for output in &target.output { match output.target { TargetType::M3u => { - if let Some(epg_path) = get_m3u_epg_file_path(config, &target.get_m3u_filename()) { - if file_utils::path_exists(&epg_path) { - return Some(epg_path); - } else { - info!("Cant find epg file for m3u target: {}", epg_path.to_str().unwrap_or("?")) - } - } + return get_epg_path_for_target_of_type(&target.name, get_m3u_epg_file_path(config, &target.get_m3u_filename())); } TargetType::Xtream => { if let Some(storage_path) = get_xtream_storage_path(config, &target.name) { - let epg_path = get_xtream_epg_file_path(&storage_path); - if file_utils::path_exists(&epg_path) { - return Some(epg_path); - } else { - info!("Cant find epg file for xtream target: {}", epg_path.to_str().unwrap_or("?")) - } + return get_epg_path_for_target_of_type(&target.name, Some(get_xtream_epg_file_path(&storage_path))); } } TargetType::Strm => {} @@ -44,39 +44,13 @@ fn get_epg_path_for_target(config: &Config, target: &ConfigTarget) -> Option, req: HttpRequest, - _app_state: web::Data, + app_state: web::Data, ) -> HttpResponse { - if let Some((user, target)) = get_user_target(&api_req, &_app_state) { - match get_epg_path_for_target(&_app_state.config, target) { + if let Some((user, target)) = get_user_target(&api_req, &app_state) { + match get_epg_path_for_target(&app_state.config, target) { None => { - // If no epg_url is provided for input, we did not process the xmltv for our channels. - // We are now delivering the original untouched xmltv. - // If you want to use xmltv then provide the url in the config to filter unnecessary content. - // If you have multiple xtream sources, the first one will be used for epg - let target_name = &target.name; - if let Some(input) = _app_state.config.get_input_for_target(target_name, &InputType::Xtream) { - let epg_url = input.epg_url.as_ref().map_or("".to_string(), |s| s.to_owned()); - let api_url = if epg_url.is_empty() { - format!("{}/xmltv.php?username={}&password={}", - input.url.as_str(), - input.username.as_ref().unwrap_or(&"".to_string()).as_str(), - input.password.as_ref().unwrap_or(&"".to_string()).as_str(), - ) - } else { epg_url.to_string() }; - if let Ok(url) = Url::parse(&api_url) { - if user.proxy == ProxyType::Redirect { - debug!("Redirecting epg request to {}", api_url); - return HttpResponse::Found().insert_header(("Location", api_url)).finish(); - } - let client = request_utils::get_client_request(Some(input), url, None); - if let Ok(response) = client.send().await { - if response.status().is_success() { - if let Ok(content) = response.text().await { - return HttpResponse::Ok().content_type(mime::TEXT_XML).body(content); - } - } - } - } + if let Some(value) = get_xmltv_raw_epg(&app_state.config, &user, &target.name).await { + return value; } } Some(epg_path) => return serve_file(&epg_path, &req, mime::TEXT_XML).await @@ -86,7 +60,56 @@ async fn xmltv_api( r#""#) } +fn get_xmltv_epg_url(input: &ConfigInput) -> Result { + let epg_url = input.epg_url.as_ref().map_or("".to_string(), |s| s.to_owned()); + if epg_url.is_empty() { + if let Some(user_info) = input.get_user_info() { + let url = user_info.base_url.as_str(); + let username = user_info.username.as_str(); + let password = user_info.password.as_str(); + Url::parse(format!("{url}/xmltv.php?username={username}&password={password}").as_str()) + } else { + Err(ParseError::EmptyHost) + } + } else { + Url::parse(epg_url.as_str()) + } +} + +async fn get_xmltv_raw_epg(config: &Config, user: &ProxyUserCredentials, target_name: &str) -> Option { + // If no epg_url is provided for input, we did not process the xmltv for our channels. + // We are now delivering the original untouched xmltv. + // If you want to use xmltv then provide the url in the config to filter unnecessary content. + // If you have multiple xtream sources, no response because of mapped ids + // if you want epg for multi xtream input, then provide epg_url. + if let Some(inputs) = config.get_inputs_for_target(target_name) { + if inputs.len() == 1 { + if let Some(&input) = inputs.first() { + if let Ok(url) = get_xmltv_epg_url(input) { + if user.proxy == ProxyType::Redirect { + debug!("Redirecting epg request to {}", url.as_str()); + return Some(HttpResponse::Found().insert_header(("Location", url.as_str())).finish()); + } + let client = request_utils::get_client_request(Some(input), url, None); + if let Ok(response) = client.send().await { + if response.status().is_success() { + if let Ok(content) = response.text().await { + return Some(HttpResponse::Ok().content_type(mime::TEXT_XML).body(content)); + } + } + } + } else { + debug!("Could not generate epg url for {target_name}") + } + } + } else { + debug!("No epg_url is provided for target {target_name}, multi input requires epg_url") + } + } + None +} + pub(crate) fn xmltv_api_register(cfg: &mut web::ServiceConfig) { - cfg.service(web::resource("/xmltv.php").route(web::get().to(xmltv_api))); - cfg.service(web::resource("/epg").route(web::get().to(xmltv_api))); + cfg.service(web::resource("/xmltv.php").route(web::get().to(xmltv_api))) + .service(web::resource("/epg").route(web::get().to(xmltv_api))); } \ No newline at end of file diff --git a/src/api/xtream_api.rs b/src/api/xtream_api.rs index 89b626835..a835a03e2 100644 --- a/src/api/xtream_api.rs +++ b/src/api/xtream_api.rs @@ -1,79 +1,38 @@ // https://github.com/tellytv/go.xtream-codes/blob/master/structs.go use std::collections::HashMap; -use std::io::Error; +use std::io::{Error}; use std::path::Path; use std::str::FromStr; - use actix_web::{HttpRequest, HttpResponse, web}; use chrono::{Duration, Local}; use log::{debug, error}; use url::Url; -use crate::api::api_model::{AppState, UserApiRequest, XtreamAuthorizationResponse, XtreamServerInfo, XtreamUserInfo}; use crate::api::api_utils::{get_user_server_info, get_user_target, get_user_target_by_credentials, serve_file, stream_response}; +use crate::api::api_model::{AppState, UserApiRequest, XtreamAuthorizationResponse, XtreamServerInfo, XtreamUserInfo}; use crate::model::api_proxy::{ProxyType, ProxyUserCredentials}; -use crate::model::config::{Config, ConfigInput, InputType}; -use crate::model::config::TargetType; +use crate::model::config::{Config, ConfigInput, ConfigTarget}; +use crate::model::config::{TargetType}; use crate::model::playlist::XtreamCluster; use crate::repository::xtream_repository; use crate::utils::{json_utils, request_utils}; -struct M3uUrlInfo { - pub base_url: String, - pub username: String, - pub password: String, -} - -fn parse_m3u_url(url: &str) -> Option { - if let Ok(url) = Url::parse(url) { - let base_url = url.origin().ascii_serialization(); - let mut username = None; - let mut password = None; - for (key, value) in url.query_pairs() { - if key.eq("username") { - username = Some(value.into_owned()); - } else if key.eq("password") { - password = Some(value.into_owned()); - } - } - if username.is_some() || password.is_some() { - return Some(M3uUrlInfo { - base_url, - username: username.as_ref().unwrap().to_owned(), - password: username.as_ref().unwrap().to_owned(), - }); - } - } - None -} - pub(crate) async fn serve_query(file_path: &Path, filter: &HashMap<&str, &str>) -> HttpResponse { let filtered = json_utils::filter_json_file(file_path, filter); HttpResponse::Ok().json(filtered) } fn get_xtream_player_api_action_url(input: &ConfigInput, action: &str) -> Option { - match input.input_type { - InputType::M3u => { - match parse_m3u_url(input.url.as_str()) { - None => None, - Some(m3u_url_info) => Some( - format!("{}/player_api.php?username={}&password={}&action={}", - m3u_url_info.base_url, - m3u_url_info.username, - m3u_url_info.password, - action - )) - } - } - InputType::Xtream => Some( - format!("{}/player_api.php?username={}&password={}&action={}", - input.url.as_str(), - input.username.as_ref().unwrap_or(&"".to_string()).as_str(), - input.password.as_ref().unwrap_or(&"".to_string()).as_str(), - action - )) + if let Some(user_info) = input.get_user_info() { + Some(format!("{}/player_api.php?username={}&password={}&action={}", + &user_info.base_url, + &user_info.username, + &user_info.password, + action + )) + } else { + None } } @@ -89,25 +48,16 @@ fn get_xtream_player_api_info_url(input: &ConfigInput, cluster: &XtreamCluster, fn get_xtream_player_api_stream_url(input: &ConfigInput, context: &str, action_path: &str) -> Option { let ctx_path = if context.is_empty() { "".to_string() } else { format!("{}/", context) }; - match input.input_type { - InputType::M3u => match parse_m3u_url(input.url.as_str()) { - None => None, - Some(m3u_url_info) => Some( - format!("{}/{}{}/{}/{}", - m3u_url_info.base_url, - ctx_path, - m3u_url_info.username, - m3u_url_info.password, - action_path - )) - } - InputType::Xtream => Some(format!("{}/{}{}/{}/{}", - input.url.as_str(), - ctx_path, - input.username.as_ref().unwrap_or(&"".to_string()).as_str(), - input.password.as_ref().unwrap_or(&"".to_string()).as_str(), - action_path + if let Some(user_info) = input.get_user_info() { + Some(format!("{}/{}{}/{}/{}", + &user_info.base_url, + ctx_path, + &user_info.username, + &user_info.password, + action_path )) + } else { + None } } @@ -143,6 +93,16 @@ fn get_user_info(user: &ProxyUserCredentials, cfg: &Config) -> XtreamAuthorizati } } +fn 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(); + (number_part, rest) + } else { + (input.to_string(), String::new()) + } +} + async fn xtream_player_api_stream( req: &HttpRequest, api_req: &web::Query, @@ -155,11 +115,23 @@ async fn xtream_player_api_stream( 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) { - if let Some(target_input) = match _app_state.config.get_input_for_target(target_name, &InputType::Xtream) { - None => _app_state.config.get_input_for_target(target_name, &InputType::M3u), - Some(inp) => Some(inp) - } { - if let Some(stream_url) = get_xtream_player_api_stream_url(target_input, context, action_path) { + 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(); + } + + 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(); @@ -232,14 +204,26 @@ async fn xtream_player_api_timeshift_stream( 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) +} + async fn xtream_get_stream_info(app_state: &AppState, target_name: &str, stream_id: i32, cluster: &XtreamCluster) -> Result { - if let Some(target_input) = app_state.config.get_input_for_target(target_name, &InputType::Xtream) { + 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, stream_id) { + 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 { @@ -263,64 +247,67 @@ async fn xtream_get_stream_info(app_state: &AppState, target_name: &str, stream_ } async fn xtream_get_stream_info_response(app_state: &AppState, user: &ProxyUserCredentials, - target_name: &str, stream_id: &str, + target: &ConfigTarget, stream_id: &str, cluster: &XtreamCluster) -> HttpResponse { - match FromStr::from_str(stream_id) { - Ok(xtream_stream_id) => { - if user.proxy == ProxyType::Redirect { - if let Some(target_input) = app_state.config.get_input_for_target(target_name, &InputType::Xtream) { - if let Some(info_url) = get_xtream_player_api_info_url(target_input, cluster, xtream_stream_id) { - return HttpResponse::Found().insert_header(("Location", info_url)).finish(); - } + let req_stream_id: i32 = match FromStr::from_str(stream_id) { + Ok(id) => id, + Err(_) => return HttpResponse::BadRequest().finish() + }; + + if user.proxy == ProxyType::Redirect && !target.is_multi_input() { + if let Some(inputs) = app_state.config.get_inputs_for_target(&target.name) { + if let Some(&input) = inputs.first() { + if let Some(info_url) = get_xtream_player_api_info_url(input, cluster, req_stream_id) { + return HttpResponse::Found().insert_header(("Location", info_url)).finish(); } } - - match xtream_get_stream_info(app_state, target_name, xtream_stream_id, cluster).await { - Ok(content) => HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body(content), - Err(_) => HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body("{info:[]}"), - } } - Err(_) => HttpResponse::BadRequest().finish() + return HttpResponse::BadRequest().finish(); + } + + match xtream_get_stream_info(app_state, &target.name, req_stream_id, cluster).await { + Ok(content) => HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body(content), + Err(_) => HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body("{info:[]}"), } } async fn xtream_get_short_epg(app_state: &AppState, user: &ProxyUserCredentials, target_name: &str, stream_id: &str, limit: &str) -> HttpResponse { - if !stream_id.is_empty() { - if let Some(target_input) = app_state.config.get_input_for_target(target_name, &InputType::Xtream) { - 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, stream_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 xtream_stream_id: i32 = match FromStr::from_str(stream_id) { + Ok(id) => id, + Err(_) => return HttpResponse::BadRequest().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() - } else { - error!("No epg_id given, short epg needs id: {}", target_name); - HttpResponse::BadRequest().finish() } - + error!("Cant find short epg with id: {}/{}", target_name, stream_id); + HttpResponse::NoContent().finish() } async fn xtream_player_api( @@ -339,12 +326,12 @@ async fn xtream_player_api( match action { "get_series_info" => { - xtream_get_stream_info_response(_app_state, &user, target_name, + xtream_get_stream_info_response(_app_state, &user, target, api_req.series_id.trim(), &XtreamCluster::Series).await } "get_vod_info" => { - xtream_get_stream_info_response(_app_state, &user, target_name, + xtream_get_stream_info_response(_app_state, &user, target, api_req.vod_id.trim(), &XtreamCluster::Video).await } @@ -417,14 +404,14 @@ async fn xtream_player_api_post(req: HttpRequest, } pub(crate) fn xtream_api_register(cfg: &mut web::ServiceConfig) { - cfg.service(web::resource("/player_api.php").route(web::get().to(xtream_player_api_get)).route(web::post().to(xtream_player_api_get))); - cfg.service(web::resource("/panel_api.php").route(web::get().to(xtream_player_api_get)).route(web::post().to(xtream_player_api_get))); - cfg.service(web::resource("/xtream").route(web::get().to(xtream_player_api_get)).route(web::post().to(xtream_player_api_post))); - cfg.service(web::resource("/{username}/{password}/{stream_id}").route(web::get().to(xtream_player_api_live_stream_alt))); - cfg.service(web::resource("/live/{username}/{password}/{stream_id}").route(web::get().to(xtream_player_api_live_stream))); - cfg.service(web::resource("/movie/{username}/{password}/{stream_id}").route(web::get().to(xtream_player_api_movie_stream))); - cfg.service(web::resource("/series/{username}/{password}/{stream_id}").route(web::get().to(xtream_player_api_series_stream))); - cfg.service(web::resource("/timeshift/{username}/{password}/{duration}/{start}{stream_id}").route(web::get().to(xtream_player_api_timeshift_stream))); + cfg.service(web::resource("/player_api.php").route(web::get().to(xtream_player_api_get)).route(web::post().to(xtream_player_api_get))) + .service(web::resource("/panel_api.php").route(web::get().to(xtream_player_api_get)).route(web::post().to(xtream_player_api_get))) + .service(web::resource("/xtream").route(web::get().to(xtream_player_api_get)).route(web::post().to(xtream_player_api_post))) + .service(web::resource("/{username}/{password}/{stream_id}").route(web::get().to(xtream_player_api_live_stream_alt))) + .service(web::resource("/live/{username}/{password}/{stream_id}").route(web::get().to(xtream_player_api_live_stream))) + .service(web::resource("/movie/{username}/{password}/{stream_id}").route(web::get().to(xtream_player_api_movie_stream))) + .service(web::resource("/series/{username}/{password}/{stream_id}").route(web::get().to(xtream_player_api_series_stream))) + .service(web::resource("/timeshift/{username}/{password}/{duration}/{start}{stream_id}").route(web::get().to(xtream_player_api_timeshift_stream))); /* TODO cfg.service(web::resource("/hlsr/{token}/{username}/{password}/{channel}/{hash}/{chunk}").route(web::get().to(xtream_player_api_hlsr_stream))); cfg.service(web::resource("/hls/{token}/{chunk}").route(web::get().to(xtream_player_api_hls_stream))); diff --git a/src/auth/authenticator.rs b/src/auth/authenticator.rs index 707aa1956..ceb494269 100644 --- a/src/auth/authenticator.rs +++ b/src/auth/authenticator.rs @@ -32,8 +32,8 @@ pub(crate) fn create_jwt(web_auth_config: &WebAuthConfig) -> Result, secret_key: &[u8]) -> bool { if let Some(auth) = bearer { let token = auth.token(); - let token_message = decode::(&token, &DecodingKey::from_secret(secret_key), &Validation::new(Algorithm::HS256)); - if let Ok(_) = token_message { + let token_message = decode::(token, &DecodingKey::from_secret(secret_key), &Validation::new(Algorithm::HS256)); + if token_message.is_ok() { return true; } } diff --git a/src/auth/password.rs b/src/auth/password.rs index e1e5accb8..e889320b7 100644 --- a/src/auth/password.rs +++ b/src/auth/password.rs @@ -13,7 +13,7 @@ fn generate_salt(length: usize) -> String { pub(crate) fn hash(password: &[u8]) -> Option { let salt = generate_salt(64); - if password.len() > 0 { + if !password.is_empty() { let config = argon2::Config::default(); if let Ok(hash) = argon2::hash_encoded(password, salt.as_bytes(), &config) { return Some(hash); diff --git a/src/filter.rs b/src/filter.rs index a38e3713b..5f7de60e9 100644 --- a/src/filter.rs +++ b/src/filter.rs @@ -70,7 +70,6 @@ pub(crate) struct RegexWithCaptures { pub captures: Vec, } - #[derive(Parser)] //#[grammar = "filter.pest"] #[grammar_inline = r#" diff --git a/src/model/config.rs b/src/model/config.rs index 817b44591..719998375 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -2,6 +2,7 @@ use std::rc::Rc; use enum_iterator::Sequence; use std::borrow::BorrowMut; use std::collections::{HashMap, HashSet}; +use std::fmt::Display; use std::fs::File; use std::io::BufRead; use std::path::PathBuf; @@ -38,7 +39,6 @@ macro_rules! valid_property { }}; } - pub(crate) fn default_as_true() -> bool { true } pub(crate) fn default_as_false() -> bool { false } @@ -279,6 +279,8 @@ pub(crate) struct ConfigTargetOptions { pub xtream_skip_live_direct_source: bool, #[serde(default = "default_as_true")] pub xtream_skip_video_direct_source: bool, + // #[serde(default = "default_as_true")] + // pub xtream_skip_series_direct_source: bool, #[serde(default = "default_as_false")] pub xtream_resolve_series: bool, #[serde(default = "default_as_two")] @@ -394,9 +396,9 @@ impl ConfigTarget { } } - // pub(crate) fn is_multi_input(&self) -> bool { - // self._multi_input - // } + pub(crate) fn is_multi_input(&self) -> bool { + self._multi_input + } pub(crate) fn filter(&self, provider: &ValueProvider) -> bool { let mut processor = MockValueProcessor {}; @@ -439,13 +441,12 @@ impl ConfigSource { Ok(index + (self.inputs.len() as u16)) } - pub(crate) fn get_input_for_target(&self, target_name: &str, input_type: &InputType) -> Option<&ConfigInput> { + pub(crate) fn get_inputs_for_target(&self, target_name: &str) -> Option> { for target in &self.targets { if target.name.eq(target_name) { - for input in &self.inputs { - if input.input_type.eq(input_type) { - return Some(input); - } + let inputs = self.inputs.iter().filter(|&i| i.enabled).collect::>(); + if !inputs.is_empty() { + return Some(inputs); } } } @@ -467,12 +468,13 @@ pub(crate) enum InputType { Xtream, } -impl ToString for InputType { - fn to_string(&self) -> String { - match self { +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) } } @@ -502,11 +504,11 @@ pub(crate) struct ConfigInputOptions { pub xtream_skip_series: bool, } -// pub(crate) struct InputUserInfo { -// pub base_url: String, -// pub username: String, -// pub password: String, -// } +pub(crate) struct InputUserInfo { + pub base_url: String, + pub username: String, + pub password: String, +} fn default_as_type_m3u() -> InputType { InputType::M3u } @@ -576,36 +578,36 @@ impl ConfigInput { Ok(()) } - // pub(crate) fn get_user_info(&self) -> Option { - // if self.input_type == InputType::Xtream { - // if self.username.is_some() || self.password.is_some() { - // return Some(InputUserInfo { - // base_url: self.url.to_owned(), - // username: self.username.as_ref().unwrap().to_owned(), - // password: self.password.as_ref().unwrap().to_owned(), - // }); - // } - // } else if let Ok(url) = url::Url::parse(&self.url) { - // let base_url = url.origin().ascii_serialization(); - // let mut username = None; - // let mut password = None; - // for (key, value) in url.query_pairs() { - // if key.eq("username") { - // username = Some(value.into_owned()); - // } else if key.eq("password") { - // password = Some(value.into_owned()); - // } - // } - // if username.is_some() || password.is_some() { - // return Some(InputUserInfo { - // base_url, - // username: username.as_ref().unwrap().to_owned(), - // password: password.as_ref().unwrap().to_owned(), - // }); - // } - // } - // None - // } + pub(crate) fn get_user_info(&self) -> Option { + if self.input_type == InputType::Xtream { + if self.username.is_some() || self.password.is_some() { + return Some(InputUserInfo { + base_url: self.url.to_owned(), + username: self.username.as_ref().unwrap().to_owned(), + password: self.password.as_ref().unwrap().to_owned(), + }); + } + } else if let Ok(url) = url::Url::parse(&self.url) { + let base_url = url.origin().ascii_serialization(); + let mut username = None; + let mut password = None; + for (key, value) in url.query_pairs() { + if key.eq("username") { + username = Some(value.into_owned()); + } else if key.eq("password") { + password = Some(value.into_owned()); + } + } + if username.is_some() || password.is_some() { + return Some(InputUserInfo { + base_url, + username: username.as_ref().unwrap().to_owned(), + password: password.as_ref().unwrap().to_owned(), + }); + } + } + None + } } #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] @@ -767,19 +769,14 @@ impl WebAuthConfig { if let Ok(file) = File::open(&userfile_path) { let mut users = vec![]; let reader = std::io::BufReader::new(file); - for line in reader.lines() { - match line { - Ok(credential) => { - let mut parts = credential.split(':'); - if let (Some(username), Some(password)) = (parts.next(), parts.next()) { - users.push(UserCredential { - username: username.trim().to_string(), - password: password.trim().to_string(), - }); - debug!("Read ui user {}", username); - } - } - Err(_) => {} + for credentials in reader.lines().map_while(Result::ok) { + let mut parts = credentials.split(':'); + if let (Some(username), Some(password)) = (parts.next(), parts.next()) { + users.push(UserCredential { + username: username.trim().to_string(), + password: password.trim().to_string(), + }); + debug!("Read ui user {}", username); } } @@ -853,9 +850,11 @@ impl Config { } } - pub(crate) fn get_input_for_target(&self, target_name: &str, input_type: &InputType) -> Option<&ConfigInput> { + pub(crate) fn get_inputs_for_target(&self, target_name: &str) -> Option> { for source in &self.sources { - if let Some(cfg) = source.get_input_for_target(target_name, input_type) { return Some(cfg); } + if let Some(cfg) = source.get_inputs_for_target(target_name) { + return Some(cfg); + } } None } @@ -878,7 +877,7 @@ impl Config { } } - pub fn get_input_by_id(&self, input_id: &u16) -> Option { + pub(crate) fn get_input_by_id(&self, input_id: &u16) -> Option { for source in &self.sources { for input in &source.inputs { if input.id == *input_id { @@ -889,7 +888,7 @@ impl Config { None } - pub fn set_mappings(&mut self, mappings: Option) -> Result<(), M3uFilterError> { + pub(crate) fn set_mappings(&mut self, mappings: Option) -> Result<(), M3uFilterError> { if let Some(mapping_list) = mappings { for source in &mut self.sources { for target in &mut source.targets { @@ -984,9 +983,7 @@ impl Config { if !web_auth.enabled { self.web_auth = None } else { - if let Err(err) = web_auth.prepare(&self._config_path, resolve_var) { - return Err(err); - } + web_auth.prepare(&self._config_path, resolve_var)? } } diff --git a/src/model/playlist.rs b/src/model/playlist.rs index 5da2c0801..d4ee92556 100644 --- a/src/model/playlist.rs +++ b/src/model/playlist.rs @@ -68,7 +68,7 @@ pub(crate) trait FieldAccessor { #[derive(Debug, Clone, Serialize, Deserialize)] pub(crate) struct PlaylistItemHeader { // stream_id is a custom field for processing - pub stream_id: u32, + pub stream_id: Rc, pub id: Rc, pub name: Rc, pub logo: Rc, @@ -100,7 +100,7 @@ macro_rules! update_fields { stringify!($prop) => { $self.$prop = Rc::new($val); true - }, + } )* _ => false, } @@ -118,38 +118,6 @@ macro_rules! get_fields { }; } -impl FieldAccessor for PlaylistItemHeader { - fn get_field(&self, field: &str) -> Option> { - let val = get_fields!(self, field, id, name, logo, logo_small, group, title, parent_code, audio_track, time_shift, rec, source, url;); - if val.is_some() { - return val; - } - match field { - "epg_channel_id" | "epg_id" => self.epg_channel_id.clone(), - _ => None - } - } - - fn set_field(&mut self, field: &str, value: &str) -> bool { - let val = String::from(value); - let updated = update_fields!(self, field, id, name, logo, logo_small, group, title, parent_code, audio_track, time_shift, rec, source, url; val); - if updated { - return updated; - } - match field { - "epg_channel_id" | "epg_id" => { - self.epg_channel_id = Some(Rc::new(value.to_owned())); - true - }, - _ => false, - } - } -} - -#[derive(Debug, Clone, Serialize, Deserialize)] -pub(crate) struct PlaylistItem { - pub header: RefCell, -} macro_rules! to_m3u_non_empty_fields { ($header:expr, $line:expr, $(($prop:ident, $field:expr)),*;) => { @@ -161,33 +129,101 @@ 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;); + if val.is_some() { + return val; + } + match field { + "epg_channel_id" | "epg_id" => self.epg_channel_id.clone(), + _ => None + } + } -impl PlaylistItem { - pub fn to_m3u(&self, target: &ConfigTarget) -> String { + 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); + if updated { + return updated; + } + match field { + "epg_channel_id" | "epg_id" => { + self.epg_channel_id = Some(Rc::new(value.to_owned())); + true + } + _ => false, + } + } +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub(crate) struct M3uPlaylistItem { + pub stream_id: Rc, + pub id: Rc, + pub name: Rc, + pub logo: Rc, + pub logo_small: Rc, + pub group: Rc, + pub title: Rc, + pub parent_code: Rc, + pub audio_track: Rc, + pub time_shift: Rc, + pub rec: Rc, + pub url: Rc, + pub epg_channel_id: Option>, +} + + +impl M3uPlaylistItem { + pub fn to_m3u(&self, target: &ConfigTarget, url: Option<&str>) -> String { let options = target.options.as_ref(); - let header = self.header.borrow(); let ignore_logo = options.map_or(false, |o| o.ignore_logo); - let mut line = format!("#EXTINF:-1 str-id=\"{}\" tvg-id=\"{}\" tvg-name=\"{}\" group-title=\"{}\"", - header.stream_id, - header.epg_channel_id.as_ref().map_or("", |o| o.as_ref()), - header.name, header.group); + let mut line = format!("#EXTINF:-1 tvg-id=\"{}\" tvg-name=\"{}\" group-title=\"{}\"", + self.epg_channel_id.as_ref().map_or("", |o| o.as_ref()), + self.name, self.group); // line = format!("{} tvg-chno=\"{}\"", line, header.chno); if !ignore_logo { - to_m3u_non_empty_fields!(header, line, (logo, "tvg-logo"), (logo_small, "tvg-logo-small");); + to_m3u_non_empty_fields!(self, line, (logo, "tvg-logo"), (logo_small, "tvg-logo-small");); } - to_m3u_non_empty_fields!(header, line, + to_m3u_non_empty_fields!(self, line, (parent_code, "parent-code"), (audio_track, "audio-track"), (time_shift, "timeshift"), (rec, "tvg-rec");); - format!("{},{}\n{}", line, header.title, header.url) + format!("{},{}\n{}", line, self.title, if url.is_none() { self.url.as_str() } else { url.unwrap() }) } } +#[derive(Debug, Clone, Serialize, Deserialize)] +pub(crate) struct PlaylistItem { + pub header: RefCell, +} + +impl PlaylistItem { + pub fn to_m3u(&self) -> M3uPlaylistItem { + let header = self.header.borrow(); + M3uPlaylistItem { + stream_id: Rc::clone(&header.stream_id), + id: Rc::clone(&header.id), + 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), + audio_track: Rc::clone(&header.audio_track), + time_shift: Rc::clone(&header.time_shift), + rec: Rc::clone(&header.rec), + url: Rc::clone(&header.url), + epg_channel_id: header.epg_channel_id.clone(), + } + } +} #[derive(Debug, Clone, Serialize, Deserialize)] pub(crate) struct PlaylistGroup { diff --git a/src/model/stats.rs b/src/model/stats.rs index 34b085ba9..b752b4987 100644 --- a/src/model/stats.rs +++ b/src/model/stats.rs @@ -1,3 +1,4 @@ +use std::fmt::Display; use crate::model::config::InputType; #[derive(Debug, Clone)] @@ -6,9 +7,9 @@ pub(crate) struct PlaylistStats { pub channel_count: usize, } -impl ToString for PlaylistStats { - fn to_string(&self) -> String { - format!("{{\"groups\": {}, \"channels\": {}}}", self.group_count, self.channel_count) +impl Display for PlaylistStats { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + write!(f, "{}", format_args!("{{\"groups\": {}, \"channels\": {}}}", self.group_count, self.channel_count)) } } @@ -21,10 +22,11 @@ pub(crate) struct InputStats { pub processed_stats: PlaylistStats, } -impl ToString for InputStats { - fn to_string(&self) -> String { - format!("{{\"name\": {}, \"type\": {}, \"errors\": {}, \"raw\": {}, \"processed\": {}}}", - self.name, self.input_type.to_string(), self.error_count, - self.raw_stats.to_string(), self.processed_stats.to_string()) +impl Display for InputStats { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + let str = format!("{{\"name\": {}, \"type\": {}, \"errors\": {}, \"raw\": {}, \"processed\": {}}}", + self.name, self.input_type, self.error_count, + self.raw_stats, self.processed_stats); + write!(f, "{}", str) } } \ No newline at end of file diff --git a/src/model/xtream.rs b/src/model/xtream.rs index 0a3805492..06c263bb3 100644 --- a/src/model/xtream.rs +++ b/src/model/xtream.rs @@ -280,8 +280,8 @@ pub(crate) struct XtreamSeriesInfoEpisodeInfo { pub duration_secs: u32, pub duration: String, pub movie_image: String, - // "video": [], - // "audio": [], + pub video: Value, + pub audio: Value, pub bitrate: u32, pub rating: f64, pub season: u32, diff --git a/src/processing/affix_processor.rs b/src/processing/affix_processor.rs new file mode 100644 index 000000000..435968f5f --- /dev/null +++ b/src/processing/affix_processor.rs @@ -0,0 +1,66 @@ +use log::{debug, Level, log_enabled}; +use crate::model::config::{AFFIX_FIELDS, ConfigInput, InputAffix}; +use crate::model::playlist::{FetchedPlaylist, FieldAccessor, PlaylistItem}; +use crate::valid_property; + +type AffixProcessor<'a> = Box; + +fn create_affix_processor(affix: &InputAffix, is_prefix: bool) -> AffixProcessor { + Box::new(move |channel: &mut PlaylistItem| { + let header = &mut channel.header.borrow_mut(); + let value = if let Some(field_value) = header.get_field(affix.field.as_str()) { + if is_prefix { + format!("{}{}", field_value.as_str(), &affix.value) + } else { + format!("{}{}", &affix.value, field_value.as_str()) + } + } else { + String::from(&affix.value) + }; + if log_enabled!(Level::Debug) { + debug!("Applying input {}: {}={}", if is_prefix {"prefix"} else {"suffix"}, &affix.field, &value); + } + header.set_field(&affix.field, value.as_str()); + }) +} + +fn validate_and_create_affix_processor(affix: Option<&InputAffix>, is_prefix: bool) -> Option { + if let Some(affix_def) = affix { + if (valid_property!(&affix_def.field.as_str(), AFFIX_FIELDS) && !affix_def.value.is_empty()) { + return Some(create_affix_processor(affix_def, is_prefix)); + } + }; + None +} + +fn get_affix_processor(input: &ConfigInput) -> Option { + if input.suffix.is_some() || input.prefix.is_some() { + let processors: Vec = vec![ + validate_and_create_affix_processor(input.prefix.as_ref(), true), + validate_and_create_affix_processor(input.suffix.as_ref(), false) + ].into_iter().flatten().collect(); + + if !processors.is_empty() { + let apply_affix: AffixProcessor = Box::new(move |channel: &mut PlaylistItem| { + for x in &processors { + x(channel); + } + }); + return Some(apply_affix); + } + } + None +} + +pub fn apply_affixes(fetched_playlists: &mut [FetchedPlaylist]) { + fetched_playlists.iter_mut().for_each(|fetched_playlist| { + let FetchedPlaylist { input, playlist, epg: _ } = fetched_playlist; + if let Some(affix_processor) = get_affix_processor(input) { + playlist.iter_mut().for_each(|group| { + group.channels.iter_mut().for_each(|channel| { + affix_processor(channel); + }); + }); + } + }); +} diff --git a/src/processing/m3u_parser.rs b/src/processing/m3u_parser.rs index 14934b998..1c4f877e8 100644 --- a/src/processing/m3u_parser.rs +++ b/src/processing/m3u_parser.rs @@ -55,8 +55,8 @@ fn skip_digit(it: &mut std::str::Chars) -> Option { fn create_empty_playlistitem_header(content: &String, url: String) -> PlaylistItemHeader { PlaylistItemHeader { - stream_id: 0, id: default_as_empty_rc_str(), + stream_id: default_as_empty_rc_str(), name: default_as_empty_rc_str(), logo: default_as_empty_rc_str(), logo_small: default_as_empty_rc_str(), @@ -103,22 +103,16 @@ fn process_header(video_suffixes: &Vec<&str>, content: &String, url: String) -> let token = token_till(&mut it, '='); if let Some(t) = token { let value = token_value(&mut it); - if t.as_str().eq("str-id") { - if let Ok(stream_id) = value.to_string().parse::() { - plih.stream_id = stream_id; - } - } else { - process_header_fields!(plih, t.as_str(), - (id, "tvg-id"), - (group, "group-title"), - (name, "tvg-name"), - (parent_code, "parent-code"), - (audio_track, "audio-track"), - (logo, "tvg-logo"), - (logo_small, "tvg-logo-small"), - (time_shift, "timeshift"), - (rec, "tvg-rec"); value) - } + process_header_fields!(plih, t.as_str(), + (id, "tvg-id"), + (group, "group-title"), + (name, "tvg-name"), + (parent_code, "parent-code"), + (audio_track, "audio-track"), + (logo, "tvg-logo"), + (logo_small, "tvg-logo-small"), + (time_shift, "timeshift"), + (rec, "tvg-rec"); value) } } } diff --git a/src/processing/mod.rs b/src/processing/mod.rs index a291631df..cda7b468a 100644 --- a/src/processing/mod.rs +++ b/src/processing/mod.rs @@ -1,5 +1,7 @@ pub(crate) mod m3u_parser; pub(crate) mod xtream_parser; pub(crate) mod playlist_processor; -pub(crate) mod playlist_watch; -pub(crate) mod xmltv_parser; \ No newline at end of file +pub(crate) mod xmltv_parser; +mod playlist_watch; +mod xtream_processor; +mod affix_processor; \ No newline at end of file diff --git a/src/processing/playlist_processor.rs b/src/processing/playlist_processor.rs index 4a16433b8..bc0ae75f5 100644 --- a/src/processing/playlist_processor.rs +++ b/src/processing/playlist_processor.rs @@ -5,109 +5,54 @@ use std::collections::{HashMap, HashSet}; use std::rc::Rc; use std::sync::{Arc, Mutex}; use std::thread; -use actix_rt::System; +use actix_rt::System; use log::{debug, error, info, Level, log_enabled}; use unidecode::unidecode; -use crate::{Config, get_errors_notify_message, model::config, valid_property}; +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, InputAffix, InputType, ProcessTargets}; +use crate::model::config::{ConfigTarget, default_as_default, InputType, ProcessTargets}; +use crate::model::config::{ItemField, ProcessingOrder, SortOrder::{Asc, Desc}, TargetType}; use crate::model::mapping::{Mapping, MappingValueProcessor}; -use crate::model::config::{AFFIX_FIELDS, ItemField, ProcessingOrder, SortOrder::{Asc, Desc}, TargetType}; -use crate::model::playlist::{FetchedPlaylist, FieldAccessor, PlaylistGroup, PlaylistItem, PlaylistItemHeader}; + +use crate::model::playlist::{FetchedPlaylist, PlaylistGroup, PlaylistItem}; use crate::model::stats::{InputStats, PlaylistStats}; -use crate::model::xmltv::{Epg}; +use crate::model::xtream::MultiXtreamMapping; use crate::processing::playlist_watch::process_group_watch; use crate::processing::xmltv_parser::flatten_tvguide; -use crate::repository::epg_repository::write_epg; -use crate::repository::m3u_repository::{write_m3u_playlist, write_strm_playlist}; -use crate::repository::xtream_repository::write_xtream_playlist; +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) }; + target.filter(&provider) +} fn filter_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option> { debug!("Filtering {} groups", playlist.len()); let mut new_playlist = Vec::new(); playlist.iter_mut().for_each(|pg| { - if log_enabled!(Level::Debug) { - debug!("Filtering group {} with {} items", pg.title, pg.channels.len()); - } - let mut channels = Vec::new(); - pg.channels.iter_mut().for_each(|pli| { - if is_valid(pli, target) { - channels.push(pli.clone()); - } - }); - if log_enabled!(Level::Debug) { - debug!("Filtered group {} has now {} items", pg.title, channels.len()); - } + let channels = pg.channels.iter() + .filter(|&pli| is_valid(pli, target)).cloned().collect::>(); + debug!("Filtered group {} has now {}/{} items", pg.title, channels.len(), pg.channels.len()); if !channels.is_empty() { new_playlist.push(PlaylistGroup { id: pg.id, title: pg.title.clone(), channels, - xtream_cluster: pg.xtream_cluster.clone() + xtream_cluster: pg.xtream_cluster.clone(), }); } }); Some(new_playlist) } -fn apply_affixes(fetched_playlists: &mut [FetchedPlaylist]) { - fetched_playlists.iter_mut().for_each(|fetched_playlist| { - let FetchedPlaylist { input, playlist, epg: _ } = fetched_playlist; - if input.suffix.is_some() || input.prefix.is_some() { - let validate_affix = |a: &Option| match a { - Some(affix) => { - valid_property!(&affix.field.as_str(), AFFIX_FIELDS) && !affix.value.is_empty() - } - _ => false - }; - - let apply_prefix = validate_affix(&input.prefix); - let apply_suffix = validate_affix(&input.suffix); - - if apply_prefix || apply_suffix { - let get_affix_applied_value = |header: &mut PlaylistItemHeader, affix: &InputAffix, prefix: bool| { - if let Some(field_value) = header.get_field(affix.field.as_str()) { - return if prefix { - format!("{}{}", &affix.value, field_value.as_str()) - } else { - format!("{}{}", field_value.as_str(), &affix.value) - }; - } - String::from(&affix.value) - }; - - playlist.iter_mut().for_each(|group| { - group.channels.iter_mut().for_each(|channel| { - if apply_suffix { - if let Some(suffix) = &input.suffix { - let value = get_affix_applied_value(&mut channel.header.borrow_mut(), suffix, false); - if log_enabled!(Level::Debug) { - debug!("Applying input suffix: {}={}", &suffix.field, &value); - } - channel.header.borrow_mut().set_field(&suffix.field, value.as_str()); - } - } - if apply_prefix { - if let Some(prefix) = &input.prefix { - let value = get_affix_applied_value(&mut channel.header.borrow_mut(), prefix, true); - if log_enabled!(Level::Debug) { - debug!("Applying input prefix: {}={}", &prefix.field, &value); - } - channel.header.borrow_mut().set_field(&prefix.field, value.as_str()); - } - } - }); - }); - } - } - }); -} - fn sort_playlist(target: &ConfigTarget, new_playlist: &mut [PlaylistGroup]) { if let Some(sort) = &target.sort { let match_as_ascii = &sort.match_as_ascii; @@ -146,12 +91,6 @@ fn sort_playlist(target: &ConfigTarget, new_playlist: &mut [PlaylistGroup]) { } } - -fn is_valid(pli: &mut PlaylistItem, target: &ConfigTarget) -> bool { - let provider = ValueProvider { pli: RefCell::new(pli) }; - target.filter(&provider) -} - fn exec_rename(pli: &mut PlaylistItem, rename: &Option>) { if let Some(renames) = rename { if !renames.is_empty() { @@ -236,46 +175,33 @@ fn map_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option let new_playlist: Vec = playlist.iter().map(|playlist_group| { let mut grp = playlist_group.clone(); let mappings = target._mapping.as_ref().unwrap(); - mappings.iter().filter(|mapping| !mapping.mapper.is_empty()).for_each(|mapping| + mappings.iter().filter(|&mapping| !mapping.mapper.is_empty()).for_each(|mapping| grp.channels = grp.channels.drain(..).map(|chan| map_channel(chan, mapping)).collect()); grp }).collect(); // if the group names are changed, restructure channels to the right groups // we use - let mut max_group_id = 0; let mut new_groups: Vec = Vec::new(); + let mut grp_id: u32 = 0; for playlist_group in new_playlist { - let mut group_id_used = false; for channel in &playlist_group.channels { let cluster = &channel.header.borrow().xtream_cluster; let title = &channel.header.borrow().group; match new_groups.iter_mut().find(|x| *x.title == **title) { Some(grp) => grp.channels.push(channel.clone()), _ => { - let new_group_id = if group_id_used { - 0 - } else if *title == playlist_group.title { - group_id_used = true; - max_group_id = max_group_id.max(playlist_group.id); - playlist_group.id - } else { - 0 - }; + grp_id += 1; new_groups.push(PlaylistGroup { - id: new_group_id, + id: grp_id, title: Rc::clone(title), channels: vec![channel.clone()], - xtream_cluster: cluster.clone() + xtream_cluster: cluster.clone(), }) } } } } - new_groups.iter_mut().filter(|g| g.id == 0).for_each(|grp| { - max_group_id += 1; - grp.id = max_group_id; - }); Some(new_groups) } else { None @@ -372,7 +298,7 @@ async fn process_source(cfg: Arc, source_idx: usize, user_targets: Arc

, user_targets: Arc) -> (Vec, Vec) { +async fn process_sources(config: Arc, user_targets: Arc) -> (Vec, Vec) { let mut handle_list = vec![]; let thread_num = config.threads; let process_parallel = thread_num > 1 && config.sources.len() > 1; @@ -415,8 +341,7 @@ pub(crate) async fn process_sources(config: Arc, user_targets: Arc Option>>; +pub type ProcessingPipe = Vec Option>>; fn get_processing_pipe(target: &ConfigTarget) -> ProcessingPipe { match &target.processing_order { @@ -429,10 +354,10 @@ fn get_processing_pipe(target: &ConfigTarget) -> ProcessingPipe { } } -pub(crate) async fn process_playlist<'a>(playlists: &mut [FetchedPlaylist<'a>], - target: &ConfigTarget, cfg: &Config, - stats: &mut HashMap, - errors: &mut Vec) -> Result<(), Vec> { +async fn process_playlist<'a>(playlists: &mut [FetchedPlaylist<'a>], + target: &ConfigTarget, cfg: &Config, + stats: &mut HashMap, + errors: &mut Vec) -> Result<(), Vec> { let pipe = get_processing_pipe(target); if log_enabled!(Level::Debug) { debug!("Processing order is {}", &target.processing_order); @@ -452,32 +377,7 @@ pub(crate) async fn process_playlist<'a>(playlists: &mut [FetchedPlaylist<'a>], new_fpl.playlist = v; } } - let (resolve_series, resolve_series_delay) = - if let Some(options) = &target.options { - (options.xtream_resolve_series && fpl.input.input_type == InputType::Xtream && target.has_output(&TargetType::M3u), - options.xtream_resolve_series_delay) - } else { - (false, 0) - }; - if resolve_series { - let mut series_playlist = download::get_xtream_playlist_series(fpl, errors, resolve_series_delay).await; - // original content saved into original list - for plg in &series_playlist { - fpl.update_playlist(plg); - } - // run processing pipe over new items - for f in &pipe { - let r = f(&mut series_playlist, target); - if let Some(v) = r { - series_playlist = v; - } - } - // assign new items to the new playlist - for plg in &series_playlist { - new_fpl.update_playlist(plg); - } - } - + playlist_resolve_series(target, errors, &pipe, fpl, &mut new_fpl).await; // stats let input_stats = stats.get_mut(&new_fpl.input.id); if let Some(stat) = input_stats { @@ -491,6 +391,49 @@ pub(crate) 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![]; @@ -524,7 +467,7 @@ pub(crate) async fn process_playlist<'a>(playlists: &mut [FetchedPlaylist<'a>], let watch_re = target._watch_re.as_ref().unwrap(); new_playlist.iter().for_each(|pl| { if watch_re.iter().any(|r| r.is_match(&pl.title)) { - process_group_watch(cfg, &target.name, pl) + process_group_watch(cfg, &target.name, pl) } }); } @@ -537,30 +480,6 @@ pub(crate) async fn process_playlist<'a>(playlists: &mut [FetchedPlaylist<'a>], } } -fn persist_playlist(playlist: &mut [PlaylistGroup], epg: Option, - target: &ConfigTarget, cfg: &Config) -> Result<(), Vec> { - let mut errors = vec![]; - for output in &target.output { - match match output.target { - TargetType::M3u => write_m3u_playlist(target, cfg, playlist, &output.filename), - TargetType::Strm => write_strm_playlist(target, cfg, playlist, &output.filename), - TargetType::Xtream => write_xtream_playlist(target, cfg, playlist) - } { - Ok(_) => { - if !playlist.is_empty() { - match write_epg(target, cfg, &epg, output) { - Ok(_) => {} - Err(err) => errors.push(err) - } - } - } - Err(err) => errors.push(err) - } - } - - if errors.is_empty() { Ok(()) } else { Err(errors) } -} - pub(crate) async fn exec_processing(cfg: Arc, targets: Arc) { let (stats, errors) = process_sources(cfg.to_owned(), targets.to_owned()).await; let stats_msg = format!("{{\"stats\": {}}}", stats.iter().map(|stat| stat.to_string()).collect::>().join("\n")); @@ -572,7 +491,7 @@ pub(crate) async fn exec_processing(cfg: Arc, targets: Arc Option { let mut channel_ids: Vec<&String> = vec![]; tv_guides.iter().for_each(|guide| { if epg.attributes.is_none() { - epg.attributes = guide.attributes.clone(); + epg.attributes.clone_from(&guide.attributes); } guide.children.iter().for_each(|c| { if c.name.as_str() == "channel" { diff --git a/src/processing/xtream_parser.rs b/src/processing/xtream_parser.rs index 4e92c013e..2e49fe2d5 100644 --- a/src/processing/xtream_parser.rs +++ b/src/processing/xtream_parser.rs @@ -39,8 +39,8 @@ pub(crate) fn parse_xtream_series_info(info: &Value, group_title: &str, input: & let result: Vec = series_info.episodes.values().flatten().map(|episode| PlaylistItem { header: RefCell::new(PlaylistItemHeader { - stream_id: 0, id: Rc::new(episode.id.to_owned()), + stream_id: Rc::new(episode.id.to_owned()), name: Rc::new(episode.title.to_owned()), logo: Rc::new(episode.info.movie_image.to_owned()), logo_small: default_as_empty_rc_str(), @@ -97,8 +97,8 @@ pub(crate) fn parse_xtream(input: &ConfigInput, let title = &grp.category_name; let item = PlaylistItem { header: RefCell::new(PlaylistItemHeader { - stream_id: stream.stream_id.unwrap_or(0) as u32, id: Rc::new(stream.get_stream_id()), + stream_id: Rc::new(stream.get_stream_id()), name: Rc::clone(&stream.name), logo: Rc::clone(&stream.stream_icon), logo_small: default_as_empty_rc_str(), @@ -143,7 +143,7 @@ pub(crate) fn parse_xtream(input: &ConfigInput, Ok(Some(group_map.values().map(|category| { let cat = category.borrow(); PlaylistGroup { - id: 0, + id: cat.category_id.parse::().unwrap_or(0), xtream_cluster: xtream_cluster.clone(), title: Rc::clone(&cat.category_name), channels: cat.channels.clone() diff --git a/src/processing/xtream_processor.rs b/src/processing/xtream_processor.rs new file mode 100644 index 000000000..039a648b8 --- /dev/null +++ b/src/processing/xtream_processor.rs @@ -0,0 +1,36 @@ +use crate::m3u_filter_error::M3uFilterError; +use crate::model::config::{ConfigTarget, InputType, TargetType}; +use crate::model::playlist::FetchedPlaylist; +use crate::processing::playlist_processor::ProcessingPipe; +use crate::utils::download; + +pub async fn playlist_resolve_series<'a>(target: &ConfigTarget, errors: &mut Vec, + pipe: &ProcessingPipe, + fpl: &mut FetchedPlaylist<'_>, + new_fpl: &mut FetchedPlaylist<'_>) { + let (resolve_series, resolve_series_delay) = + if let Some(options) = &target.options { + (options.xtream_resolve_series && fpl.input.input_type == InputType::Xtream && target.has_output(&TargetType::M3u), + options.xtream_resolve_series_delay) + } else { + (false, 0) + }; + if resolve_series { + let mut series_playlist = download::get_xtream_playlist_series(fpl, errors, resolve_series_delay).await; + // original content saved into original list + for plg in &series_playlist { + fpl.update_playlist(plg); + } + // run processing pipe over new items + for f in pipe { + let r = f(&mut series_playlist, target); + if let Some(v) = r { + series_playlist = v; + } + } + // assign new items to the new playlist + for plg in &series_playlist { + new_fpl.update_playlist(plg); + } + } +} \ No newline at end of file diff --git a/src/repository/kodi_repository.rs b/src/repository/kodi_repository.rs new file mode 100644 index 000000000..4897be6b0 --- /dev/null +++ b/src/repository/kodi_repository.rs @@ -0,0 +1,129 @@ +use std::fs::File; +use std::io::Write; +use chrono::Datelike; +use log::error; +use crate::create_m3u_filter_error_result; +use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; +use crate::model::config::{Config, ConfigTarget}; +use crate::model::playlist::PlaylistGroup; +use crate::utils::file_utils; + +struct KodiStyle { + year: regex::Regex, + season: regex::Regex, + episode: regex::Regex, + whitespace: regex::Regex, +} + +fn sanitize_for_filename(text: &str, underscore_whitespace: bool) -> String { + return text.chars().filter(|c| c.is_alphanumeric() || c.is_whitespace()) + .map(|c| if underscore_whitespace { if c.is_whitespace() { '_' } else { c } } else { c }) + .collect::(); +} + +fn kodi_style_rename_year(name: &String, style: &KodiStyle) -> (String, Option) { + let current_date = chrono::Utc::now(); + let cur_year = current_date.year(); + match style.year.find(name) { + Some(m) => { + let s_year = &name[m.start()..m.end()]; + let t_year: i32 = s_year.parse().unwrap(); + if t_year > 1900 && t_year <= cur_year { + let new_name = format!("{}{}", &name[0..m.start()], &name[m.end()..]); + return (new_name, Some(String::from(s_year))); + } + (String::from(name), Some(cur_year.to_string())) + } + _ => (String::from(name), Some(cur_year.to_string())), + } +} + +fn kodi_style_rename_season(name: &String, style: &KodiStyle) -> (String, Option) { + match style.season.find(name) { + Some(m) => { + let s_season = &name[m.start()..m.end()]; + let season = Some(String::from(&s_season[1..])); + let new_name = format!("{}{}", &name[0..m.start()], &name[m.end()..]); + (new_name, season) + } + _ => (String::from(name), Some(String::from("01"))), + } +} + +fn kodi_style_rename_episode(name: &String, style: &KodiStyle) -> (String, Option) { + match style.episode.find(name) { + Some(m) => { + let s_episode = &name[m.start()..m.end()]; + let episode = Some(String::from(&s_episode[1..])); + let new_name = format!("{}{}", &name[0..m.start()], &name[m.end()..]); + (new_name, episode) + } + _ => (String::from(name), None), + } +} + +fn kodi_style_rename(name: &String, style: &KodiStyle) -> String { + let (work_name_1, year) = kodi_style_rename_year(name, style); + let (work_name_2, season) = kodi_style_rename_season(&work_name_1, style); + let (work_name_3, episode) = kodi_style_rename_episode(&work_name_2, style); + if year.is_some() && season.is_some() && episode.is_some() { + let formatted = format!("{} ({}) S{}E{}", work_name_3, year.unwrap(), season.unwrap(), episode.unwrap()); + return String::from(style.whitespace.replace_all(formatted.as_str(), " ").as_ref()); + } + String::from(name) +} + + +pub(crate) fn write_strm_playlist(target: &ConfigTarget, cfg: &Config, new_playlist: &[PlaylistGroup], filename: &Option) -> Result<(), M3uFilterError> { + if !new_playlist.is_empty() { + if filename.is_none() { + return Err(M3uFilterError::new(M3uFilterErrorKind::Notify, "write strm playlist failed: ".to_string())); + } + let underscore_whitespace = target.options.as_ref().map_or(false, |o| o.underscore_whitespace); + let cleanup = target.options.as_ref().map_or(false, |o| o.cleanup); + let kodi_style = target.options.as_ref().map_or(false, |o| o.kodi_style); + + if let Some(path) = file_utils::get_file_path(&cfg.working_dir, Some(std::path::PathBuf::from(&filename.as_ref().unwrap()))) { + if cleanup { + let _ = std::fs::remove_dir_all(&path); + } + if let Err(e) = std::fs::create_dir_all(&path) { + error!("cant create directory: {:?}", &path); + return create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "failed to write strm playlist: {}", e); + }; + for pg in new_playlist { + for pli in &pg.channels { + let header = &pli.header.borrow(); + let dir_path = path.join(sanitize_for_filename(&header.group, underscore_whitespace)); + if let Err(e) = std::fs::create_dir_all(&dir_path) { + error!("cant create directory: {:?}", &path); + return create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "failed to write strm playlist: {}", e); + }; + let mut file_name = sanitize_for_filename(&header.title, underscore_whitespace); + if kodi_style { + let style = KodiStyle { + season: regex::Regex::new(r"[Ss]\d\d").unwrap(), + episode: regex::Regex::new(r"[Ee]\d\d").unwrap(), + year: regex::Regex::new(r"\d\d\d\d").unwrap(), + whitespace: regex::Regex::new(r"\s+").unwrap(), + }; + file_name = kodi_style_rename(&file_name, &style); + } + let file_path = dir_path.join(format!("{}.strm", file_name)); + match File::create(&file_path) { + Ok(mut strm_file) => { + match file_utils::check_write(strm_file.write_all(header.url.as_bytes())) { + Ok(_) => (), + Err(e) => return create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "failed to write strm playlist: {}", e), + } + } + Err(err) => { + return create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "failed to write strm playlist: {}", err); + } + } + } + } + } + } + Ok(()) +} \ No newline at end of file diff --git a/src/repository/m3u_repository.rs b/src/repository/m3u_repository.rs index 938c6ff4f..64c3cafde 100644 --- a/src/repository/m3u_repository.rs +++ b/src/repository/m3u_repository.rs @@ -3,17 +3,15 @@ use std::io::{Error, ErrorKind, Write, SeekFrom, Seek, Read}; use std::path::Path; use std::rc::Rc; -use chrono::Datelike; use log::error; use crate::{create_m3u_filter_error_result}; use crate::api::api_utils::get_user_server_info; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; -use crate::model::api_proxy::ProxyUserCredentials; +use crate::model::api_proxy::{ProxyType, ProxyUserCredentials}; use crate::model::config::{Config, ConfigTarget}; -use crate::model::playlist::{PlaylistGroup, PlaylistItemType}; -use crate::processing::m3u_parser::consume_m3u; -use crate::utils::file_reader::FileReader; +use crate::model::playlist::{M3uPlaylistItem, PlaylistGroup, PlaylistItemType}; +use crate::repository::repository_utils::IndexRecord; use crate::utils::file_utils; macro_rules! cant_write_result { @@ -22,85 +20,12 @@ macro_rules! cant_write_result { } } -fn check_write(res: std::io::Result<()>) -> Result<(), std::io::Error> { - match res { - Ok(_) => Ok(()), - Err(_) => Err(std::io::Error::new(std::io::ErrorKind::Other, "Unable to write file")), - } -} - -fn sanitize_for_filename(text: &str, underscore_whitespace: bool) -> String { - return text.chars().filter(|c| c.is_alphanumeric() || c.is_whitespace()) - .map(|c| if underscore_whitespace { if c.is_whitespace() { '_' } else { c } } else { c }) - .collect::(); -} - -struct KodiStyle { - year: regex::Regex, - season: regex::Regex, - episode: regex::Regex, - whitespace: regex::Regex, -} - -fn kodi_style_rename_year(name: &String, style: &KodiStyle) -> (String, Option) { - let current_date = chrono::Utc::now(); - let cur_year = current_date.year(); - match style.year.find(name) { - Some(m) => { - let s_year = &name[m.start()..m.end()]; - let t_year: i32 = s_year.parse().unwrap(); - if t_year > 1900 && t_year <= cur_year { - let new_name = format!("{}{}", &name[0..m.start()], &name[m.end()..]); - return (new_name, Some(String::from(s_year))); - } - (String::from(name), Some(cur_year.to_string())) - } - _ => (String::from(name), Some(cur_year.to_string())), - } -} - -fn kodi_style_rename_season(name: &String, style: &KodiStyle) -> (String, Option) { - match style.season.find(name) { - Some(m) => { - let s_season = &name[m.start()..m.end()]; - let season = Some(String::from(&s_season[1..])); - let new_name = format!("{}{}", &name[0..m.start()], &name[m.end()..]); - (new_name, season) - } - _ => (String::from(name), Some(String::from("01"))), - } -} - -fn kodi_style_rename_episode(name: &String, style: &KodiStyle) -> (String, Option) { - match style.episode.find(name) { - Some(m) => { - let s_episode = &name[m.start()..m.end()]; - let episode = Some(String::from(&s_episode[1..])); - let new_name = format!("{}{}", &name[0..m.start()], &name[m.end()..]); - (new_name, episode) - } - _ => (String::from(name), None), - } -} - -fn kodi_style_rename(name: &String, style: &KodiStyle) -> String { - let (work_name_1, year) = kodi_style_rename_year(name, style); - let (work_name_2, season) = kodi_style_rename_season(&work_name_1, style); - let (work_name_3, episode) = kodi_style_rename_episode(&work_name_2, style); - if year.is_some() && season.is_some() && episode.is_some() { - let formatted = format!("{} ({}) S{}E{}", work_name_3, year.unwrap(), season.unwrap(), episode.unwrap()); - return String::from(style.whitespace.replace_all(formatted.as_str(), " ").as_ref()); - } - String::from(name) -} - -pub(crate) fn get_m3u_file_paths(cfg: &Config, filename: &Option) -> Option<(std::path::PathBuf, std::path::PathBuf, std::path::PathBuf)> { +pub(crate) fn get_m3u_file_paths(cfg: &Config, filename: &Option) -> Option<(std::path::PathBuf, std::path::PathBuf)> { match file_utils::get_file_path(&cfg.working_dir, Some(std::path::PathBuf::from(&filename.as_ref().unwrap()))) { Some(m3u_path) => { let extension = m3u_path.extension().map(|ext| format!("{}_", ext.to_str().unwrap_or(""))).unwrap_or("".to_owned()); - let url_path = m3u_path.with_extension(format!("{}url", &extension)); - let index_path = m3u_path.with_extension(format!("{}idx_url", &extension)); - Some((m3u_path, url_path, index_path)) + let index_path = m3u_path.with_extension(format!("{}idx", &extension)); + Some((m3u_path, index_path)) } None => None } @@ -111,17 +36,12 @@ 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, url_path: &Path, idx_path: &Path) -> Result<(File, File, File), M3uFilterError> { +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(url_path) { - Ok(url_file) => { - match File::create(idx_path) { - Ok(idx_file) => Ok((m3u_file, url_file, idx_file)), - Err(e) => cant_write_result!(&idx_path, e), - } - } - Err(e) => cant_write_result!(&url_path, e), + 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), @@ -135,105 +55,37 @@ pub(crate) fn write_m3u_playlist(target: &ConfigTarget, cfg: &Config, new_playli M3uFilterErrorKind::Notify, format!("write m3u playlist for target {} failed: No filename set", target.name))); } - if let Some((m3u_path, url_path, idx_path)) = get_m3u_file_paths(cfg, filename) { - match create_m3u_files(&m3u_path, &url_path, &idx_path) { - Ok((mut m3u_file, mut m3u_url_file, mut m3u_idx_file)) => { - match check_write(m3u_file.write_all(b"#EXTM3U\n")) { - Ok(_) => (), - Err(e) => return cant_write_result!(&m3u_path, e), - } - let mut id_counter: u32 = 0; + if let Some((m3u_path, idx_path)) = get_m3u_file_paths(cfg, filename) { + match create_m3u_files(&m3u_path, &idx_path) { + Ok((mut m3u_file, mut m3u_idx_file)) => { let mut idx_offset: u32 = 0; - for pg in new_playlist { - for pli in &pg.channels { - if pli.header.borrow().item_type == PlaylistItemType::SeriesInfo { - // we skip series info, because this is only necessary when writing xtream files - continue; - } - let url = pli.header.borrow().url.as_str().to_owned(); - let url_bytes = url.as_bytes(); - let bytes_to_write: u32 = url_bytes.len() as u32; - match check_write(m3u_url_file.write_all(url_bytes)) { + let m3u_playlist = new_playlist.iter() + .flat_map(|pg| &pg.channels) + .filter(|&pli| pli.header.borrow().item_type != PlaylistItemType::SeriesInfo) + .map(|pli| pli.to_m3u()).collect::>(); + let mut stream_id: u32 = 1; + for mut m3u in m3u_playlist { + m3u.stream_id = Rc::new(stream_id.to_string()); + if let Ok(encoded) = bincode::serialize(&m3u) { + match file_utils::check_write(m3u_file.write_all(&encoded)) { Ok(_) => { - match check_write(m3u_idx_file.write_all(&idx_offset.to_le_bytes())) { - Ok(_) => (), - Err(e) => return cant_write_result!(&m3u_path, e), + let bytes_written = encoded.len() as u32; + let combined_bytes: [u8; 8] = 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); } - idx_offset += bytes_to_write; + idx_offset += bytes_written; + stream_id += 1; + } + Err(err) => { + return cant_write_result!(&m3u_path, err); } - Err(e) => return cant_write_result!(&m3u_path, e), - } - - id_counter += 1; - pli.header.borrow_mut().stream_id = id_counter; - let content = pli.to_m3u(target); - match check_write(m3u_file.write_all(content.as_bytes())) { - Ok(_) => (), - Err(e) => return cant_write_result!(&m3u_path, e), - } - match check_write(m3u_file.write_all(b"\n")) { - Ok(_) => (), - Err(e) => return cant_write_result!(&m3u_path, e), } } } } - Err(e) => return cant_write_result!(&m3u_path, e), - } - } - } - Ok(()) -} - -pub(crate) fn write_strm_playlist(target: &ConfigTarget, cfg: &Config, new_playlist: &[PlaylistGroup], filename: &Option) -> Result<(), M3uFilterError> { - if !new_playlist.is_empty() { - if filename.is_none() { - return Err(M3uFilterError::new(M3uFilterErrorKind::Notify, "write strm playlist failed: ".to_string())); - } - let underscore_whitespace = target.options.as_ref().map_or(false, |o| o.underscore_whitespace); - let cleanup = target.options.as_ref().map_or(false, |o| o.cleanup); - let kodi_style = target.options.as_ref().map_or(false, |o| o.kodi_style); - - if let Some(path) = file_utils::get_file_path(&cfg.working_dir, Some(std::path::PathBuf::from(&filename.as_ref().unwrap()))) { - if cleanup { - let _ = std::fs::remove_dir_all(&path); - } - if let Err(e) = std::fs::create_dir_all(&path) { - error!("cant create directory: {:?}", &path); - return create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "failed to write strm playlist: {}", e); - }; - for pg in new_playlist { - for pli in &pg.channels { - let header = &pli.header.borrow(); - let dir_path = path.join(sanitize_for_filename(&header.group, underscore_whitespace)); - if let Err(e) = std::fs::create_dir_all(&dir_path) { - error!("cant create directory: {:?}", &path); - return create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "failed to write strm playlist: {}", e); - }; - let mut file_name = sanitize_for_filename(&header.title, underscore_whitespace); - if kodi_style { - let style = KodiStyle { - season: regex::Regex::new(r"[Ss]\d\d").unwrap(), - episode: regex::Regex::new(r"[Ee]\d\d").unwrap(), - year: regex::Regex::new(r"\d\d\d\d").unwrap(), - whitespace: regex::Regex::new(r"\s+").unwrap(), - }; - file_name = kodi_style_rename(&file_name, &style); - } - let file_path = dir_path.join(format!("{}.strm", file_name)); - match File::create(&file_path) { - Ok(mut strm_file) => { - match check_write(strm_file.write_all(header.url.as_bytes())) { - Ok(_) => (), - Err(e) => return create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "failed to write strm playlist: {}", e), - } - } - Err(err) => { - return create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "failed to write strm playlist: {}", err); - } - } - } + Err(err) => return cant_write_result!(&m3u_path, err), } } } @@ -243,70 +95,68 @@ pub(crate) fn write_strm_playlist(target: &ConfigTarget, cfg: &Config, new_playl pub(crate) fn 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, url_path, idx_path)) = get_m3u_file_paths(cfg, &filename) { - if m3u_path.exists() && url_path.exists() && idx_path.exists() { - let server_info = get_user_server_info(cfg, user); - let url = format!("{}/m3u-stream/{}/{}", server_info.get_base_url(), user.username, user.password); - let mut result = vec![]; - result.push("#EXTM3U".to_string()); - match File::open(&m3u_path) { - Ok(m3u_file) => { - let reader = FileReader::new(m3u_file); - consume_m3u(cfg, reader, |item| { - let stream_id = item.header.borrow().stream_id; - item.header.borrow_mut().url = Rc::new(format!("{}/{}", url, stream_id)); - result.push(item.to_m3u(target)) - }); - } - Err(err) => { - error!("Could not open file {}: {}", &m3u_path.to_str().unwrap(), err); - return None; - } - } - return Some(result.join("\n")); - } else if m3u_path.exists() { - match std::fs::read_to_string(&m3u_path) { - Ok(result) => return Some(result), - Err(err) => { - error!("Could not open file {}: {}", &m3u_path.to_str().unwrap(), err); - return None; - } + if let Some((m3u_path, idx_path)) = get_m3u_file_paths(cfg, &filename) { + if m3u_path.exists() && idx_path.exists() { + match std::fs::read(&m3u_path) { + Ok(encoded_m3u) => { + match std::fs::read(&idx_path) { + Ok(encoded_idx) => { + let mut cursor = 0; + let size = encoded_idx.len(); + + let server_info = get_user_server_info(cfg, user); + let url = format!("{}/m3u-stream/{}/{}", server_info.get_base_url(), user.username, user.password); + let mut result = vec![]; + result.push("#EXTM3U".to_string()); + + 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(m3u_pli) = bincode::deserialize::(&encoded_m3u[start_offset..end_offset]) { + match user.proxy { + ProxyType::Reverse => { + let stream_id = Rc::clone(&m3u_pli.stream_id); + result.push(m3u_pli.to_m3u(target, Some(format!("{}/{}", url, stream_id).as_str()))); + } + ProxyType::Redirect => { + result.push(m3u_pli.to_m3u(target, None)); + } + } + } else { + error!("Could not deserialize item {}", &m3u_path.to_str().unwrap()); + } + } + return Some(result.join("\n")); + }, + Err(err) => error!("Could not open file {}: {}", &idx_path.to_str().unwrap(), err), + } + }, + Err(err) => error!("Could not open file {}: {}", &m3u_path.to_str().unwrap(), err), } } + } else { + error!("Could not open file for target {}", &target.name); } } None } -pub(crate) fn get_m3u_url_for_stream_id(stream_id: u32, url_path: &Path, idx_path: &Path) -> Result { +pub(crate) fn get_m3u_item_for_stream_id(stream_id: u32, m3u_path: &Path, idx_path: &Path) -> Result { if stream_id < 1 { return Err(Error::new(ErrorKind::Other, "id should start with 1")); } - if url_path.exists() && idx_path.exists() { - let index = (stream_id - 1) as u64; - let mapping_size = 4; // u32 - let offset = mapping_size * index; + if m3u_path.exists() && idx_path.exists() { + let offset: u64 = IndexRecord::get_index_offset(stream_id - 1) as u64; let mut idx_file = File::open(idx_path)?; - idx_file.seek(SeekFrom::Start(offset))?; - let mut url_offset_bytes = [0u8; 4]; - idx_file.read_exact(&mut url_offset_bytes)?; - let url_offset = u32::from_le_bytes(url_offset_bytes); - let url_end_offset = match idx_file.read_exact(&mut url_offset_bytes) { - Ok(_) => u32::from_le_bytes(url_offset_bytes), - Err(_) => 0 - }; - - let mut url_file = File::open(url_path)?; - url_file.seek(SeekFrom::Start(url_offset as u64))?; - let mut buf = String::new(); - - if (if url_end_offset > 0 { - url_file.take((url_end_offset - url_offset) as u64).read_to_string(&mut buf) - } else { - url_file.read_to_string(&mut buf) - }).is_ok() { - return Ok(buf); - }; + let mut m3u_file = File::open(m3u_path)?; + let index_record = IndexRecord::from_file(&mut idx_file, offset)?; + m3u_file.seek(SeekFrom::Start(index_record.index as u64))?; + let mut buffer : Vec = vec![0; index_record.size as usize]; + m3u_file.read_exact(&mut buffer)?; + if let Ok(m3u_pli) = bincode::deserialize::(&buffer[..]) { + return Ok(m3u_pli); + } } - Err(Error::new(ErrorKind::Other, format!("Failed to read m3u url form stream-id {}", stream_id))) + Err(Error::new(ErrorKind::Other, format!("Failed to read m3u for stream-id {}", stream_id))) } \ No newline at end of file diff --git a/src/repository/mod.rs b/src/repository/mod.rs index 1d070961a..f66d0374f 100644 --- a/src/repository/mod.rs +++ b/src/repository/mod.rs @@ -1,3 +1,7 @@ +pub(crate) mod playlist_repository; pub(crate) mod m3u_repository; pub(crate) mod xtream_repository; -pub(crate) mod epg_repository; \ No newline at end of file +pub(crate) mod epg_repository; +pub(crate)mod kodi_repository; + +mod repository_utils; \ No newline at end of file diff --git a/src/repository/playlist_repository.rs b/src/repository/playlist_repository.rs new file mode 100644 index 000000000..2169fa869 --- /dev/null +++ b/src/repository/playlist_repository.rs @@ -0,0 +1,32 @@ +use crate::m3u_filter_error::M3uFilterError; +use crate::model::config::{Config, ConfigTarget, TargetType}; +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::xtream_repository::write_xtream_playlist; + +pub(crate) fn persist_playlist(playlist: &mut [PlaylistGroup], epg: Option, + target: &ConfigTarget, cfg: &Config) -> Result<(), Vec> { + let mut errors = vec![]; + for output in &target.output { + match match output.target { + TargetType::M3u => write_m3u_playlist(target, cfg, playlist, &output.filename), + TargetType::Strm => write_strm_playlist(target, cfg, playlist, &output.filename), + TargetType::Xtream => write_xtream_playlist(target, cfg, playlist) + } { + Ok(_) => { + if !playlist.is_empty() { + match write_epg(target, cfg, &epg, output) { + Ok(_) => {} + Err(err) => errors.push(err) + } + } + } + Err(err) => errors.push(err) + } + } + + if errors.is_empty() { Ok(()) } else { Err(errors) } +} \ No newline at end of file diff --git a/src/repository/repository_utils.rs b/src/repository/repository_utils.rs new file mode 100644 index 000000000..d657c00de --- /dev/null +++ b/src/repository/repository_utils.rs @@ -0,0 +1,54 @@ +use std::convert::TryInto; +use std::io::Error; +use std::fs::File; +use std::io::{Read, SeekFrom, Seek}; + +/** +We write the structs with bincode::encode to a file. +To access each entry we need a index file where we can find the +Entries with offset and size of the encoded struct. +This is used for the index file where first entry is the index +of the encoded file, and size is the size of the encoded struct. +*/ +pub struct IndexRecord { + pub index: u32, + pub size: u32, +} + +impl IndexRecord { + pub fn new(index: u32, size: u32) -> IndexRecord { + IndexRecord { index, size } + } + + 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]; + 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); + Ok(IndexRecord { index, size }) + } + + pub fn from_bytes(bytes: &[u8], cursor: &mut usize) -> IndexRecord { + 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(); + *cursor += 4; + let index = u32::from_le_bytes(index_bytes); + let size = u32::from_le_bytes(size_bytes); + IndexRecord { index, size } + } + + pub fn to_bytes(&self) -> [u8; 8] { + 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]; + 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 } +} \ No newline at end of file diff --git a/src/repository/xtream_repository.rs b/src/repository/xtream_repository.rs index c9c48a16c..b59ae72cc 100644 --- a/src/repository/xtream_repository.rs +++ b/src/repository/xtream_repository.rs @@ -1,8 +1,8 @@ use std::cell::Ref; use std::collections::{BTreeMap, HashMap}; -use std::fs; +use std::{fs, io}; use std::fs::{File, OpenOptions}; -use std::io::{BufReader, BufWriter, Error, Read, Seek, SeekFrom, Write}; +use std::io::{BufReader, BufWriter, Error, ErrorKind, Read, Seek, SeekFrom, Write}; use std::iter::FromIterator; use std::path::{Path, PathBuf}; use log::{error}; @@ -13,12 +13,12 @@ use crate::model::playlist::{PlaylistGroup, PlaylistItemHeader, PlaylistItemType 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::utils::file_utils; use crate::utils::json_utils::iter_json_array; type IndexTree = BTreeMap; - pub(crate) static COL_CAT_LIVE: &str = "cat_live"; pub(crate) static COL_CAT_SERIES: &str = "cat_series"; pub(crate) static COL_CAT_VOD: &str = "cat_vod"; @@ -41,7 +41,6 @@ const SERIES_STREAM_FIELDS: &[&str] = &[ "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(' ', "_")))) } @@ -114,161 +113,169 @@ fn write_xtream_info(app_state: &AppState, target_name: &str, stream_id: i32, cl Ok(()) } -pub(crate) fn write_xtream_playlist(target: &ConfigTarget, cfg: &Config, playlist: &mut [PlaylistGroup]) -> Result<(), M3uFilterError> { - if let Some(path) = get_xtream_storage_path(cfg, &target.name) { +fn ensure_xtream_storage_path(cfg: &Config, target_name: &str) -> Result { + if let Some(path) = get_xtream_storage_path(cfg, target_name) { if fs::create_dir_all(&path).is_err() { - let msg = format!("Failed to save, can't create directory {}", &path.to_str().unwrap()); + let msg = format!("Failed to save xtream data, can't create directory {}", &path.to_str().unwrap()); return Err(M3uFilterError::new(M3uFilterErrorKind::Notify, msg)); } + Ok(path) + } else { + let msg = format!("Failed to save xtream data, can't create directory for target {target_name}"); + Err(M3uFilterError::new(M3uFilterErrorKind::Notify, msg)) + } +} - 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)); +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 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 vod_map = HashMap::::new(); + let mut series_map = HashMap::::new(); - let mut channel_num: i32 = 0; - let mut errors = Vec::new(); + let mut channel_num: i32 = 0; + let mut errors = Vec::new(); - // preserve category_ids - let (max_cat_id, existing_cat_ids) = load_old_category_ids(&path); - let mut cat_id_counter = max_cat_id; - for plg in playlist.iter_mut() { - if !&plg.channels.is_empty() { - let cat_id = existing_cat_ids.get(plg.title.as_ref()).unwrap_or_else(|| { - cat_id_counter += 1; - &cat_id_counter - }); - plg.id = *cat_id; + // preserve category_ids + let (max_cat_id, existing_cat_ids) = load_old_category_ids(&path); + let mut cat_id_counter = max_cat_id; + for plg in playlist.iter_mut() { + if !&plg.channels.is_empty() { + let cat_id = existing_cat_ids.get(plg.title.as_ref()).unwrap_or_else(|| { + cat_id_counter += 1; + &cat_id_counter + }); + plg.id = *cat_id; - match &plg.xtream_cluster { - XtreamCluster::Live => &mut cat_live_col, - XtreamCluster::Series => &mut cat_series_col, - XtreamCluster::Video => &mut cat_vod_col, - }.push( - json!({ + match &plg.xtream_cluster { + 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 })); - 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; - } - 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())), - ]); + 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; + } + 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())); + 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()) + }); } - 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())); + 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())); } - document.insert("custom_sid".to_string(), Value::String("".to_string())); - } - XtreamCluster::Series => { - document.insert("series_id".to_string(), stream_id_value); - } - }; + 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()); + 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)); } - - 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)); } } } - } - 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) { - Ok(()) => {} - Err(err) => { - errors.push(format!("Persisting collection failed: {}: {}", &col_path.to_str().unwrap(), err)); + 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) { + Ok(()) => {} + Err(err) => { + errors.push(format!("Persisting collection failed: {}: {}", &col_path.to_str().unwrap(), err)); + } } } + if !errors.is_empty() { + return create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "{}", errors.join("\n")); + } } - if !errors.is_empty() { - return create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "{}", errors.join("\n")); - } - } else { - return create_m3u_filter_error_result!(M3uFilterErrorKind::Notify, "Persisting playlist failed: {}.db", &target.name); + Err(err) => return Err(err) } - Ok(()) } @@ -282,7 +289,7 @@ fn load_old_category_ids(path: &Path) -> (u32, HashMap) { if col_path.exists() { if let Ok(file) = File::open(col_path) { let reader = BufReader::new(file); - for entry in iter_json_array::>(reader).flatten() { + for entry in iter_json_array::>(reader).flatten() { if let Some(item) = entry.as_object() { if let Some(category_id_value) = item.get("category_id") { if let Some(category_id) = category_id_value.as_str() { @@ -357,7 +364,7 @@ pub(crate) fn xtream_get_collection_path(cfg: &Config, target_name: &str, collec return Ok((Some(col_path), None)); } } - Err(Error::new(std::io::ErrorKind::Other, format!("Cant find collection: {}/{}", target_name, collection_name))) + Err(Error::new(ErrorKind::Other, format!("Cant find collection: {}/{}", target_name, collection_name))) } fn load_index(path: &Path) -> Option { @@ -370,7 +377,7 @@ fn load_index(path: &Path) -> Option { } } -fn write_index(path: &PathBuf, index: &IndexTree) -> std::io::Result<()> { +fn write_index(path: &PathBuf, index: &IndexTree) -> io::Result<()> { let encoded = bincode::serialize(index).unwrap(); fs::write(path, encoded) } @@ -404,10 +411,10 @@ pub(crate) async fn xtream_get_stored_stream_info( if let Some((offset, size)) = idx_map.get(&stream_id) { let mut reader = BufReader::new(File::open(&col_path).unwrap()); if let Ok(bytes) = seek_read(&mut reader, *offset as u64, *size) { - let mut decomp: Vec = Vec::new(); - let _ = lzma_rs::lzma_decompress(&mut bytes.as_slice(), &mut decomp); + let mut decompressed: Vec = Vec::new(); + let _ = lzma_rs::lzma_decompress(&mut bytes.as_slice(), &mut decompressed); drop(shared_lock); - return Ok(String::from_utf8(decomp).unwrap()); + return Ok(String::from_utf8(decompressed).unwrap()); } } } @@ -443,4 +450,55 @@ pub(crate) async fn xtream_persist_stream_info( drop(shared_lock); } } +} + +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 })); + } + } + Ok(None) } \ No newline at end of file diff --git a/src/test.rs b/src/test.rs index ad0ccbfd2..1b484ab4b 100644 --- a/src/test.rs +++ b/src/test.rs @@ -1,6 +1,8 @@ #[cfg(test)] mod tests { use crate::filter::get_filter; + use crate::model::xtream::MultiXtreamMapping; + use crate::repository::xtream_repository::{read_xtream_mapping, write_xtream_mapping}; #[test] fn test_filter() { diff --git a/src/utils/config_reader.rs b/src/utils/config_reader.rs index 3498d3279..582d50b32 100644 --- a/src/utils/config_reader.rs +++ b/src/utils/config_reader.rs @@ -30,7 +30,7 @@ pub(crate) fn read_mappings(args_mapping: Option, cfg: &mut Config) -> R pub(crate) fn read_api_proxy_config(args_api_proxy_config: Option, cfg: &mut Config) { let api_proxy_config_file: String = args_api_proxy_config.unwrap_or(file_utils::get_default_api_proxy_config_path(cfg._config_path.as_str())); - cfg._api_proxy_file_path = api_proxy_config_file.to_owned(); + api_proxy_config_file.clone_into(&mut cfg._api_proxy_file_path); let api_proxy_config = read_api_proxy(api_proxy_config_file.as_str(), true); match api_proxy_config { None => { diff --git a/src/utils/file_utils.rs b/src/utils/file_utils.rs index 45e525f23..ce7aeb910 100644 --- a/src/utils/file_utils.rs +++ b/src/utils/file_utils.rs @@ -77,6 +77,7 @@ pub(crate) fn get_working_path(wd: &String) -> String { String::from(current_dir.to_str().unwrap_or(".")) } else { let work_path = std::path::PathBuf::from(wd); + let _ = fs::create_dir_all(&work_path); let wdpath = match fs::metadata(&work_path) { Ok(md) => { if md.is_dir() && !md.permissions().readonly() { @@ -164,3 +165,10 @@ pub(crate) fn path_exists(file_path: &Path) -> bool { } false } + +pub(crate) fn check_write(res: std::io::Result<()>) -> Result<(), std::io::Error> { + match res { + Ok(_) => Ok(()), + Err(_) => Err(std::io::Error::new(std::io::ErrorKind::Other, "Unable to write file")), + } +} diff --git a/src/utils/mod.rs b/src/utils/mod.rs index 9fb33a6ed..aa7bb018c 100644 --- a/src/utils/mod.rs +++ b/src/utils/mod.rs @@ -5,4 +5,4 @@ pub (crate) mod string_utils; pub (crate) mod json_utils; pub (crate) mod config_reader; pub (crate) mod multi_file_reader; -pub (crate) mod file_reader; \ No newline at end of file +// pub (crate) mod file_reader; \ No newline at end of file