diff --git a/CHANGELOG.md b/CHANGELOG.md index 8841a6cc7..91ae9a9e0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -23,7 +23,7 @@ server: - Added Active clients count (for reverse proxy mode users) which is now displayed in `/status` and can be logged with setting `active_clients: true` under `log`section in `config.yml` - Fixed iptv player using live tv stream without `/live/` context. -- Added `log_level`to `log` config. Priority: CLI-Argument, Env-Var, Config, Default(`info`) +- Added `log_level` to `log` config. Priority: CLI-Argument, Env-Var, Config, Default(`info`) ```yaml log: sanitize_sensitive_info: false @@ -35,6 +35,7 @@ web_ui_enabled: true - Added new option to `input` `xtream_live_stream_without_extension`. Default is `false`. Some providers don't like `.ts` extension, some providers need it. Now you can disable or enable it for a provider. - Added `path` to `api-proxy.yml` server config for simpler front reverse-proxy configuration (like nginx) +- added `hls` handling. # 2.1.3 (2025-01-26) - Hotfix 2.1.2, forgot to update the stream api code. diff --git a/frontend/src/component/api-proxy-view/api-proxy-view.tsx b/frontend/src/component/api-proxy-view/api-proxy-view.tsx index 7624c67da..40a624d97 100644 --- a/frontend/src/component/api-proxy-view/api-proxy-view.tsx +++ b/frontend/src/component/api-proxy-view/api-proxy-view.tsx @@ -15,11 +15,10 @@ const SERVER_INFO_FIELDS = [ {name: 'protocol', label: 'Protocol', fieldType: FormFieldType.SINGLE_SELECT, options:[{value: 'http', label:'http'}, {value: 'https', label:'https'}]}, {name: 'host', label: 'Host', fieldType: FormFieldType.TEXT}, - {name: 'http_port', label: 'HTTP port', fieldType: FormFieldType.NUMBER, validator: isNumber}, - {name: 'https_port', label: 'HTTPS port', fieldType: FormFieldType.NUMBER,validator: isNumber}, - {name: 'rtmp_port', label: 'RTMP port', fieldType: FormFieldType.NUMBER,validator: isNumber}, + {name: 'port', label: 'Port', fieldType: FormFieldType.NUMBER, validator: isNumber}, {name: 'timezone', label: 'Timezone', fieldType: FormFieldType.TEXT}, {name: 'message', label: 'Message', fieldType: FormFieldType.TEXT}, + {name: 'path', label: 'Path', fieldType: FormFieldType.TEXT}, ]; interface ApiProxyViewProps { diff --git a/frontend/src/model/server-config.ts b/frontend/src/model/server-config.ts index bff367204..70cb5c3af 100644 --- a/frontend/src/model/server-config.ts +++ b/frontend/src/model/server-config.ts @@ -151,11 +151,10 @@ export interface ApiProxyServerInfo { name: string; protocol: string; host: string; - http_port: string; - https_port: string; - rtmp_port: string; + port: string; timezone: string; message: string; + path: string; } export interface ApiProxyConfig { diff --git a/src/api/api_utils.rs b/src/api/api_utils.rs index 731afe457..4475becd8 100644 --- a/src/api/api_utils.rs +++ b/src/api/api_utils.rs @@ -31,6 +31,44 @@ use std::sync::Arc; use url::Url; use crate::api::model::active_client_stream::ActiveClientStream; +#[macro_export] +macro_rules! try_option_bad_request { + ($option:expr, $msg_is_error:expr, $msg:expr) => { + match $option { + Some(value) => value, + None => { + if $msg_is_error {error!("{}", $msg);} else {debug!("{}", $msg);} + return HttpResponse::BadRequest().finish(); + } + } + }; + ($option:expr) => { + match $option { + Some(value) => value, + None => return HttpResponse::BadRequest().finish(), + } + }; +} + +#[macro_export] +macro_rules! try_result_bad_request { + ($option:expr, $msg_is_error:expr, $msg:expr) => { + match $option { + Ok(value) => value, + Err(_) => { + if $msg_is_error {error!("{}", $msg);} else {debug!("{}", $msg);} + return HttpResponse::BadRequest().finish(); + } + } + }; + ($option:expr) => { + match $option { + Ok(value) => value, + Err(_) => return HttpResponse::BadRequest().finish(), + } + }; +} + pub async fn serve_file(file_path: &Path, req: &HttpRequest, mime_type: mime::Mime) -> HttpResponse { if file_path.exists() { if let Ok(file) = actix_files::NamedFile::open_async(file_path).await { @@ -250,3 +288,11 @@ pub async fn resource_response(app_state: &AppState, resource_url: &str, req: &H } HttpResponse::BadRequest().finish() } + +pub fn separate_number_and_remainder(input: &str) -> (String, Option) { + input.rfind('.').map_or_else(|| (input.to_string(), None), |dot_index| { + let number_part = input[..dot_index].to_string(); + let rest = input[dot_index..].to_string(); + (number_part, if rest.len() < 2 { None } else { Some(rest) }) + }) +} diff --git a/src/api/hls_api.rs b/src/api/hls_api.rs new file mode 100644 index 000000000..f3c1bdb6d --- /dev/null +++ b/src/api/hls_api.rs @@ -0,0 +1,89 @@ +use std::sync::Arc; +use actix_web::{web, HttpRequest, HttpResponse}; +use actix_web::web::Data; +use log::{debug, error}; +use crate::api::api_utils::{get_user_target_by_credentials, stream_response}; +use crate::api::model::app_state::AppState; +use crate::api::model::request::UserApiRequest; +use crate::model::api_proxy::ProxyUserCredentials; +use crate::model::config::{ConfigInput, TargetType}; +use crate::model::playlist::{PlaylistEntry, PlaylistItemType, XtreamCluster}; +use crate::processing::hls_parser::{rewrite_hls, M3U_HLSR_PREFIX}; +use crate::{try_option_bad_request, try_result_bad_request}; +use crate::repository::{m3u_repository, xtream_repository}; +use crate::repository::playlist_repository::HLS_EXT; +use crate::utils::request_utils; +use crate::utils::request_utils::{replace_extension, sanitize_sensitive_info}; + +pub(in crate::api) async fn handle_hls_stream_request(app_state: &Data, user: &ProxyUserCredentials, pli: &dyn PlaylistEntry, input: &ConfigInput, target_type: TargetType) -> HttpResponse { + let url = replace_extension(&pli.get_provider_url(), HLS_EXT); + match request_utils::download_text_content(Arc::clone(&app_state.http_client), input, &url, None).await { + Ok(content) => { + let hls_content = rewrite_hls(&content, pli.get_virtual_id(), user, &target_type); + HttpResponse::Ok().content_type("application/x-mpegurl").body(hls_content) + } + Err(err) => { + error!("Failed to download m3u8 {}", sanitize_sensitive_info(err.to_string().as_str())); + HttpResponse::NoContent().finish() + } + } +} + +async fn hls_api_stream( + req: &HttpRequest, + api_req: &web::Query, + path: web::Path<(String, String, String, String, String, String)>, + app_state: &web::Data, + target_type: TargetType +) -> HttpResponse { + let (_token, username, password, channel, _hash, _chunk) = path.into_inner(); + let (_user, target) = try_option_bad_request!( + get_user_target_by_credentials(&username, &password, api_req, app_state), + false, + format!("Could not find any user {username}")); + + let target_name = &target.name; + let virtual_id: u32 = try_result_bad_request!(channel.parse()); + let (pli_url, input_name) = if target_type == TargetType::Xtream { + let (pli, _ ) = try_result_bad_request!(xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, None).await, true, format!("Failed to read xtream item for stream id {}", virtual_id)); + (pli.url, pli.input_name) + } else { + let pli = try_result_bad_request!(m3u_repository::m3u_get_item_for_stream_id(virtual_id, &app_state.config, target).await, true, format!("Failed to read xtream item for stream id {}", virtual_id)); + (pli.url, pli.input_name) + }; + let input = try_option_bad_request!(app_state.config.get_input_by_name(&input_name), true, format!("Cant find input for target {target_name}, context {}, stream_id {virtual_id}", XtreamCluster::Live)); + // let input_username = input.username.as_ref().map_or("", |v| v); + // let input_password = input.password.as_ref().map_or("", |v| v); + // let input_url = input.url.as_str(); + + // we don't respond as hlsr, we take the original stream, because the location could be different and then it does not work + // The next problem is, different url to same channel causes to fail stream share. + // let stream_url = format!("{input_url}/hlsr/{token}/{input_username}/{input_password}/{}/{hash}/{chunk}", pli.provider_id); + stream_response(app_state, &pli_url, req, Some(input), PlaylistItemType::Live, target).await +} + +async fn hls_api_stream_xtream( + req: HttpRequest, + api_req: web::Query, + path: web::Path<(String, String, String, String, String, String)>, + app_state: web::Data, +) -> HttpResponse { + hls_api_stream(&req, &api_req, path, &app_state, TargetType::Xtream).await +} + +async fn hls_api_stream_m3u( + req: HttpRequest, + api_req: web::Query, + path: web::Path<(String, String, String, String, String, String)>, + app_state: web::Data, +) -> HttpResponse { + hls_api_stream(&req, &api_req, path, &app_state, TargetType::M3u).await +} + + +pub fn hls_api_register(cfg: &mut web::ServiceConfig) { + cfg.service(web::resource("/hlsr/{token}/{username}/{password}/{channel}/{hash}/{chunk}").route(web::get().to(hls_api_stream_xtream))); + cfg.service(web::resource(format!("/{M3U_HLSR_PREFIX}/{{token}}/{{username}}/{{password}}/{{channel}}/{{hash}}/{{chunk}}")).route(web::get().to(hls_api_stream_m3u))); + //cfg.service(web::resource("/hls/{token}/{stream}").route(web::get().to(xtream_player_api_hls_stream))); + //cfg.service(web::resource("/play/{token}/{type}").route(web::get().to(xtream_player_api_play_stream))); +} \ No newline at end of file diff --git a/src/api/m3u_api.rs b/src/api/m3u_api.rs index f16be0144..0719761ed 100644 --- a/src/api/m3u_api.rs +++ b/src/api/m3u_api.rs @@ -3,16 +3,18 @@ use bytes::Bytes; use futures::stream; use log::{debug, error}; -use crate::api::api_utils::{get_user_target, get_user_target_by_credentials, resource_response, stream_response}; +use crate::api::api_utils::{get_user_target, get_user_target_by_credentials, resource_response, separate_number_and_remainder, stream_response}; +use crate::api::hls_api::handle_hls_stream_request; use crate::api::model::app_state::AppState; use crate::api::model::request::UserApiRequest; use crate::model::api_proxy::ProxyType; use crate::model::config::TargetType; -use crate::model::playlist::FieldGetAccessor; -use crate::repository::m3u_playlist_iterator::{M3U_STREAM_PATH, M3U_RESOURCE_PATH}; -use crate::repository::m3u_repository::{m3u_get_file_paths, m3u_get_item_for_stream_id, m3u_load_rewrite_playlist}; -use crate::repository::storage::get_target_storage_path; -use crate::utils::request_utils::sanitize_sensitive_info; +use crate::model::playlist::{FieldGetAccessor, XtreamCluster}; +use crate::repository::m3u_playlist_iterator::{M3U_RESOURCE_PATH, M3U_STREAM_PATH}; +use crate::repository::m3u_repository::{m3u_get_item_for_stream_id, m3u_load_rewrite_playlist}; +use crate::utils::request_utils::{replace_extension, sanitize_sensitive_info}; +use crate::{debug_if_enabled, try_option_bad_request, try_result_bad_request}; +use crate::repository::playlist_repository::HLS_EXT; async fn m3u_api( api_req: &UserApiRequest, @@ -57,20 +59,15 @@ async fn m3u_api_stream( app_state: web::Data, ) -> HttpResponse { let (username, password, stream_id) = path.into_inner(); - let Ok(m3u_stream_id) = stream_id.parse::() else { return HttpResponse::BadRequest().finish() }; + let (action_stream_id, stream_ext) = separate_number_and_remainder(&stream_id); + let virtual_id: u32 = try_result_bad_request!(action_stream_id.trim().parse()); let Some((user, target)) = get_user_target_by_credentials(&username, &password, &api_req, &app_state) else { return HttpResponse::BadRequest().finish() }; if !target.has_output(&TargetType::M3u) { return HttpResponse::BadRequest().finish(); } - let Some(target_path) = get_target_storage_path(&app_state.config, target.name.as_str()) else { - error!("Failed to get target path for {}", target.name); - return HttpResponse::BadRequest().finish(); - }; - - let (m3u_path, idx_path) = m3u_get_file_paths(&target_path); - let m3u_item = match m3u_get_item_for_stream_id(&app_state.config, m3u_stream_id, &m3u_path, &idx_path).await { + let m3u_item = match m3u_get_item_for_stream_id(virtual_id, &app_state.config, target).await { Ok(item) => item, Err(err) => { error!("Failed to get m3u url: {}", sanitize_sensitive_info(err.to_string().as_str())); @@ -78,10 +75,19 @@ async fn m3u_api_stream( } }; + let is_hls_request = stream_ext.as_deref() == Some(HLS_EXT); + if user.proxy == ProxyType::Redirect { - let stream_url = m3u_item.url; - debug!("Redirecting stream request to {}", sanitize_sensitive_info(&stream_url)); - return HttpResponse::Found().insert_header(("Location", stream_url.to_string())).finish(); + let redirect_url = if is_hls_request { &replace_extension(&m3u_item.url, "m3u8") } else { &m3u_item.url }; + debug_if_enabled!("Redirecting m3u stream request to {}", sanitize_sensitive_info(redirect_url)); + return HttpResponse::Found().insert_header(("Location", redirect_url.as_str())).finish(); + } + // Reverse proxy mode + if is_hls_request { + let target_name = &target.name; + let input = try_option_bad_request!(app_state.config.get_input_by_name(m3u_item.input_name.as_str()), true, + format!("Cant find input for target {target_name}, context {}, stream_id {virtual_id}", XtreamCluster::Live)); + return handle_hls_stream_request(&app_state, &user, &m3u_item, input, TargetType::M3u).await; } stream_response(&app_state, m3u_item.url.as_str(), &req, None, m3u_item.item_type, target).await @@ -100,14 +106,7 @@ async fn m3u_api_resource( if !target.has_output(&TargetType::M3u) { return HttpResponse::BadRequest().finish(); } - - let Some(target_path) = get_target_storage_path(&app_state.config, target.name.as_str()) else { - error!("Failed to get target path for {}", target.name); - return HttpResponse::BadRequest().finish(); - }; - - let (m3u_path, idx_path) = m3u_get_file_paths(&target_path); - let m3u_item = match m3u_get_item_for_stream_id(&app_state.config, m3u_stream_id, &m3u_path, &idx_path).await { + let m3u_item = match m3u_get_item_for_stream_id(m3u_stream_id, &app_state.config, target).await { Ok(item) => item, Err(err) => { error!("Failed to get m3u url: {}", sanitize_sensitive_info(err.to_string().as_str())); diff --git a/src/api/main_api.rs b/src/api/main_api.rs index 9b6b682b6..38e470df9 100644 --- a/src/api/main_api.rs +++ b/src/api/main_api.rs @@ -9,6 +9,7 @@ use std::io::ErrorKind; use std::path::{PathBuf}; use std::sync::Arc; use std::sync::atomic::{AtomicUsize, Ordering}; +use crate::api::hls_api::hls_api_register; use crate::api::m3u_api::m3u_api_register; use crate::api::model::app_state::AppState; use crate::api::model::download::DownloadQueue; @@ -177,6 +178,7 @@ pub async fn start_server(cfg: Arc, targets: Arc) -> fut .configure(xtream_api_register) .configure(m3u_api_register) .configure(xmltv_api_register) + .configure(hls_api_register) .configure(|srvcfg| { if web_ui_enabled { srvcfg.configure(index_register(&web_dir_path)); diff --git a/src/api/mod.rs b/src/api/mod.rs index 737f4af48..ca4cfc3fe 100644 --- a/src/api/mod.rs +++ b/src/api/mod.rs @@ -8,4 +8,5 @@ mod xmltv_api; mod scheduler; mod web_index; -pub(crate) mod model; \ No newline at end of file +pub(crate) mod model; +mod hls_api; \ No newline at end of file diff --git a/src/api/model/client_stream.rs b/src/api/model/client_stream.rs index b4ce7b807..6b8f29c68 100644 --- a/src/api/model/client_stream.rs +++ b/src/api/model/client_stream.rs @@ -4,9 +4,9 @@ use std::pin::Pin; use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::{Arc}; use std::task::{Poll}; -use log::debug; use futures::{Stream}; use crate::api::model::stream_error::StreamError; +use crate::trace_if_enabled; use crate::utils::atomic_once_flag::AtomicOnceFlag; use crate::utils::request_utils::sanitize_sensitive_info; @@ -57,7 +57,7 @@ impl Stream for ClientStream { impl Drop for ClientStream { fn drop(&mut self) { - debug!("Client disconnected {}", sanitize_sensitive_info(&self.url)); + trace_if_enabled!("Client disconnected {}", sanitize_sensitive_info(&self.url)); self.close_signal.notify(); } } \ No newline at end of file diff --git a/src/api/model/xtream.rs b/src/api/model/xtream.rs index 55b767205..d204b8f61 100644 --- a/src/api/model/xtream.rs +++ b/src/api/model/xtream.rs @@ -26,10 +26,8 @@ pub struct XtreamUserInfo { pub struct XtreamServerInfo { pub url: String, pub port: String, - pub https_port: String, - pub server_protocol: String, - // http, https - pub rtmp_port: String, + pub path: Option, + pub protocol: String, // http, https pub timezone: String, pub timestamp_now: i64, pub time_now: String, //"2021-06-28 17:07:37" @@ -60,10 +58,9 @@ impl XtreamAuthorizationResponse { }, server_info: XtreamServerInfo { url: server_info.host.clone(), - port: if server_info.protocol == "http" { server_info.port.clone() } else { String::new() }, - https_port: if server_info.protocol == "https" { server_info.port.clone() } else { String::new() }, - server_protocol: server_info.protocol.clone(), - rtmp_port: String::new(), + port: server_info.port.clone(), + protocol: server_info.protocol.clone(), + path: server_info.path.clone(), timezone: server_info.timezone.to_string(), timestamp_now: now.timestamp(), time_now: now.format("%Y-%m-%d %H:%M:%S").to_string(), diff --git a/src/api/xmltv_api.rs b/src/api/xmltv_api.rs index 3f3c47fe0..cad9bf0e4 100644 --- a/src/api/xmltv_api.rs +++ b/src/api/xmltv_api.rs @@ -2,7 +2,7 @@ use std::fs::File; use std::path::{Path, PathBuf}; use actix_web::{HttpRequest, HttpResponse, web, http::header}; -use log::{error, info}; +use log::{error, trace}; use quick_xml::{Reader, Writer}; use flate2::write::GzEncoder; use flate2::Compression; @@ -42,7 +42,7 @@ fn get_epg_path_for_target_of_type(target_name: &str, epg_path: PathBuf) -> Opti if file_utils::path_exists(&epg_path) { return Some(epg_path); } - info!("Cant find epg file for {target_name} target: {}", epg_path.to_str().unwrap_or("?")); + trace!("Cant find epg file for {target_name} target: {}", epg_path.to_str().unwrap_or("?")); None } diff --git a/src/api/xtream_api.rs b/src/api/xtream_api.rs index 9d1825c64..a86ef9b72 100644 --- a/src/api/xtream_api.rs +++ b/src/api/xtream_api.rs @@ -1,6 +1,6 @@ // https://github.com/tellytv/go.xtream-codes/blob/master/structs.go -use crate::Arc; +use crate::{trace_if_enabled, try_option_bad_request, try_result_bad_request, Arc}; use std::collections::HashMap; use std::fmt::{Display, Formatter}; use std::path::Path; @@ -14,7 +14,8 @@ use futures::Stream; use log::{debug, error, warn}; use serde_json::{Map, Value}; -use crate::api::api_utils::{get_user_target, get_user_target_by_credentials, resource_response, serve_file, stream_response}; +use crate::api::api_utils::{get_user_target, get_user_target_by_credentials, resource_response, separate_number_and_remainder, serve_file, stream_response}; +use crate::api::hls_api::handle_hls_stream_request; use crate::api::model::app_state::AppState; use crate::api::model::request::UserApiRequest; use crate::api::model::xtream::XtreamAuthorizationResponse; @@ -24,18 +25,15 @@ use crate::model::config::TargetType; use crate::model::config::{Config, ConfigInput, ConfigTarget}; use crate::model::playlist::{get_backdrop_path_value, FieldGetAccessor, PlaylistEntry, PlaylistItemType, XtreamCluster, XtreamPlaylistItem}; use crate::model::xtream::{INFO_RESOURCE_PREFIX, INFO_RESOURCE_PREFIX_EPISODE, PROP_BACKDROP_PATH, SEASON_RESOURCE_PREFIX}; +use crate::repository::playlist_repository::HLS_EXT; use crate::repository::storage::{get_target_storage_path, hex_encode}; use crate::repository::target_id_mapping::TargetIdMapping; use crate::repository::xtream_repository; use crate::repository::xtream_repository::{TAG_EPISODES, TAG_INFO_DATA, TAG_SEASONS_DATA}; -use crate::utils::hash_utils::generate_playlist_uuid; -use crate::utils::json_utils::get_u32_from_serde_value; -use crate::utils::request_utils::{extract_extension_from_url, sanitize_sensitive_info}; -use crate::utils::xtream_utils::{ACTION_GET_LIVE_CATEGORIES, ACTION_GET_LIVE_STREAMS, - ACTION_GET_SERIES, ACTION_GET_SERIES_CATEGORIES, - ACTION_GET_SERIES_INFO, ACTION_GET_VOD_CATEGORIES, - ACTION_GET_VOD_INFO, - ACTION_GET_VOD_STREAMS}; +use crate::utils::hash_utils::{generate_playlist_uuid}; +use crate::utils::json_utils::{get_u32_from_serde_value}; +use crate::utils::request_utils::{extract_extension_from_url, replace_extension, sanitize_sensitive_info}; +use crate::utils::xtream_utils::{create_vod_info_from_item, ACTION_GET_LIVE_CATEGORIES, ACTION_GET_LIVE_STREAMS, ACTION_GET_SERIES, ACTION_GET_SERIES_CATEGORIES, ACTION_GET_SERIES_INFO, ACTION_GET_VOD_CATEGORIES, ACTION_GET_VOD_INFO, ACTION_GET_VOD_STREAMS}; use crate::utils::{json_utils, request_utils, xtream_utils}; use crate::{debug_if_enabled, info_err}; @@ -47,41 +45,6 @@ const TAG_CATEGORY_ID: &str = "category_id"; const TAG_STREAM_ID: &str = "stream_id"; const TAG_EPG_LISTINGS: &str = "epg_listings"; -macro_rules! try_option_bad_request { - ($option:expr, $msg_is_error:expr, $msg:expr) => { - match $option { - Some(value) => value, - None => { - if $msg_is_error {error!("{}", $msg);} else {debug!("{}", $msg);} - return HttpResponse::BadRequest().finish(); - } - } - }; - ($option:expr) => { - match $option { - Some(value) => value, - None => return HttpResponse::BadRequest().finish(), - } - }; -} -macro_rules! try_result_bad_request { - ($option:expr, $msg_is_error:expr, $msg:expr) => { - match $option { - Ok(value) => value, - Err(_) => { - if $msg_is_error {error!("{}", $msg);} else {debug!("{}", $msg);} - return HttpResponse::BadRequest().finish(); - } - } - }; - ($option:expr) => { - match $option { - Ok(value) => value, - Err(_) => return HttpResponse::BadRequest().finish(), - } - }; -} - #[derive(Debug)] enum XtreamApiStreamContext { LiveAlt, @@ -161,14 +124,6 @@ fn get_user_info(user: &ProxyUserCredentials, cfg: &Config) -> XtreamAuthorizati XtreamAuthorizationResponse::new(&server_info, user) } -fn xtream_api_request_separate_number_and_remainder(input: &str) -> (String, Option) { - input.rfind('.').map_or_else(|| (input.to_string(), None), |dot_index| { - let number_part = input[..dot_index].to_string(); - let rest = input[dot_index..].to_string(); - (number_part, if rest.len() < 2 { None } else { Some(rest) }) - }) -} - async fn xtream_player_api_stream( req: &HttpRequest, api_req: &web::Query, @@ -181,7 +136,7 @@ async fn xtream_player_api_stream( debug!("Target has no xtream output {}", target_name); return HttpResponse::BadRequest().finish(); } - let (action_stream_id, stream_ext) = xtream_api_request_separate_number_and_remainder(stream_req.stream_id); + let (action_stream_id, stream_ext) = separate_number_and_remainder(stream_req.stream_id); let virtual_id: u32 = try_result_bad_request!(action_stream_id.trim().parse()); let (pli, mapping) = try_result_bad_request!(xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, None).await, true, format!("Failed to read xtream item for stream id {}", virtual_id)); let input = try_option_bad_request!(app_state.config.get_input_by_name(pli.input_name.as_str()), true, format!("Cant find input for target {target_name}, context {}, stream_id {virtual_id}", stream_req.context)); @@ -191,6 +146,8 @@ async fn xtream_player_api_stream( return HttpResponse::Found().insert_header(("Location", pli.url.to_string())).finish(); } + let is_hls_request = stream_ext.as_deref() == Some(HLS_EXT); + if user.proxy == ProxyType::Redirect { if pli.xtream_cluster == XtreamCluster::Series { let ext = stream_ext.unwrap_or_else(String::new); @@ -200,8 +157,15 @@ async fn xtream_player_api_stream( let stream_url = format!("{url}/series/{username}/{password}/{}{ext}", mapping.provider_id); return HttpResponse::Found().insert_header(("Location", stream_url)).finish(); } - debug_if_enabled!("Redirecting stream request to {}", sanitize_sensitive_info(&pli.url)); - HttpResponse::Found().insert_header(("Location", pli.url.as_str())).finish(); + + let redirect_url = if is_hls_request { &replace_extension(&pli.url, "m3u8") } else { &pli.url }; + debug_if_enabled!("Redirecting stream request to {}", sanitize_sensitive_info(redirect_url)); + return HttpResponse::Found().insert_header(("Location", redirect_url.as_str())).finish(); + } + + // Reverse proxy mode + if is_hls_request { + return handle_hls_stream_request(app_state, &user, &pli, input, TargetType::Xtream).await; } let extension = stream_ext.unwrap_or_else( @@ -218,10 +182,11 @@ async fn xtream_player_api_stream( true, format!("Cant find stream url for target {target_name}, context {}, stream_id {virtual_id}", stream_req.context)); - debug_if_enabled!("Streaming stream request from {}", sanitize_sensitive_info(&stream_url)); + trace_if_enabled!("Streaming stream request from {}", sanitize_sensitive_info(&stream_url)); stream_response(app_state, &stream_url, req, Some(input), pli.item_type, target).await } + fn get_doc_id_and_field_name(input: &str) -> Option<(u32, &str)> { if let Some(pos) = input.find('_') { let (number_part, rest) = input.split_at(pos); @@ -374,10 +339,10 @@ async fn xtream_player_api_resource( None => HttpResponse::NotFound().finish(), Some(url) => { if user.proxy == ProxyType::Redirect { - debug!("Redirecting resource request to {}", sanitize_sensitive_info(&url)); + trace_if_enabled!("Redirecting resource request to {}", sanitize_sensitive_info(&url)); HttpResponse::Found().insert_header(("Location", url.as_str())).finish() } else { - debug_if_enabled!("Resource request to {}", sanitize_sensitive_info(&url)); + trace_if_enabled!("Resource request to {}", sanitize_sensitive_info(&url)); resource_response(app_state, url.as_str(), req, None).await } } @@ -456,24 +421,30 @@ async fn xtream_get_stream_info_response(app_state: &AppState, user: &ProxyUserC Err(_) => return HttpResponse::BadRequest().finish() }; - if let Ok((pli, _)) = xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, Some(cluster)).await { - let input_name = Rc::clone(&pli.input_name); - if let Some(input) = app_state.config.get_input_by_name(input_name.as_str()) { - if let Some(info_url) = xtream_utils::get_xtream_player_api_info_url(input, cluster, pli.provider_id) { - // Redirect is only possible for live streams, vod and series info needs to be modified - if user.proxy == ProxyType::Redirect && cluster == XtreamCluster::Live { - return HttpResponse::Found().insert_header(("Location", info_url)).finish(); - } else if let Ok(content) = xtream_utils::get_xtream_stream_info(Arc::clone(&app_state.http_client), &app_state.config, user, input, target, &pli, info_url.as_str(), cluster).await { - return HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body(content); + if let Ok((pli, virtual_record)) = xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, Some(cluster)).await { + if pli.provider_id > 0 { + let input_name = Rc::clone(&pli.input_name); + if let Some(input) = app_state.config.get_input_by_name(input_name.as_str()) { + if let Some(info_url) = xtream_utils::get_xtream_player_api_info_url(input, cluster, pli.provider_id) { + // Redirect is only possible for live streams, vod and series info needs to be modified + if user.proxy == ProxyType::Redirect && cluster == XtreamCluster::Live { + return HttpResponse::Found().insert_header(("Location", info_url)).finish(); + } else if let Ok(content) = xtream_utils::get_xtream_stream_info(Arc::clone(&app_state.http_client), &app_state.config, user, input, target, &pli, info_url.as_str(), cluster).await { + return HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body(content); + } } } } + + return match cluster { + XtreamCluster::Video => { + let content = create_vod_info_from_item(user, &pli, virtual_record.last_updated); + HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body(content) + } + XtreamCluster::Live | XtreamCluster::Series => HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body("{}"), + }; } - match cluster { - XtreamCluster::Live => HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body("{}"), - XtreamCluster::Video => HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body("{info:[]}"), - XtreamCluster::Series => HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body("[]"), - } + HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body("{}") } async fn xtream_get_short_epg(app_state: &AppState, user: &ProxyUserCredentials, target: &ConfigTarget, stream_id: &str, limit: &str) -> HttpResponse { @@ -485,24 +456,26 @@ async fn xtream_get_short_epg(app_state: &AppState, user: &ProxyUserCredentials, }; if let Ok((pli, _)) = xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, None).await { - let input_name = Rc::clone(&pli.input_name); - if let Some(input) = app_state.config.get_input_by_name(input_name.as_str()) { - if let Some(action_url) = xtream_utils::get_xtream_player_api_action_url(input, ACTION_GET_SHORT_EPG) { - let mut info_url = format!("{action_url}&{TAG_STREAM_ID}={}", pli.provider_id); - if !(limit.is_empty() || limit.eq("0")) { - info_url = format!("{info_url}&limit={limit}"); - } - if user.proxy == ProxyType::Redirect { - return HttpResponse::Found().insert_header(("Location", info_url)).finish(); - } - - return match request_utils::download_text_content(Arc::clone(&app_state.http_client), input, info_url.as_str(), None).await { - Ok(content) => HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body(content), - Err(err) => { - error!("Failed to download epg {}", sanitize_sensitive_info(err.to_string().as_str())); - HttpResponse::NoContent().finish() + if pli.provider_id > 0 { + let input_name = Rc::clone(&pli.input_name); + if let Some(input) = app_state.config.get_input_by_name(input_name.as_str()) { + if let Some(action_url) = xtream_utils::get_xtream_player_api_action_url(input, ACTION_GET_SHORT_EPG) { + let mut info_url = format!("{action_url}&{TAG_STREAM_ID}={}", pli.provider_id); + if !(limit.is_empty() || limit.eq("0")) { + info_url = format!("{info_url}&limit={limit}"); } - }; + if user.proxy == ProxyType::Redirect { + return HttpResponse::Found().insert_header(("Location", info_url)).finish(); + } + + return match request_utils::download_text_content(Arc::clone(&app_state.http_client), input, info_url.as_str(), None).await { + Ok(content) => HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body(content), + Err(err) => { + error!("Failed to download epg {}", sanitize_sensitive_info(err.to_string().as_str())); + HttpResponse::NoContent().finish() + } + }; + } } } } @@ -541,28 +514,44 @@ async fn xtream_get_catchup_response(app_state: &AppState, target: &ConfigTarget let mut doc: Map = try_result_bad_request!(serde_json::from_str(&content)); let epg_listings = try_option_bad_request!(doc.get_mut(TAG_EPG_LISTINGS).and_then(Value::as_array_mut)); let target_path = try_option_bad_request!(get_target_storage_path(&app_state.config, target.name.as_str())); - let mut target_id_mapping = TargetIdMapping::new(&target_path); - + let mut target_id_mapping = { + let _file_lock = match app_state.config.file_locks.read_lock(&target_path).await { + Ok(lock) => lock, + Err(err) => { + error!("Could not get lock for id mapping for target {} err:{err}", target.name); + return HttpResponse::InternalServerError().finish(); + } + }; + TargetIdMapping::new(&target_path) + }; for epg_list_item in epg_listings.iter_mut().filter_map(Value::as_object_mut) { // TODO epg_id if let Some(catchup_provider_id) = epg_list_item.get(TAG_ID).and_then(Value::as_str).and_then(|id| id.parse::().ok()) { let uuid = generate_playlist_uuid(&hex_encode(&pli.get_uuid()), &catchup_provider_id.to_string(), &pli.url); - let virtual_id = target_id_mapping.get_virtual_id(uuid, catchup_provider_id, PlaylistItemType::Catchup, pli.provider_id); + let virtual_id = target_id_mapping.get_and_update_virtual_id(uuid, catchup_provider_id, PlaylistItemType::Catchup, pli.provider_id); epg_list_item.insert(TAG_ID.to_string(), Value::String(virtual_id.to_string())); } } - if let Err(err) = target_id_mapping.persist() { - error!("Failed to write catchup id mapping {err}"); - return HttpResponse::BadRequest().finish(); - } - + { + let _file_lock = match app_state.config.file_locks.write_lock(&target_path).await { + Ok(lock) => lock, + Err(err) => { + error!("Could not get lock for id mapping for target {} err:{err}", target.name); + return HttpResponse::InternalServerError().finish(); + } + }; + if let Err(err) = target_id_mapping.persist() { + error!("Failed to write catchup id mapping {err}"); + return HttpResponse::BadRequest().finish(); + } + }; serde_json::to_string(&doc).map_or_else(|_| HttpResponse::BadRequest().finish(), |result| HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body(result)) } -macro_rules! skip_response_if_flag_set { +macro_rules! skip_json_response_if_flag_set { ($flag:expr, $stmt:expr) => { if $flag { - return HttpResponse::NoContent().finish(); + return HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body("[]"); } return $stmt; }; @@ -606,10 +595,10 @@ async fn xtream_player_api( match action { ACTION_GET_SERIES_INFO => { - skip_response_if_flag_set!(skip_series, xtream_get_stream_info_response(app_state, &user, target, api_req.series_id.trim(), XtreamCluster::Series).await); + skip_json_response_if_flag_set!(skip_series, xtream_get_stream_info_response(app_state, &user, target, api_req.series_id.trim(), XtreamCluster::Series).await); } ACTION_GET_VOD_INFO => { - skip_response_if_flag_set!(skip_vod, xtream_get_stream_info_response(app_state, &user, target, api_req.vod_id.trim(), XtreamCluster::Video).await); + skip_json_response_if_flag_set!(skip_vod, xtream_get_stream_info_response(app_state, &user, target, api_req.vod_id.trim(), XtreamCluster::Video).await); } ACTION_GET_EPG | ACTION_GET_SHORT_EPG => { return xtream_get_short_epg( @@ -617,7 +606,7 @@ async fn xtream_player_api( ).await; } ACTION_GET_CATCHUP_TABLE => { - skip_response_if_flag_set!(skip_live, xtream_get_catchup_response(app_state, target, api_req.stream_id.trim(), api_req.start.trim(), api_req.end.trim()).await); + skip_json_response_if_flag_set!(skip_live, xtream_get_catchup_response(app_state, target, api_req.stream_id.trim(), api_req.start.trim(), api_req.end.trim()).await); } _ => {} } @@ -653,12 +642,15 @@ async fn xtream_player_api( } Err(err) => { error!("Failed response for xtream target: {} action: {} error: {}", &target.name, action, err); - HttpResponse::NoContent().finish() + // Some players fail on NoContent, so we return an empty array + HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body("[]") + // HttpResponse::NoContent().finish() } } } None => { - HttpResponse::NoContent().finish() + // Some players fail on NoContent, so we return an empty array + HttpResponse::Ok().content_type(mime::APPLICATION_JSON).body("[]") } } } else { @@ -746,9 +738,4 @@ pub fn xtream_api_register(cfg: &mut web::ServiceConfig) { ("live", xtream_player_api_live_resource), ("movie", xtream_player_api_movie_resource), ("series", xtream_player_api_series_resource)]); - /* 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))); - cfg.service(web::resource("/play/{token}/{type}").route(web::get().to(xtream_player_api_play_stream))); - */ } \ No newline at end of file diff --git a/src/m3u_filter_error.rs b/src/m3u_filter_error.rs index 1bd27e3f3..7f710df5b 100644 --- a/src/m3u_filter_error.rs +++ b/src/m3u_filter_error.rs @@ -12,7 +12,7 @@ macro_rules! get_errors_notify_message { .filter(|&err| err.kind == M3uFilterErrorKind::Notify) .map(|err| err.message.as_str()) .collect::>() - .join("\n"); + .join("\r\n"); if $size > 0 && text.len() > std::cmp::max($size - 3, 3) { Some(format!("{}...", text.get(0..$size).unwrap())) } else { diff --git a/src/model/playlist.rs b/src/model/playlist.rs index 4556db699..c085b475f 100644 --- a/src/model/playlist.rs +++ b/src/model/playlist.rs @@ -11,6 +11,7 @@ use crate::utils::json_utils::{get_string_from_serde_value, get_u64_from_serde_v use serde::{Deserialize, Serialize}; use serde_json::{Map, Value}; use crate::utils::hash_utils::{generate_playlist_uuid, get_provider_id}; +use crate::utils::request_utils::extract_extension_from_url; // https://de.wikipedia.org/wiki/M3U // https://siptv.eu/howto/playlist.html @@ -76,7 +77,7 @@ impl TryFrom for XtreamCluster { type Error = String; fn try_from(item_type: PlaylistItemType) -> Result { match item_type { - PlaylistItemType::Live => Ok(Self::Live), + PlaylistItemType::Live | PlaylistItemType::LiveHls | PlaylistItemType::LiveUnknown => Ok(Self::Live), PlaylistItemType::Video => Ok(Self::Video), PlaylistItemType::Series => Ok(Self::Series), _ => Err(format!("Cant convert {item_type}")), @@ -384,6 +385,15 @@ impl XtreamPlaylistItem { pub fn to_doc(&self, url: &str, options: &XtreamMappingOptions, user: &ProxyUserCredentials) -> Value { xtream_playlistitem_to_document(self, url, options, user) } + + pub fn get_additional_property(&self, field: &str) -> Option { + if let Some(json) = self.additional_properties.as_ref() { + if let Ok(Value::Object(props)) = serde_json::from_str(json) { + return props.get(field).cloned(); + } + } + None + } } impl PlaylistEntry for XtreamPlaylistItem { @@ -501,6 +511,39 @@ impl PlaylistItem { pub fn to_xtream(&self) -> XtreamPlaylistItem { let header = self.header.borrow(); let provider_id = header.id.parse::().unwrap_or_default(); + let mut additional_properties = None; + if header.xtream_cluster != XtreamCluster::Live { + let add_ext = match header.get_additional_property("container_extension") { + None => true, + Some(ext) => ext.as_str().map_or(true, str::is_empty) + }; + if add_ext { + if let Some(cont_ext) = extract_extension_from_url(&header.url) { + let ext = if let Some(stripped) = cont_ext.strip_prefix('.') { stripped } else { cont_ext }; + let mut result = match header.additional_properties.as_ref() { + None => Map::new(), + Some(props) => { + if let Value::Object(map) = props { + map.clone() + } else { + Map::new() + } + } + }; + result.insert("container_extension".to_string(), Value::String(ext.to_string())); + additional_properties = serde_json::to_string(&Value::Object(result)).ok(); + } + } + } + if additional_properties.is_none() { + additional_properties = header.additional_properties.as_ref().and_then(|props| { + serde_json::to_string(props).ok() + }); + } + // let additional_properties = header.additional_properties.as_ref().and_then(|props| { + // serde_json::to_string(props).ok() + // }); + XtreamPlaylistItem { virtual_id: header.virtual_id, provider_id, @@ -514,7 +557,7 @@ impl PlaylistItem { url: Rc::clone(&header.url), epg_channel_id: header.epg_channel_id.clone(), xtream_cluster: header.xtream_cluster, - additional_properties: header.additional_properties.as_ref().and_then(|props| serde_json::to_string(props).ok()), + additional_properties, item_type: header.item_type, category_id: header.category_id, input_name: Rc::clone(&header.input_name), diff --git a/src/processing/hls_parser.rs b/src/processing/hls_parser.rs new file mode 100644 index 000000000..c17c88a42 --- /dev/null +++ b/src/processing/hls_parser.rs @@ -0,0 +1,51 @@ +use crate::model::api_proxy::ProxyUserCredentials; +use std::str; +use crate::model::config::TargetType; + +// /hlsr/{token}/{username}/{password}/{channel}/{hash}/{chunk} +#[derive(Debug)] +pub struct HlsrPath { + token: String, + // username: String, + // password: String, + // channel: String, + hash: String, + chunk: String, +} +fn parse_hlsr_path(input: &str) -> Option { + let parts: Vec<&str> = input.split('/').collect(); + + if parts.len() != 8 || !parts[0].is_empty() || parts[1] != "hlsr" { + return None; + } + + Some(HlsrPath { + token: parts[2].to_string(), + // username: parts[3].to_string(), + // password: parts[4].to_string(), + // channel: parts[5].to_string(), + hash: parts[6].to_string(), + chunk: parts[7].to_string(), + }) +} + +pub const M3U_HLSR_PREFIX: &str = "mhlsr"; + +pub fn rewrite_hls_url(stream_id: u32, username: &str, password: &str, hlsr: &HlsrPath, target_type: &TargetType) -> String { + let prefix = if *target_type == TargetType::Xtream { "hlsr" } else { M3U_HLSR_PREFIX }; + format!("/{prefix}/{}/{username}/{password}/{stream_id}/{}/{}", hlsr.token, hlsr.hash, hlsr.chunk) +} + +pub fn rewrite_hls(content: &str, virtual_id: u32, user: &ProxyUserCredentials, target_type: &TargetType) -> String { + content.lines().map(|line| { + if line.starts_with('#') { + line.to_string() + } else { + match parse_hlsr_path(line) { + None => line.to_string(), + Some(hlsr) => rewrite_hls_url(virtual_id, &user.username, &user.password, &hlsr, target_type) + } + } + }).collect::>() + .join("\r\n") +} diff --git a/src/processing/m3u_parser.rs b/src/processing/m3u_parser.rs index 01f159714..4307f17cd 100644 --- a/src/processing/m3u_parser.rs +++ b/src/processing/m3u_parser.rs @@ -201,10 +201,13 @@ where { let mut sort_order: Vec> = vec![]; let mut sort_order_idx: usize = 0; - let mut group_map: std::collections::HashMap, usize> = std::collections::HashMap::new(); + let mut group_map: std::collections::HashMap = std::collections::HashMap::new(); consume_m3u(cfg, input, lines, |item| { // keep the original sort order for groups and group the playlist items - let key = Rc::clone(&item.header.borrow().group); + let key = { + let header = item.header.borrow(); + format!("{}{}", &header.xtream_cluster, &header.group) + }; match group_map.entry(key) { std::collections::hash_map::Entry::Vacant(v) => { v.insert(sort_order_idx); diff --git a/src/processing/mod.rs b/src/processing/mod.rs index 7ac87204b..0b67a5980 100644 --- a/src/processing/mod.rs +++ b/src/processing/mod.rs @@ -7,3 +7,4 @@ mod xtream_processor; mod affix_processor; mod xtream_processor_vod; mod xtream_processor_series; +pub mod hls_parser; diff --git a/src/processing/xtream_parser.rs b/src/processing/xtream_parser.rs index 4ff1eee31..50a3a4983 100644 --- a/src/processing/xtream_parser.rs +++ b/src/processing/xtream_parser.rs @@ -79,7 +79,7 @@ pub fn parse_xtream_series_info(info: &Value, group_title: &str, series_name: &s } } -fn create_xtream_url(xtream_cluster: XtreamCluster, url: &str, username: &str, password: &str, +pub fn create_xtream_url(xtream_cluster: XtreamCluster, url: &str, username: &str, password: &str, stream: &XtreamStream, live_stream_without_extension: bool) -> Rc { if stream.direct_source.is_empty() { let stream_base_url = match xtream_cluster { diff --git a/src/repository/m3u_repository.rs b/src/repository/m3u_repository.rs index 5beb186f8..1329bf40a 100644 --- a/src/repository/m3u_repository.rs +++ b/src/repository/m3u_repository.rs @@ -10,7 +10,7 @@ use crate::model::config::{Config, ConfigTarget}; use crate::model::playlist::{M3uPlaylistItem, PlaylistGroup, PlaylistItem, PlaylistItemType}; use crate::repository::indexed_document::{IndexedDocumentDirectAccess, IndexedDocumentWriter}; use crate::repository::m3u_playlist_iterator::M3uPlaylistIterator; -use crate::repository::storage::{FILE_SUFFIX_DB, FILE_SUFFIX_INDEX}; +use crate::repository::storage::{get_target_storage_path, FILE_SUFFIX_DB, FILE_SUFFIX_INDEX}; use crate::utils::file_utils; use crate::utils::file_utils::file_writer; @@ -89,12 +89,14 @@ pub async fn m3u_load_rewrite_playlist( } -pub async fn m3u_get_item_for_stream_id(cfg: &Config, stream_id: u32, m3u_path: &Path, idx_path: &Path) -> Result { +pub async fn m3u_get_item_for_stream_id(stream_id: u32, cfg: &Config, target: &ConfigTarget) -> Result { if stream_id < 1 { return Err(str_to_io_error("id should start with 1")); } { - let _file_lock = cfg.file_locks.read_lock(m3u_path).await?; - IndexedDocumentDirectAccess::read_indexed_item::(m3u_path, idx_path, &stream_id) + let target_path = get_target_storage_path(cfg, target.name.as_str()).ok_or_else(|| str_to_io_error(&format!("Could not find path for target {}", &target.name)))?; + let (m3u_path, idx_path) = m3u_get_file_paths(&target_path); + let _file_lock = cfg.file_locks.read_lock(&m3u_path).await?; + IndexedDocumentDirectAccess::read_indexed_item::(&m3u_path, &idx_path, &stream_id) } } \ No newline at end of file diff --git a/src/repository/playlist_repository.rs b/src/repository/playlist_repository.rs index ca0338c05..5c9de6618 100644 --- a/src/repository/playlist_repository.rs +++ b/src/repository/playlist_repository.rs @@ -1,7 +1,6 @@ use crate::info_err; use crate::m3u_filter_error::{M3uFilterError, M3uFilterErrorKind}; use crate::model::config::{Config, ConfigTarget, TargetType}; -use crate::model::playlist::PlaylistItemType::LiveUnknown; use crate::model::playlist::{PlaylistGroup, PlaylistItemType}; use crate::model::xmltv::Epg; use crate::repository::epg_repository::epg_write; @@ -11,6 +10,8 @@ use crate::repository::storage::{ensure_target_storage_path, get_target_id_mappi use crate::repository::target_id_mapping::TargetIdMapping; use crate::repository::xtream_repository::xtream_write_playlist; +pub const HLS_EXT: &str = ".m3u8"; + pub async fn persist_playlist(playlist: &mut [PlaylistGroup], epg: Option<&Epg>, target: &ConfigTarget, cfg: &Config) -> Result<(), Vec> { let mut errors = vec![]; @@ -37,11 +38,15 @@ pub async fn persist_playlist(playlist: &mut [PlaylistGroup], epg: Option<&Epg>, let mut header = channel.header.borrow_mut(); let provider_id = header.get_provider_id().unwrap_or_default(); if provider_id == 0 { - header.item_type = if header.url.ends_with(".m3u8") { PlaylistItemType::LiveHls } else { LiveUnknown }; + header.item_type = match (header.url.ends_with(HLS_EXT), header.item_type) { + (true, _) => PlaylistItemType::LiveHls, + (false, PlaylistItemType::Live) => PlaylistItemType::LiveUnknown, + _ => header.item_type, + }; } let uuid = header.get_uuid(); let item_type = header.item_type; - header.virtual_id = target_id_mapping.get_virtual_id(**uuid, provider_id, item_type, 0); + header.virtual_id = target_id_mapping.get_and_update_virtual_id(**uuid, provider_id, item_type, 0); } } @@ -51,17 +56,11 @@ pub async fn persist_playlist(playlist: &mut [PlaylistGroup], epg: Option<&Epg>, TargetType::Xtream => xtream_write_playlist(target, cfg, playlist).await, TargetType::Strm => kodi_write_strm_playlist(target, cfg, playlist, output).await, }; - if let Err(err) = result { errors.push(err); - } else { - if let Err(err) = target_id_mapping.persist() { - errors.push(info_err!(err.to_string())); - } - if !playlist.is_empty() { - if let Err(err) = epg_write(target, cfg, &target_path, epg, output) { - errors.push(err); - } + } else if !playlist.is_empty() { + if let Err(err) = epg_write(target, cfg, &target_path, epg, output) { + errors.push(err); } } } diff --git a/src/repository/target_id_mapping.rs b/src/repository/target_id_mapping.rs index e858cc1ac..5b7054dbf 100644 --- a/src/repository/target_id_mapping.rs +++ b/src/repository/target_id_mapping.rs @@ -71,7 +71,20 @@ impl TargetIdMapping { } } - pub fn get_virtual_id(&mut self, uuid: UUIDType, provider_id: u32, item_type: PlaylistItemType, parent_virtual_id: u32) -> u32 { + // pub fn get_virtual_id(&mut self, uuid: UUIDType, provider_id: u32, item_type: PlaylistItemType, parent_virtual_id: u32) -> u32 { + // match self.by_uuid.get(&uuid) { + // None => { + // self.dirty = true; + // self.virtual_id_counter += 1; + // let record = VirtualIdRecord::new(provider_id, self.virtual_id_counter, item_type, parent_virtual_id, uuid); + // self.by_virtual_id.insert(self.virtual_id_counter, record); + // self.virtual_id_counter + // } + // Some(virtual_id) => *virtual_id + // } + // } + + pub fn get_and_update_virtual_id(&mut self, uuid: UUIDType, provider_id: u32, item_type: PlaylistItemType, parent_virtual_id: u32) -> u32 { match self.by_uuid.get(&uuid) { None => { self.dirty = true; @@ -80,7 +93,16 @@ impl TargetIdMapping { self.by_virtual_id.insert(self.virtual_id_counter, record); self.virtual_id_counter } - Some(record) => *record + Some(virtual_id) => { + if let Some(record) = self.by_virtual_id.query(virtual_id) { + if record.provider_id != provider_id || record.item_type != item_type || record.parent_virtual_id != parent_virtual_id { + let new_record = VirtualIdRecord::new(provider_id, *virtual_id, item_type, parent_virtual_id, uuid); + self.by_virtual_id.insert(*virtual_id, new_record); + self.dirty = true; + } + } + *virtual_id + } } } @@ -99,4 +121,20 @@ impl Drop for TargetIdMapping { error!("Failed to persist target id mapping {:?} err:{err}", &self.path); } } +} + +#[cfg(test)] +mod tests { + use std::path::PathBuf; + use crate::repository::bplustree::BPlusTree; + use crate::repository::target_id_mapping::{VirtualIdRecord}; + + #[test] + fn test_id_mapping() { + let path = PathBuf::from("/home/euzuner/projects/m3u-test/settings/m3u-catbox/data/m3u/id_mapping.db"); + let mapping = BPlusTree::::load(&path); + mapping.unwrap().traverse(|keys, values| { + println!("{keys:?} {values:?}"); + }); + } } \ No newline at end of file diff --git a/src/repository/xtream_playlist_iterator.rs b/src/repository/xtream_playlist_iterator.rs index 23d9ea6b1..a07a3ce9f 100644 --- a/src/repository/xtream_playlist_iterator.rs +++ b/src/repository/xtream_playlist_iterator.rs @@ -28,6 +28,9 @@ impl XtreamPlaylistIterator { ) -> Result { if let Some(storage_path) = xtream_get_storage_path(config, target.name.as_str()) { let (xtream_path, idx_path) = xtream_get_file_paths(&storage_path, cluster); + if !xtream_path.exists() || !idx_path.exists() { + return Err(info_err!(format!("No {cluster} entries found for target {}", &target.name))); + } let file_lock = config.file_locks.read_lock(&xtream_path).await .map_err(|err| info_err!(format!("Could not lock document {xtream_path:?}: {err}")))?; diff --git a/src/repository/xtream_repository.rs b/src/repository/xtream_repository.rs index 1699fc284..965bcca77 100644 --- a/src/repository/xtream_repository.rs +++ b/src/repository/xtream_repository.rs @@ -160,7 +160,7 @@ fn get_map_item_as_str(map: &serde_json::Map, key: &str) -> Optio fn load_old_category_ids(path: &Path) -> (u32, HashMap) { let mut result: HashMap = HashMap::new(); let mut max_id: u32 = 0; - for cat in [COL_CAT_LIVE, COL_CAT_VOD, COL_CAT_SERIES] { + for (cluster, cat) in [(XtreamCluster::Live, COL_CAT_LIVE), (XtreamCluster::Video, COL_CAT_VOD), (XtreamCluster::Series, COL_CAT_SERIES)] { let col_path = get_collection_path(path, cat); if col_path.exists() { if let Ok(file) = File::open(col_path) { @@ -169,7 +169,7 @@ fn load_old_category_ids(path: &Path) -> (u32, HashMap) { if let Some(category_id) = entry.get(TAG_CATEGORY_ID).and_then(get_u32_from_serde_value) { if let Value::Object(item) = entry { if let Some(category_name) = get_map_item_as_str(&item, TAG_CATEGORY_NAME) { - result.insert(category_name, category_id); + result.insert(format!("{cluster}{category_name}"), category_id); max_id = max_id.max(category_id); } } @@ -224,19 +224,20 @@ pub async fn xtream_write_playlist( ) -> Result<(), M3uFilterError> { let path = ensure_xtream_storage_path(cfg, target.name.as_str())?; let mut errors = Vec::new(); - 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::with_capacity(1_000); + let mut cat_series_col = Vec::with_capacity(1_000); + let mut cat_vod_col = Vec::with_capacity(1_000); + let mut live_col = Vec::with_capacity(50_000); + let mut series_col = Vec::with_capacity(10_000); + let mut vod_col = Vec::with_capacity(10_000); // 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(|| { + let cat_key = format!("{}{}", plg.xtream_cluster, &plg.title); + let cat_id = existing_cat_ids.get(&cat_key).unwrap_or_else(|| { cat_id_counter += 1; &cat_id_counter }); @@ -254,30 +255,36 @@ pub async fn xtream_write_playlist( for pli in &plg.channels { let mut header = pli.header.borrow_mut(); - let col = match header.item_type { - PlaylistItemType::LiveUnknown | PlaylistItemType::LiveHls => { - header.category_id = *cat_id; - Some(&mut live_col) - } - _ => { - if header.get_provider_id().is_some() { - header.category_id = *cat_id; - Some(match header.xtream_cluster { - XtreamCluster::Live => &mut live_col, - XtreamCluster::Series => &mut series_col, - XtreamCluster::Video => &mut vod_col, - }) - } else { - let title = header.title.as_str(); - errors.push(format!("Channel does not have an id: {title}")); - None - } - } + header.category_id = *cat_id; + let col = match header.xtream_cluster { + XtreamCluster::Live => &mut live_col, + XtreamCluster::Series => &mut series_col, + XtreamCluster::Video => &mut vod_col, }; + + // let col = match header.item_type { + // PlaylistItemType::LiveUnknown | PlaylistItemType::LiveHls => { + // header.category_id = *cat_id; + // Some(&mut live_col) + // } + // _ => { + // if header.get_provider_id().is_some() { + // header.category_id = *cat_id; + // Some(match header.xtream_cluster { + // XtreamCluster::Live => &mut live_col, + // XtreamCluster::Series => &mut series_col, + // XtreamCluster::Video => &mut vod_col, + // }) + // // } else { + // let title = header.title.as_str(); + // errors.push(format!("Channel does not have an id: {title}")); + // errors.push(format!("Channel does not have an id: {title}")); + // None + // } + // } + // }; drop(header); - if let Some(pl) = col { - pl.push(pli); - } + col.push(pli); } } } @@ -699,7 +706,7 @@ async fn rewrite_xtream_series_info

( if let Some(episode_provider_id) = episode.get(TAG_ID).and_then(get_u32_from_serde_value) { let uuid = generate_playlist_uuid(&hex_encode(&pli.get_uuid()), &episode_provider_id.to_string(), &provider_url); - let episode_virtual_id = target_id_mapping.get_virtual_id( + let episode_virtual_id = target_id_mapping.get_and_update_virtual_id( uuid, episode_provider_id, PlaylistItemType::Series, @@ -720,6 +727,9 @@ async fn rewrite_xtream_series_info

( } } + if let Err(err) = target_id_mapping.persist() { + error!("{}", err.to_string()); + } drop(target_id_mapping); } let result = serde_json::to_string(&doc).map_err(|_| str_to_io_error("Failed to serialize updated series info"))?; diff --git a/src/utils/request_utils.rs b/src/utils/request_utils.rs index 4d10aa3e1..c5bbb1a01 100644 --- a/src/utils/request_utils.rs +++ b/src/utils/request_utils.rs @@ -425,9 +425,22 @@ pub fn classify_content_type(headers: &[(String, String)]) -> MimeCategory { }) } +pub fn replace_extension(path: &str, new_ext: &str) -> String { + let ext = if let Some(stripped) = new_ext.strip_prefix('.') { stripped } else { new_ext }; + if let Some(pos) = path.rfind('/') { + if let Some(dot_pos) = path[pos..].rfind('.') { + let dot_index = pos + dot_pos; + return format!("{}{}.{}", &path[..dot_index], "", ext); + } + } else if let Some(dot_pos) = path.rfind('.') { + return format!("{}{}.{}", &path[..dot_pos], "", ext); + } + format!("{path}.{ext}") +} + #[cfg(test)] mod tests { - use crate::utils::request_utils::{sanitize_sensitive_info}; + use crate::utils::request_utils::{replace_extension, sanitize_sensitive_info}; #[test] fn test_url_mask() { @@ -437,4 +450,20 @@ mod tests { println!("{masked}") } + #[test] + fn test_replace_ext() { + let tests = [ + "test.txt", + "folder/test.txt", + "folder/subfolder/file", + "/absolute/path/to/file.tar.gz", + "/home/user/script", + "no_extension", + "some/path/file.with.dots", + ]; + + for test in &tests { + println!("{} -> {}", test, replace_extension(test, ".mp4")); + } + } } diff --git a/src/utils/xtream_utils.rs b/src/utils/xtream_utils.rs index eb3ff0a94..ecb25337b 100644 --- a/src/utils/xtream_utils.rs +++ b/src/utils/xtream_utils.rs @@ -1,7 +1,7 @@ use crate::Arc; use crate::m3u_filter_error::{str_to_io_error, M3uFilterError}; use crate::model::config::{Config, ConfigInput, ConfigTarget}; -use crate::model::playlist::{PlaylistEntry, PlaylistGroup, XtreamCluster}; +use crate::model::playlist::{PlaylistEntry, PlaylistGroup, XtreamCluster, XtreamPlaylistItem}; use crate::processing::{xtream_parser}; use crate::repository::xtream_repository::{rewrite_xtream_series_info_content, rewrite_xtream_vod_info_content, xtream_get_input_info}; use crate::repository::xtream_repository; @@ -9,7 +9,10 @@ use crate::utils::{request_utils}; use log::{info, warn}; use std::cmp::Ordering; use std::io::{Error}; -use crate::model::api_proxy::{ProxyUserCredentials}; +use crate::model::api_proxy::{ProxyType, ProxyUserCredentials}; +use crate::utils::json_utils::get_string_from_serde_value; +use crate::utils::request_utils::extract_extension_from_url; + pub const ACTION_GET_SERIES_INFO: &str = "get_series_info"; pub const ACTION_GET_VOD_INFO: &str = "get_vod_info"; pub const ACTION_GET_LIVE_INFO: &str = "get_live_info"; @@ -136,7 +139,8 @@ pub async fn get_xtream_playlist(client: Arc, input: &ConfigInp let username = input.username.as_ref().map_or("", |v| v); let password = input.password.as_ref().map_or("", |v| v); - let base_url = format!("{}/player_api.php?username={}&password={}", input.url, username, password); + + let base_url = get_xtream_stream_url_base(&input.url, username, password); if let Err(err) = request_utils::get_input_json_content(Arc::clone(&client), input, base_url.as_str(), None).await { warn!("Failed to login xtream account {username} {err}"); @@ -185,4 +189,29 @@ pub async fn get_xtream_playlist(client: Arc, input: &ConfigInp plg.id = grp_id; } (playlist_groups, errors) +} + +pub fn create_vod_info_from_item(user: &ProxyUserCredentials, pli: &XtreamPlaylistItem, last_updated: i64) -> String { + let category_id = pli.category_id; + let stream_id = if user.proxy == ProxyType::Redirect { pli.virtual_id } else { pli.provider_id }; + let name = &pli.name; + let extension = pli.get_additional_property("container_extension") + .map_or_else(|| extract_extension_from_url(&pli.url).map_or_else (String::new, std::string::ToString::to_string), + |v| get_string_from_serde_value(&v).map_or_else(String::new, |v| v)); + let added = last_updated / 1000; + format!(r#"{{ + "info": {{}}, + "movie_data": {{ + "added": "{added}", + "category_id": {category_id}, + "category_ids": [ + {category_id} + ], + "container_extension": "{extension}", + "custom_sid": "", + "direct_source": "", + "name": "{name}", + "stream_id": {stream_id} + }} +}}"#) } \ No newline at end of file diff --git a/test/rest-api.http b/test/rest-api.http index c2f2a49e2..81b0dea83 100644 --- a/test/rest-api.http +++ b/test/rest-api.http @@ -36,16 +36,19 @@ GET {{local}}/player_api.php?username={{username}}&password={{password}} GET {{local}}/player_api.php?username={{username}}&password={{password}}&action=get_live_categories ### xtream live_streams -GET {{local}}/player_api.php?username={{username}}&password={{password}}&action=get_live_streams&category_id=1 +GET {{local}}/player_api.php?username={{username}}&password={{password}}&action=get_live_streams&category_id=7 ### xtream vod_categories GET {{local}}/player_api.php?username={{username}}&password={{password}}&action=get_vod_categories +### xtream vod streams for category +GET {{local}}/player_api.php?username={{username}}&password={{password}}&action=get_vod_streams&category_id=6 + ### xtream vod streams -GET {{local}}/player_api.php?username={{username}}&password={{password}}&action=get_vod_streams&category_id=29 +GET {{silver}}/player_api.php?username={{username}}&password={{password}}&action=get_vod_streams ### xtream vod info -GET {{local}}/player_api.php?username={{username}}&password={{password}}&action=get_vod_info&vod_id=38497 +GET {{silver}}/player_api.php?username={{username}}&password={{password}}&action=get_vod_info&vod_id=5564 ### xtream series_categories GET {{local}}/player_api.php?username={{username}}&password={{password}}&action=get_series_categories