diff --git a/CHANGELOG.md b/CHANGELOG.md index 3d02844c7..6eec6d87b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,6 @@ # Changelog # 2.2.3 (2023-04-xx) +- hls reverse proxy - input alias definition for same provider with same content but different credentials ```yaml - sources: diff --git a/Cargo.lock b/Cargo.lock index f1114b9bd..b7b382ed4 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -174,7 +174,7 @@ dependencies = [ "bytes", "form_urlencoded", "futures-util", - "http 1.2.0", + "http 1.3.1", "http-body", "http-body-util", "hyper", @@ -206,7 +206,7 @@ checksum = "df1362f362fd16024ae199c1970ce98f9661bf5ef94b9808fee734bc3698b733" dependencies = [ "bytes", "futures-util", - "http 1.2.0", + "http 1.3.1", "http-body", "http-body-util", "mime", @@ -943,7 +943,7 @@ dependencies = [ "fnv", "futures-core", "futures-sink", - "http 1.2.0", + "http 1.3.1", "indexmap", "slab", "tokio", @@ -976,9 +976,9 @@ dependencies = [ [[package]] name = "http" -version = "1.2.0" +version = "1.3.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f16ca2af56261c99fba8bac40a10251ce8188205a4c448fbb745a2e4daa76fea" +checksum = "f4a85d31aea989eead29a3aaf9e1115a180df8282431156e533de47660892565" dependencies = [ "bytes", "fnv", @@ -992,18 +992,18 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1efedce1fb8e6913f23e0c92de8e62cd5b772a67e7b3946df930a62566c93184" dependencies = [ "bytes", - "http 1.2.0", + "http 1.3.1", ] [[package]] name = "http-body-util" -version = "0.1.2" +version = "0.1.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "793429d76616a256bcb62c2a2ec2bed781c8307e797e2598c50010f2bee2544f" +checksum = "b021d93e26becf5dc7e1b75b1bed1fd93124b374ceb73f43d4d4eafec896a64a" dependencies = [ "bytes", - "futures-util", - "http 1.2.0", + "futures-core", + "http 1.3.1", "http-body", "pin-project-lite", ] @@ -1036,7 +1036,7 @@ dependencies = [ "futures-channel", "futures-util", "h2", - "http 1.2.0", + "http 1.3.1", "http-body", "httparse", "httpdate", @@ -1054,7 +1054,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2d191583f3da1305256f22463b9bb0471acad48a4e534a5218b9963e9c1f59b2" dependencies = [ "futures-util", - "http 1.2.0", + "http 1.3.1", "hyper", "hyper-util", "rustls", @@ -1090,7 +1090,7 @@ dependencies = [ "bytes", "futures-channel", "futures-util", - "http 1.2.0", + "http 1.3.1", "http-body", "hyper", "pin-project-lite", @@ -1943,9 +1943,9 @@ dependencies = [ [[package]] name = "quote" -version = "1.0.39" +version = "1.0.40" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c1f1914ce909e1658d9907913b4b91947430c7d9be598b15a1912935b8c04801" +checksum = "1885c039570dc00dcb4ff087a89e185fd56bae234ddc7f056a945bf36467248d" dependencies = [ "proc-macro2", ] @@ -2050,9 +2050,9 @@ checksum = "2b15c43186be67a4fd63bee50d0303afffcef381492ebe2c5d87f324e1b8815c" [[package]] name = "reqwest" -version = "0.12.12" +version = "0.12.14" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "43e734407157c3c2034e0258f5e4473ddb361b1e85f95a66690d67264d7cd1da" +checksum = "989e327e510263980e231de548a33e63d34962d29ae61b467389a1a09627a254" dependencies = [ "base64 0.22.1", "bytes", @@ -2061,7 +2061,7 @@ dependencies = [ "futures-core", "futures-util", "h2", - "http 1.2.0", + "http 1.3.1", "http-body", "http-body-util", "hyper", @@ -2102,9 +2102,9 @@ dependencies = [ [[package]] name = "ring" -version = "0.17.13" +version = "0.17.14" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "70ac5d832aa16abd7d1def883a8545280c20a60f523a370aa3a9617c2b8550ee" +checksum = "a4689e6c2294d81e88dc6261c768b63bc4fcdb852be6d1352498b114f61383b7" dependencies = [ "cc", "cfg-if", @@ -2582,9 +2582,9 @@ checksum = "1f3ccbac311fea05f86f61904b462b55fb3df8837a366dfc601a0161d0532f20" [[package]] name = "tokio" -version = "1.44.0" +version = "1.44.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9975ea0f48b5aa3972bf2d888c238182458437cc2a19374b81b25cdf1023fb3a" +checksum = "f382da615b842244d4b8738c82ed1275e6c5dd90c459a30941cd07080b06c91a" dependencies = [ "backtrace", "bytes", @@ -2642,9 +2642,9 @@ dependencies = [ [[package]] name = "tokio-util" -version = "0.7.13" +version = "0.7.14" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d7fcaa8d55a2bdd6b83ace262b016eca0d79ee02818c5c1bcdf0305114081078" +checksum = "6b9590b93e6fcc1739458317cccd391ad3955e2bde8913edf6f95f9e65a8f034" dependencies = [ "bytes", "futures-core", @@ -2681,7 +2681,7 @@ dependencies = [ "bytes", "futures-core", "futures-util", - "http 1.2.0", + "http 1.3.1", "http-body", "http-body-util", "http-range-header", @@ -3064,32 +3064,31 @@ checksum = "6dccfd733ce2b1753b03b6d3c65edf020262ea35e20ccdf3e288043e6dd620e3" [[package]] name = "windows-registry" -version = "0.2.0" +version = "0.4.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e400001bb720a623c1c69032f8e3e4cf09984deec740f007dd2b03ec864804b0" +checksum = "4286ad90ddb45071efd1a66dfa43eb02dd0dfbae1545ad6cc3c51cf34d7e8ba3" dependencies = [ "windows-result", "windows-strings", - "windows-targets 0.52.6", + "windows-targets 0.53.0", ] [[package]] name = "windows-result" -version = "0.2.0" +version = "0.3.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1d1043d8214f791817bab27572aaa8af63732e11bf84aa21a45a78d6c317ae0e" +checksum = "06374efe858fab7e4f881500e6e86ec8bc28f9462c47e5a9941a0142ad86b189" dependencies = [ - "windows-targets 0.52.6", + "windows-link", ] [[package]] name = "windows-strings" -version = "0.1.0" +version = "0.3.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4cd9b125c486025df0eabcb585e62173c6c9eddcec5d117d3b6e8c30e2ee4d10" +checksum = "87fa48cc5d406560701792be122a10132491cff9d0aeb23583cc2dcafc847319" dependencies = [ - "windows-result", - "windows-targets 0.52.6", + "windows-link", ] [[package]] @@ -3143,13 +3142,29 @@ dependencies = [ "windows_aarch64_gnullvm 0.52.6", "windows_aarch64_msvc 0.52.6", "windows_i686_gnu 0.52.6", - "windows_i686_gnullvm", + "windows_i686_gnullvm 0.52.6", "windows_i686_msvc 0.52.6", "windows_x86_64_gnu 0.52.6", "windows_x86_64_gnullvm 0.52.6", "windows_x86_64_msvc 0.52.6", ] +[[package]] +name = "windows-targets" +version = "0.53.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b1e4c7e8ceaaf9cb7d7507c974735728ab453b67ef8f18febdd7c11fe59dca8b" +dependencies = [ + "windows_aarch64_gnullvm 0.53.0", + "windows_aarch64_msvc 0.53.0", + "windows_i686_gnu 0.53.0", + "windows_i686_gnullvm 0.53.0", + "windows_i686_msvc 0.53.0", + "windows_x86_64_gnu 0.53.0", + "windows_x86_64_gnullvm 0.53.0", + "windows_x86_64_msvc 0.53.0", +] + [[package]] name = "windows_aarch64_gnullvm" version = "0.48.5" @@ -3162,6 +3177,12 @@ version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "32a4622180e7a0ec044bb555404c800bc9fd9ec262ec147edd5989ccd0c02cd3" +[[package]] +name = "windows_aarch64_gnullvm" +version = "0.53.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "86b8d5f90ddd19cb4a147a5fa63ca848db3df085e25fee3cc10b39b6eebae764" + [[package]] name = "windows_aarch64_msvc" version = "0.48.5" @@ -3174,6 +3195,12 @@ version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "09ec2a7bb152e2252b53fa7803150007879548bc709c039df7627cabbd05d469" +[[package]] +name = "windows_aarch64_msvc" +version = "0.53.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c7651a1f62a11b8cbd5e0d42526e55f2c99886c77e007179efff86c2b137e66c" + [[package]] name = "windows_i686_gnu" version = "0.48.5" @@ -3186,12 +3213,24 @@ version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8e9b5ad5ab802e97eb8e295ac6720e509ee4c243f69d781394014ebfe8bbfa0b" +[[package]] +name = "windows_i686_gnu" +version = "0.53.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c1dc67659d35f387f5f6c479dc4e28f1d4bb90ddd1a5d3da2e5d97b42d6272c3" + [[package]] name = "windows_i686_gnullvm" version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0eee52d38c090b3caa76c563b86c3a4bd71ef1a819287c19d586d7334ae8ed66" +[[package]] +name = "windows_i686_gnullvm" +version = "0.53.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9ce6ccbdedbf6d6354471319e781c0dfef054c81fbc7cf83f338a4296c0cae11" + [[package]] name = "windows_i686_msvc" version = "0.48.5" @@ -3204,6 +3243,12 @@ version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "240948bc05c5e7c6dabba28bf89d89ffce3e303022809e73deaefe4f6ec56c66" +[[package]] +name = "windows_i686_msvc" +version = "0.53.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "581fee95406bb13382d2f65cd4a908ca7b1e4c2f1917f143ba16efe98a589b5d" + [[package]] name = "windows_x86_64_gnu" version = "0.48.5" @@ -3216,6 +3261,12 @@ version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "147a5c80aabfbf0c7d901cb5895d1de30ef2907eb21fbbab29ca94c5b08b1a78" +[[package]] +name = "windows_x86_64_gnu" +version = "0.53.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2e55b5ac9ea33f2fc1716d1742db15574fd6fc8dadc51caab1c16a3d3b4190ba" + [[package]] name = "windows_x86_64_gnullvm" version = "0.48.5" @@ -3228,6 +3279,12 @@ version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "24d5b23dc417412679681396f2b49f3de8c1473deb516bd34410872eff51ed0d" +[[package]] +name = "windows_x86_64_gnullvm" +version = "0.53.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0a6e035dd0599267ce1ee132e51c27dd29437f63325753051e71dd9e42406c57" + [[package]] name = "windows_x86_64_msvc" version = "0.48.5" @@ -3240,6 +3297,12 @@ version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "589f6da84c646204747d1270a2a5661ea66ed1cced2631d546fdfb155959f9ec" +[[package]] +name = "windows_x86_64_msvc" +version = "0.53.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "271414315aff87387382ec3d271b52d7ae78726f5d44ac98b4f4030c91880486" + [[package]] name = "winnow" version = "0.6.26" diff --git a/src/api/endpoints/hls_api.rs b/src/api/endpoints/hls_api.rs index 1426d5870..851cb13df 100644 --- a/src/api/endpoints/hls_api.rs +++ b/src/api/endpoints/hls_api.rs @@ -1,13 +1,11 @@ -use crate::api::api_utils::{get_user_target_by_credentials, stream_response}; -use crate::api::api_utils::{try_option_bad_request, try_result_bad_request}; +use crate::api::api_utils::{stream_response}; +use crate::api::api_utils::{try_option_bad_request}; 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::parser::hls::{rewrite_hls, M3U_HLSR_PREFIX}; +use crate::model::playlist::{PlaylistItemType, XtreamCluster}; +use crate::processing::parser::hls::rewrite_hls; use crate::repository::playlist_repository::HLS_EXT; -use crate::repository::{m3u_repository, xtream_repository}; use crate::utils::network::request; use crate::utils::network::request::{replace_extension, sanitize_sensitive_info}; use axum::response::IntoResponse; @@ -21,16 +19,22 @@ struct HlsApiPathParams { token: String, username: String, password: String, - channel: String, - hash: String, - chunk: String, + stream_id: u32, + chunk: u32, } -pub(in crate::api) async fn handle_hls_stream_request(app_state: &Arc, user: &ProxyUserCredentials, pli: &dyn PlaylistEntry, input: &ConfigInput, target_type: TargetType) -> impl axum::response::IntoResponse + Send { - let url = replace_extension(&pli.get_provider_url(), HLS_EXT); +pub(in crate::api) async fn handle_hls_stream_request(app_state: &Arc, + user: &ProxyUserCredentials, + hls_url: &str, + virtual_id: u32, + input: &ConfigInput, + target_type: TargetType) -> impl axum::response::IntoResponse + Send { + let url = replace_extension(hls_url, HLS_EXT); + let server_info = app_state.config.get_user_server_info(user).await; match request::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); + let (hls_entry, hls_content) = rewrite_hls(&server_info.get_base_url(), &content, hls_url, virtual_id, user, &target_type, &input.name); + app_state.hls_cache.add_entry(hls_entry).await; axum::response::Response::builder() .status(axum::http::StatusCode::OK) .header(axum::http::header::CONTENT_TYPE, "application/x-mpegurl") @@ -46,62 +50,40 @@ pub(in crate::api) async fn handle_hls_stream_request(app_state: &Arc, } async fn hls_api_stream( - req_headers: &axum::http::HeaderMap, - api_req: &UserApiRequest, - params: HlsApiPathParams, - app_state: &Arc, - target_type: TargetType, + req_headers: axum::http::HeaderMap, + axum::extract::Path(params): axum::extract::Path, + axum::extract::State(app_state): axum::extract::State>, ) -> impl axum::response::IntoResponse + Send { let (user, target) = try_option_bad_request!( - get_user_target_by_credentials(¶ms.username, ¶ms.password, api_req, app_state).await, - false, + app_state.config.get_target_for_user(¶ms.username, ¶ms.password).await, false, format!("Could not find any user {}", params.username)); - if !user.has_permissions(app_state).await { + if !user.has_permissions(&app_state).await { return axum::http::StatusCode::FORBIDDEN.into_response(); } + let Some(hls_entry) = app_state.hls_cache.get_entry(¶ms.token).await else { return axum::http::StatusCode::BAD_REQUEST.into_response(); }; + let Some(hls_url) = hls_entry.get_chunk_url(params.chunk) else { return axum::http::StatusCode::BAD_REQUEST.into_response(); }; let target_name = &target.name; - let virtual_id: u32 = try_result_bad_request!(params.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), 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(); + let virtual_id = params.stream_id; + let input = try_option_bad_request!(app_state.config.get_input_by_name(&hls_entry.input_name), true, format!("Cant find input for target {target_name}, context {}, stream_id {virtual_id}", XtreamCluster::Live)); - // 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_headers, Some(input), PlaylistItemType::Live, target, &user).await.into_response() -} + if hls_url.ends_with(HLS_EXT) { + return handle_hls_stream_request(&app_state, &user, hls_url, virtual_id, input, hls_entry.target_type.clone()).await.into_response(); + } -async fn hls_api_stream_xtream( - req_headers: axum::http::HeaderMap, - axum::extract::Query(api_req): axum::extract::Query, - axum::extract::Path(params): axum::extract::Path, - axum::extract::State(app_state): axum::extract::State>, -) -> impl axum::response::IntoResponse + Send { - hls_api_stream(&req_headers, &api_req, params, &app_state, TargetType::Xtream).await.into_response() -} - -async fn hls_api_stream_m3u( - req_headers: axum::http::HeaderMap, - axum::extract::Query(api_req): axum::extract::Query, - axum::extract::Path(params): axum::extract::Path, - axum::extract::State(app_state): axum::extract::State>, -) -> impl axum::response::IntoResponse + Send { - hls_api_stream(&req_headers, &api_req, params, &app_state, TargetType::M3u).await.into_response() + // let (pli_url, input_name) = if hls_entry.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), 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) + // }; + stream_response(&app_state, hls_url, &req_headers, Some(input), PlaylistItemType::Live, target, &user).await.into_response() } pub fn hls_api_register() -> axum::Router> { axum::Router::new() - .route("/hlsr/{token}/{username}/{password}/{channel}/{hash}/{chunk}", axum::routing::get(hls_api_stream_xtream)) - .route(&format!("/{M3U_HLSR_PREFIX}/{{token}}/{{username}}/{{password}}/{{channel}}/{{hash}}/{{chunk}}"), axum::routing::get(hls_api_stream_m3u)) + .route("/hls/{token}/{username}/{password}/{stream_id}/{chunk}", axum::routing::get(hls_api_stream)) //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/endpoints/m3u_api.rs b/src/api/endpoints/m3u_api.rs index edaafbd57..18c1649dc 100644 --- a/src/api/endpoints/m3u_api.rs +++ b/src/api/endpoints/m3u_api.rs @@ -4,7 +4,7 @@ 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, XtreamCluster}; +use crate::model::playlist::{FieldGetAccessor, PlaylistItemType, 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::repository::playlist_repository::HLS_EXT; @@ -87,7 +87,7 @@ async fn m3u_api_stream( let input = app_state.config.get_input_by_name(m3u_item.input_name.as_str()); - let is_hls_request = stream_ext.as_deref() == Some(HLS_EXT); + let is_hls_request = m3u_item.item_type == PlaylistItemType::LiveHls || stream_ext.as_deref() == Some(HLS_EXT); if user.proxy == ProxyType::Redirect { let redirect_url = if is_hls_request { &replace_extension(&m3u_item.url, "m3u8") } else { &m3u_item.url }; @@ -98,9 +98,8 @@ async fn m3u_api_stream( // Reverse proxy mode if is_hls_request { let target_name = &target.name; - let hls_input = try_option_bad_request!(input, 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, hls_input, TargetType::M3u).await.into_response(); + let hls_input = try_option_bad_request!(input, 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.url, m3u_item.virtual_id, hls_input, TargetType::M3u).await.into_response(); } stream_response(&app_state, m3u_item.url.as_str(), &req_headers, input, m3u_item.item_type, target, &user).await.into_response() diff --git a/src/api/endpoints/web_index.rs b/src/api/endpoints/web_index.rs index f9d2622cf..5d98dc15a 100644 --- a/src/api/endpoints/web_index.rs +++ b/src/api/endpoints/web_index.rs @@ -96,13 +96,3 @@ pub fn index_register(web_dir_path: &Path) -> axum::Router> { .route("/", axum::routing::get(index)) .fallback(axum::routing::get_service(tower_http::services::ServeDir::new(web_dir_path)))) } -// pub fn index_register(web_dir_path: &Path) -> impl Fn(&mut web::ServiceConfig) + '_ { -// 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("", web_dir_path))); -// } -// } \ No newline at end of file diff --git a/src/api/endpoints/xmltv_api.rs b/src/api/endpoints/xmltv_api.rs index 078069c0a..a71e359fa 100644 --- a/src/api/endpoints/xmltv_api.rs +++ b/src/api/endpoints/xmltv_api.rs @@ -5,7 +5,6 @@ use axum::response::IntoResponse; use chrono::{Duration, NaiveDateTime, TimeDelta}; use flate2::write::GzEncoder; use flate2::Compression; -// use actix_web::{http::header, web, HttpRequest, HttpResponse}; use log::{error, trace}; use quick_xml::events::{BytesStart, Event}; use quick_xml::{Reader, Writer}; diff --git a/src/api/endpoints/xtream_api.rs b/src/api/endpoints/xtream_api.rs index d1b0dbd3f..88a5b6e26 100644 --- a/src/api/endpoints/xtream_api.rs +++ b/src/api/endpoints/xtream_api.rs @@ -34,7 +34,7 @@ use crate::repository::{user_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::network::request::{extract_extension_from_url, replace_extension, sanitize_sensitive_info}; +use crate::utils::network::request::{extract_extension_from_url, sanitize_sensitive_info}; use crate::utils::network::xtream::{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; use crate::utils::debug_if_enabled; @@ -159,12 +159,7 @@ async fn xtream_player_api_stream( let (pli, mapping) = try_result_bad_request!(xtream_repository::xtream_get_item_for_stream_id(virtual_id, &app_state.config, target, None), true, format!("Failed to read xtream item for stream id {}", virtual_id)); let input = try_option_bad_request!(app_state.config.get_input_by_name(pli.input_name.as_str()), true, format!("Cant find input for target {target_name}, context {}, stream_id {virtual_id}", stream_req.context)); - if pli.item_type == PlaylistItemType::LiveHls { - debug_if_enabled!("Redirecting stream request to {}", sanitize_sensitive_info(&pli.url)); - return redirect(&pli.url).into_response(); - } - - let is_hls_request = stream_ext.as_deref() == Some(HLS_EXT); + let is_hls_request = pli.item_type == PlaylistItemType::LiveHls || stream_ext.as_deref() == Some(HLS_EXT); if user.proxy == ProxyType::Redirect { if pli.xtream_cluster == XtreamCluster::Series { @@ -176,14 +171,19 @@ async fn xtream_player_api_stream( return redirect(&stream_url).into_response(); } - 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 redirect(redirect_url.as_str()).into_response(); + // if pli.item_type == PlaylistItemType::LiveHls { + // let redirect_url = &replace_extension(&pli.url, "m3u8"); + // debug_if_enabled!("Redirecting stream request to {}", sanitize_sensitive_info(redirect_url)); + // return redirect(redirect_url).into_response(); + // } + + debug_if_enabled!("Redirecting stream request to {}", sanitize_sensitive_info(&pli.url)); + return redirect(&pli.url).into_response(); } // Reverse proxy mode if is_hls_request { - return handle_hls_stream_request(app_state, &user, &pli, input, TargetType::Xtream).await.into_response(); + return handle_hls_stream_request(app_state, &user, &pli.url, pli.virtual_id, input, TargetType::Xtream).await.into_response(); } let extension = stream_ext.unwrap_or_else( diff --git a/src/api/main_api.rs b/src/api/main_api.rs index b5e30d906..ad61d39cc 100644 --- a/src/api/main_api.rs +++ b/src/api/main_api.rs @@ -27,6 +27,7 @@ use std::sync::Arc; use tokio::sync::Mutex; use std::future::IntoFuture; use crate::api::model::event_manager::EventManager; +use crate::api::model::hls_cache::HlsCache; fn get_web_dir_path(web_ui_enabled: bool, web_root: &str) -> Result { let web_dir = web_root.to_string(); @@ -124,8 +125,9 @@ fn create_shared_data(cfg: &Arc) -> AppState { AppState { config: Arc::clone(cfg), http_client: Arc::new(reqwest::Client::new()), - downloads: Arc::from(DownloadQueue::new()), + downloads: Arc::new(DownloadQueue::new()), cache, + hls_cache: HlsCache::garbage_collected(), shared_stream_manager: Arc::new(SharedStreamManager::new()), active_users, active_provider, @@ -303,34 +305,4 @@ pub async fn start_server(cfg: Arc, targets: Arc) -> fut let router: axum::Router<()> = router.with_state(shared_data.clone()); let listener = tokio::net::TcpListener::bind(format!("{host}:{port}")).await?; axum::serve(listener, router).into_future().await - - // HttpServer::new(move || { - // App::new() - // .wrap(Logger::default()) - // .wrap(Cors::default() - // .supports_credentials() - // .allow_any_origin() - // .allowed_methods(vec!["GET", "POST", "OPTIONS", "HEAD"]) - // .allow_any_header() - // .max_age(3600)) - // .app_data(shared_data.clone()) - // // .wrap(Condition::new(web_auth_enabled, ErrorHandlers::new().handler(StatusCode::UNAUTHORIZED, handle_unauthorized))) - // .configure(|srvcfg| { - // if web_ui_enabled { - // srvcfg.service(actix_files::Files::new("/static", web_dir_path.join("static"))); - // srvcfg.configure(v1_api_register(web_auth_enabled)); - // } - // srvcfg.service(web::resource("/healthcheck").route(web::get().to(healthcheck))); - // srvcfg.service(web::resource("/status").route(web::get().to(status))); - // }) - // .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)); - // } - // }) - // }).bind(format!("{host}:{port}"))?.run().await } diff --git a/src/api/model/app_state.rs b/src/api/model/app_state.rs index 96996743f..3b0a574a2 100644 --- a/src/api/model/app_state.rs +++ b/src/api/model/app_state.rs @@ -4,6 +4,7 @@ use crate::api::model::active_provider_manager::ActiveProviderManager; use crate::api::model::active_user_manager::ActiveUserManager; use crate::api::model::download::DownloadQueue; use crate::api::model::event_manager::EventManager; +use crate::api::model::hls_cache::HlsCache; use crate::api::model::streams::shared_stream_manager::SharedStreamManager; use crate::model::config::{Config}; use crate::model::hdhomerun_config::HdHomeRunDeviceConfig; @@ -16,6 +17,7 @@ pub struct AppState { pub http_client: Arc, pub downloads: Arc, pub cache: Arc>>, + pub hls_cache: Arc, pub shared_stream_manager: Arc, pub active_users: Arc, pub active_provider: Arc, diff --git a/src/api/model/hls_cache.rs b/src/api/model/hls_cache.rs new file mode 100644 index 000000000..ca038e514 --- /dev/null +++ b/src/api/model/hls_cache.rs @@ -0,0 +1,61 @@ +use log::error; +use std::collections::HashMap; +use std::str::FromStr; +use std::sync::Arc; +use std::time::Duration; +use chrono::Local; +use cron::Schedule; +use tokio::sync::RwLock; +use tokio::time::Instant; +use crate::utils::sys_utils::exit; +use crate::model::hls::HlsEntry; + +const EXPIRE_DURATION: u64 = 600; // 10 minutes + +fn start_garbage_collector(cache: &Arc) { + let cache_clone = cache.clone(); + tokio::spawn(async move { + match Schedule::from_str("0 */15 * * * * *") { + Ok(schedule) => { + let offset = *Local::now().offset(); + loop { + let mut upcoming = schedule.upcoming(offset).take(1); + if let Some(datetime) = upcoming.next() { + tokio::time::sleep_until(tokio::time::Instant::from(crate::api::scheduler::datetime_to_instant(datetime))).await; + cache_clone.gc().await; + } + } + } + Err(err) => exit!("Failed to start scheduler: {}", err) + } + }); +} + +pub struct HlsCache { + pub entries: RwLock>, +} + +impl HlsCache { + pub fn garbage_collected() -> Arc { + let cache = Arc::new(Self { + entries: RwLock::new(HashMap::new()), + }); + + start_garbage_collector(&cache); + cache + } + + pub async fn add_entry(&self, entry: HlsEntry) { + self.entries.write().await.insert(entry.token.to_string(), entry); + } + + pub async fn get_entry(&self, token: &str) -> Option{ + self.entries.read().await.get(token).cloned() + } + + pub async fn gc(&self) { + let threshold = Instant::now() - Duration::from_secs(EXPIRE_DURATION); + // Remove all expired elements + self.entries.write().await.retain(|_, entry| entry.ts > threshold); + } +} \ No newline at end of file diff --git a/src/api/model/mod.rs b/src/api/model/mod.rs index 13ab78872..51439073d 100644 --- a/src/api/model/mod.rs +++ b/src/api/model/mod.rs @@ -9,3 +9,4 @@ pub(crate) mod streams; pub(in crate::api) mod active_user_manager; pub(in crate::api) mod active_provider_manager; pub(in crate::api) mod event_manager; +pub(in crate::api) mod hls_cache; diff --git a/src/api/model/streams/provider_stream_factory.rs b/src/api/model/streams/provider_stream_factory.rs index 43c2c993f..8d2a35e56 100644 --- a/src/api/model/streams/provider_stream_factory.rs +++ b/src/api/model/streams/provider_stream_factory.rs @@ -421,7 +421,7 @@ pub async fn create_provider_stream(cfg: &Config, // use std::sync::Arc; // use crate::model::config::Config; // -// #[actix_rt::test] +// #[tokio::test] // async fn test_stream() { // let app = App::new().route("/test", web::get().to(test_stream_handler)); // let server = test::init_service(app).await; diff --git a/src/api/scheduler.rs b/src/api/scheduler.rs index 5dbd32ecf..8f1846286 100644 --- a/src/api/scheduler.rs +++ b/src/api/scheduler.rs @@ -8,7 +8,7 @@ use crate::utils::sys_utils::exit; use crate::model::config::{Config, ProcessTargets}; use crate::processing::processor::playlist::exec_processing; -fn datetime_to_instant(datetime: DateTime) -> Instant { +pub fn datetime_to_instant(datetime: DateTime) -> Instant { // Convert DateTime to SystemTime let target_system_time: SystemTime = datetime.into(); @@ -44,11 +44,12 @@ pub async fn start_scheduler(client: Arc, expression: &str, con mod tests { use std::str::FromStr; use std::sync::atomic::{AtomicU8, Ordering}; + use std::time::Instant; use chrono::Local; use cron::Schedule; use crate::api::scheduler::datetime_to_instant; - #[actix_rt::test] + #[tokio::test] async fn test_run_scheduler() { // Define a cron expression that runs every second let expression = "0/1 * * * * * *"; // every second @@ -63,7 +64,7 @@ mod tests { loop { let mut upcoming = schedule.upcoming(offset).take(1); if let Some(datetime) = upcoming.next() { - tokio::time::sleep_until(actix_rt::time::Instant::from(datetime_to_instant(datetime))).await; + tokio::time::sleep_until(tokio::time::Instant::from(datetime_to_instant(datetime))).await; run_me(); } if runs.load(Ordering::SeqCst) == 6 { diff --git a/src/main.rs b/src/main.rs index afa407d91..4f2a71149 100644 --- a/src/main.rs +++ b/src/main.rs @@ -25,14 +25,9 @@ use env_logger::Builder; use log::{error, info, LevelFilter}; const LOG_ERROR_LEVEL_MOD: &[&str] = &[ - "actix_web::middleware::logger", "reqwest::async_impl::client", "reqwest::connect", "hyper_util::client", - "actix_server::worker", - "actix_server::server", - "actix_server::builder", - "actix_server::accept", ]; diff --git a/src/model/config.rs b/src/model/config.rs index d99717a96..e065c473c 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -52,6 +52,7 @@ macro_rules! valid_property { pub use valid_property; use crate::m3u_filter_error::{create_m3u_filter_error_result, handle_m3u_filter_error_result, handle_m3u_filter_error_result_list}; use crate::model::hdhomerun_config::HdHomeRunConfig; +use crate::utils::file::config_reader::resolve_env_var; use crate::utils::string_utils::get_trimmed_string; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Sequence, PartialEq, Eq, Hash)] @@ -419,12 +420,13 @@ impl ConfigTarget { return create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "unique target name is required for xtream type output: {}", self.name); } } - TargetOutput::M3u(_) => { + TargetOutput::M3u(m3u_output) => { m3u_cnt += 1; + m3u_output.filename = m3u_output.filename.as_ref().map(|s| resolve_env_var(s.trim())); } TargetOutput::Strm(strm_output) => { strm_cnt += 1; - strm_output.directory = strm_output.directory.trim().to_string(); + strm_output.directory = resolve_env_var(strm_output.directory.trim()); if strm_output.directory.trim().is_empty() { return create_m3u_filter_error_result!(M3uFilterErrorKind::Info, "directory is required for strm type: {}", self.name); } @@ -741,6 +743,7 @@ impl ConfigInput { if self.url.is_empty() { return Err(info_err!("url for input is mandatory".to_string())); } + self.url = resolve_env_var(&self.url); self.username = get_trimmed_string(&self.username); self.password = get_trimmed_string(&self.password); check_input_credentials!(self, self.input_type); diff --git a/src/model/hls.rs b/src/model/hls.rs new file mode 100644 index 000000000..b9283cc85 --- /dev/null +++ b/src/model/hls.rs @@ -0,0 +1,20 @@ +use std::collections::HashMap; +use tokio::time::Instant; +use crate::model::config::TargetType; + +#[derive(Clone)] +pub struct HlsEntry { + pub ts: Instant, + pub token: String, + pub target_type: TargetType, + pub input_name: String, + pub virtual_id: u32, + pub chunk: u32, + pub chunks: HashMap, +} + +impl HlsEntry { + pub fn get_chunk_url(&self, chunk: u32) -> Option<&String> { + self.chunks.get(&chunk) + } +} \ No newline at end of file diff --git a/src/model/mod.rs b/src/model/mod.rs index 2ce38cc60..5bb16c1eb 100644 --- a/src/model/mod.rs +++ b/src/model/mod.rs @@ -7,4 +7,5 @@ pub mod xmltv; pub mod xtream; pub mod healthcheck; pub mod playlist_categories; -pub mod hdhomerun_config; \ No newline at end of file +pub mod hdhomerun_config; +pub mod hls; \ No newline at end of file diff --git a/src/processing/parser/hls.rs b/src/processing/parser/hls.rs index c17c88a42..367e39f1e 100644 --- a/src/processing/parser/hls.rs +++ b/src/processing/parser/hls.rs @@ -1,51 +1,45 @@ use crate::model::api_proxy::ProxyUserCredentials; -use std::str; use crate::model::config::TargetType; +use crate::model::hls::HlsEntry; +use crate::repository::storage::hash_string_as_hex; +use std::collections::HashMap; +use std::str; +use tokio::time::Instant; +use crate::utils::string_utils::replace_after_last_slash; -// /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(); +pub const HLS_PREFIX: &str = "hls"; - if parts.len() != 8 || !parts[0].is_empty() || parts[1] != "hlsr" { - return None; +pub fn rewrite_hls(base_url: &str, content: &str, hls_url: &str, virtual_id: u32, user: &ProxyUserCredentials, + target_type: &TargetType, input_name: &str) -> (HlsEntry, String) { + let token = hash_string_as_hex(hls_url); + let username = &user.username; + let password = &user.password; + let mut chunk: u32 = 1; + let mut chunks = HashMap::new(); + let mut result = Vec::new(); + for line in content.lines() { + if line.starts_with('#') { + result.push(line.to_string()); + } else { + let url = if line.starts_with("http") { + line.to_string() + } else { + replace_after_last_slash(hls_url, line) + }; + chunks.insert(chunk, url); + result.push(format!("{base_url}/{HLS_PREFIX}/{token}/{username}/{password}/{virtual_id}/{chunk}")); + chunk += 1; + } } - 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") + let hls = HlsEntry { + ts: Instant::now(), + token, + target_type: target_type.clone(), + input_name: input_name.to_string(), + virtual_id, + chunk, + chunks, + }; + (hls, result.join("\r\n")) } diff --git a/src/repository/kodi_repository.rs b/src/repository/kodi_repository.rs index c23d4fed3..a7e7f4a8c 100644 --- a/src/repository/kodi_repository.rs +++ b/src/repository/kodi_repository.rs @@ -966,7 +966,7 @@ async fn remove_empty_dirs(root_path: PathBuf) { // use crate::repository::kodi_repository::remove_empty_dirs; // use std::path::PathBuf; // -// #[actix_web::test] +// #[tokio::test] // async fn test_empty_dirs() { // remove_empty_dirs(PathBuf::from("/tmp/hello")).await; // } diff --git a/src/repository/m3u_playlist_iterator.rs b/src/repository/m3u_playlist_iterator.rs index 14e1b665a..175e34ad2 100644 --- a/src/repository/m3u_playlist_iterator.rs +++ b/src/repository/m3u_playlist_iterator.rs @@ -103,16 +103,13 @@ impl M3uPlaylistIterator { // TODO hls and unknown reverse proxy entry.map(|(mut m3u_pli, _has_next)| { - let rewrite_urls = match m3u_pli.item_type { - PlaylistItemType::LiveHls => None, - _ => if match &self.proxy_type { - ProxyType::Reverse => true, - ProxyType::Redirect => self.mask_redirect_url, - } { - Some((self.get_stream_url(&m3u_pli, self.include_type_in_url), if self.rewrite_resource { Some(self.get_resource_url(&m3u_pli)) } else { None })) - } else { - None - } + let rewrite_urls = if match &self.proxy_type { + ProxyType::Reverse => true, + ProxyType::Redirect => self.mask_redirect_url, + } { + Some((self.get_stream_url(&m3u_pli, self.include_type_in_url), if self.rewrite_resource { Some(self.get_resource_url(&m3u_pli)) } else { None })) + } else { + None }; let url = m3u_pli.url.to_string(); let (stream_url, resource_url) = rewrite_urls diff --git a/src/utils/file/config_reader.rs b/src/utils/file/config_reader.rs index 620601dcf..1b2310cbb 100644 --- a/src/utils/file/config_reader.rs +++ b/src/utils/file/config_reader.rs @@ -144,6 +144,9 @@ pub fn save_main_config(file_path: &str, backup_dir: &str, config: &ConfigDto) - static ENV_REGEX: LazyLock = LazyLock::new(|| Regex::new(r"\$\{env:(?P[a-zA-Z_][a-zA-Z0-9_]*)}").unwrap()); pub fn resolve_env_var(value: &str) -> String { + if value.is_empty() { + return String::new(); + } ENV_REGEX.replace_all(value, |caps: ®ex::Captures| { let var_name = &caps["var"]; env::var(var_name).unwrap_or_else(|_| format!("${{env:{var_name}}}")) diff --git a/src/utils/string_utils.rs b/src/utils/string_utils.rs index 76e964e76..2f8d77ebb 100644 --- a/src/utils/string_utils.rs +++ b/src/utils/string_utils.rs @@ -41,3 +41,11 @@ pub fn get_trimmed_string(value: &Option) -> Option { } None } + +pub fn replace_after_last_slash(input: &str, replacement: &str) -> String { + match input.rsplitn(2, '/').collect::>().as_slice() { + [_after, before] => format!("{before}/{replacement}"), + [_only] => replacement.to_string(), // if there is no slash, replace complete + _ => input.to_string(), // fallback, should never happen + } +} \ No newline at end of file