diff --git a/CHANGELOG.md b/CHANGELOG.md index b41f9258c..db97a98a0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -86,6 +86,15 @@ url: ['http://localhost:3001/xmltv.php?epg_id=1', 'http://localhost:3001/xmltv.p - `reverse[live,vod]` -> series redirect, others reverse - `/status` api endpoint moved to `/api/v1/status` for auth protection - fixed multi provider VOD seek problem (provider cycle on seek request prevented playback) +- hdhomerun supports now basic auth like +you need to enable auth in config +```yaml +hdhomerun: + enabled: true + auth: true + devices: + - name: hdhr1 +``` # 2.2.5 (2025-03-27) - fixed web ui playlist regexp search diff --git a/Cargo.lock b/Cargo.lock index dd7987ae1..b789a7172 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -108,9 +108,9 @@ dependencies = [ [[package]] name = "anyhow" -version = "1.0.97" +version = "1.0.98" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "dcfed56ad506cb2c684a14971b8861fdc3baaaae314b9e5f9bb532cbe3ba7a4f" +checksum = "e16d2d3311acee920a9eb8d33b8cbc1787ce4a264e85f964c2404b969bdcd487" [[package]] name = "arrayref" @@ -137,9 +137,9 @@ dependencies = [ [[package]] name = "async-compression" -version = "0.4.22" +version = "0.4.23" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "59a194f9d963d8099596278594b3107448656ba73831c9d8c783e613ce86da64" +checksum = "b37fc50485c4f3f736a4fb14199f6d5f5ba008d7f28fe710306c92780f004c07" dependencies = [ "brotli", "flate2", @@ -301,9 +301,9 @@ dependencies = [ [[package]] name = "blake3" -version = "1.8.1" +version = "1.8.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "389a099b34312839e16420d499a9cad9650541715937ffbdd40d36f49e77eeb3" +checksum = "3888aaa89e4b2a40fca9848e400f6a658a5a3978de7be858e209cafa8be9a4a0" dependencies = [ "arrayref", "arrayvec", @@ -323,9 +323,9 @@ dependencies = [ [[package]] name = "brotli" -version = "7.0.0" +version = "8.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "cc97b8f16f944bba54f0433f07e30be199b6dc2bd25937444bbad560bcea29bd" +checksum = "cf19e729cdbd51af9a397fb9ef8ac8378007b797f8273cfbfdf45dcaa316167b" dependencies = [ "alloc-no-stdlib", "alloc-stdlib", @@ -334,9 +334,9 @@ dependencies = [ [[package]] name = "brotli-decompressor" -version = "4.0.2" +version = "5.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "74fa05ad7d803d413eb8380983b092cbbaf9a85f151b871360e7b00cd7060b37" +checksum = "874bb8112abecc98cbd6d81ea4fa7e94fb9449648c93cc89aa40c81c24d7de03" dependencies = [ "alloc-no-stdlib", "alloc-stdlib", @@ -362,9 +362,9 @@ checksum = "a2698f953def977c68f935bb0dfa959375ad4638570e969e2f1e9f433cbf1af6" [[package]] name = "cc" -version = "1.2.18" +version = "1.2.19" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "525046617d8376e3db1deffb079e91cef90a89fc3ca5c185bbf8c9ecdd15cd5c" +checksum = "8e3a13707ac958681c13b39b458c073d0d9bc8a22cb1b2f4c8e55eb72c13f362" dependencies = [ "jobserver", "libc", @@ -399,9 +399,9 @@ dependencies = [ [[package]] name = "clap" -version = "4.5.35" +version = "4.5.37" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d8aa86934b44c19c50f87cc2790e19f54f7a67aedb64101c2e1a2e5ecfb73944" +checksum = "eccb054f56cbd38340b380d4a8e69ef1f02f1af43db2f0cc817a4774d80ae071" dependencies = [ "clap_builder", "clap_derive", @@ -409,9 +409,9 @@ dependencies = [ [[package]] name = "clap_builder" -version = "4.5.35" +version = "4.5.37" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2414dbb2dd0695280da6ea9261e327479e9d37b0630f6b53ba2a11c60c679fd9" +checksum = "efd9466fac8543255d3b1fcad4762c5e116ffe808c8a3043d4263cd4fd4862a2" dependencies = [ "anstream", "anstyle", @@ -973,9 +973,9 @@ dependencies = [ [[package]] name = "getrandom" -version = "0.2.15" +version = "0.2.16" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c4567c8db10ae91089c99af84c68c38da3ec2f087c3f82960bcdbf3656b6f4d7" +checksum = "335ff9f135e4384c8150d6f27c6daed433577f86b4750418338c01a1a2528592" dependencies = [ "cfg-if", "js-sys", @@ -1029,9 +1029,9 @@ dependencies = [ [[package]] name = "h2" -version = "0.4.8" +version = "0.4.9" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5017294ff4bb30944501348f6f8e42e6ad28f42c8bbef7a74029aff064a4e3c2" +checksum = "75249d144030531f8dee69fe9cea04d3edf809a017ae445e2abdff6629e86633" dependencies = [ "atomic-waker", "bytes", @@ -1439,9 +1439,9 @@ checksum = "4a5f13b858c8d314ee3e8f639011f7ccefe71f97f96e50151fb991f267928e2c" [[package]] name = "jiff" -version = "0.2.6" +version = "0.2.10" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1f33145a5cbea837164362c7bd596106eb7c5198f97d1ba6f6ebb3223952e488" +checksum = "5a064218214dc6a10fbae5ec5fa888d80c45d611aba169222fc272072bf7aef6" dependencies = [ "jiff-static", "log", @@ -1452,9 +1452,9 @@ dependencies = [ [[package]] name = "jiff-static" -version = "0.2.6" +version = "0.2.10" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "43ce13c40ec6956157a3635d97a1ee2df323b263f09ea14165131289cb0f5c19" +checksum = "199b7932d97e325aff3a7030e141eafe7f2c6268e1d1b24859b753a627f45254" dependencies = [ "proc-macro2", "quote", @@ -1504,9 +1504,9 @@ checksum = "bbd2bcb4c963f2ddae06a2efc7e9f3591312473c50c6685e1f298068316e66fe" [[package]] name = "libc" -version = "0.2.171" +version = "0.2.172" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c19937216e9d3aa9956d9bb8dfc0b0c8beb6058fc4f7a4dc4d850edf86a237d6" +checksum = "d750af042f7ef4f724306de029d18836c26c1765a54a6a3f094cbd23a7267ffa" [[package]] name = "libnghttp2-sys" @@ -2022,9 +2022,9 @@ dependencies = [ [[package]] name = "proc-macro2" -version = "1.0.94" +version = "1.0.95" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a31971752e70b8b2686d7e46ec17fb38dad4051d94024c88df49b667caea9c84" +checksum = "02b3e5e68a3a1a02aad3ec490a98007cbc13c37cbe84a3cd7b8e406d76e7f778" dependencies = [ "unicode-ident", ] @@ -2075,9 +2075,9 @@ dependencies = [ [[package]] name = "quinn-proto" -version = "0.11.10" +version = "0.11.11" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b820744eb4dc9b57a3398183639c511b5a26d2ed702cedd3febaa1393caa22cc" +checksum = "bcbafbbdbb0f638fe3f35f3c56739f77a8a1d070cb25603226c83339b391472b" dependencies = [ "bytes", "getrandom 0.3.2", @@ -2124,13 +2124,12 @@ checksum = "74765f6d916ee2faa39bc8e68e4f3ed8949b48cccdac59983d287a7cb71ce9c5" [[package]] name = "rand" -version = "0.9.0" +version = "0.9.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3779b94aeb87e8bd4e834cee3650289ee9e0d5677f976ecdb6d219e5f4f6cd94" +checksum = "9fbfd9d094a40bf3ae768db9361049ace4c0e04a4fd6b359518bd7b73a73dd97" dependencies = [ "rand_chacha", "rand_core", - "zerocopy", ] [[package]] @@ -2279,7 +2278,7 @@ checksum = "a4689e6c2294d81e88dc6261c768b63bc4fcdb852be6d1352498b114f61383b7" dependencies = [ "cc", "cfg-if", - "getrandom 0.2.15", + "getrandom 0.2.16", "libc", "untrusted", "windows-sys 0.52.0", @@ -2287,13 +2286,13 @@ dependencies = [ [[package]] name = "rpassword" -version = "7.3.1" +version = "7.4.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "80472be3c897911d0137b2d2b9055faf6eeac5b14e324073d83bc17b191d7e3f" +checksum = "66d4c8b64f049c6721ec8ccec37ddfc3d641c4a7fca57e8f2a89de509c73df39" dependencies = [ "libc", "rtoolbox", - "windows-sys 0.48.0", + "windows-sys 0.59.0", ] [[package]] @@ -2313,12 +2312,12 @@ dependencies = [ [[package]] name = "rtoolbox" -version = "0.0.2" +version = "0.0.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c247d24e63230cdb56463ae328478bd5eac8b8faa8c69461a77e8e323afac90e" +checksum = "a7cc970b249fbe527d6e02e0a227762c9108b2f49d81094fe357ffc6d14d7f6f" dependencies = [ "libc", - "windows-sys 0.48.0", + "windows-sys 0.52.0", ] [[package]] @@ -2856,9 +2855,9 @@ dependencies = [ [[package]] name = "tokio-util" -version = "0.7.14" +version = "0.7.15" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6b9590b93e6fcc1739458317cccd391ad3955e2bde8913edf6f95f9e65a8f034" +checksum = "66a539a9ad6d5d281510d5bd368c973d636c02dbf8a67300bfb6b950696ad7df" dependencies = [ "bytes", "futures-core", diff --git a/Cargo.toml b/Cargo.toml index e07aa87c6..6a474c402 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -32,8 +32,8 @@ jsonwebtoken = "9.3" rust-argon2 = "2.1" futures = "0.3" path-clean = "1" -pest = "2.7" -pest_derive = "2.7" +pest = "2.8" +pest_derive = "2.8" enum-iterator = "2" openssl = { version = "*", features = ["vendored"] } #https://docs.rs/openssl/0.10.34/openssl/#vendored deunicode = "1.6.1" @@ -43,9 +43,9 @@ env_logger = "0.11" rustelebot = "0.3" bincode = { version = "2.0.1", features = ["std", "serde"] } rand = "0.9" -rpassword = "7.3" +rpassword = "7.4" flate2 = "1" -blake3 = "1.7" +blake3 = "1.8" bytes = "1.10" tokio-stream = { version = "0.1", features = ["sync"] } tokio = { version = "1.44", features = ["rt-multi-thread", "parking_lot", "fs"] } diff --git a/frontend/public/assets/favicon.ico b/frontend/public/assets/favicon.ico index 01eb505bd..8d861c4c8 100644 Binary files a/frontend/public/assets/favicon.ico and b/frontend/public/assets/favicon.ico differ diff --git a/src/api/api_utils.rs b/src/api/api_utils.rs index 3bd95de3d..cdad615ab 100644 --- a/src/api/api_utils.rs +++ b/src/api/api_utils.rs @@ -27,6 +27,7 @@ use crate::utils::network::request::{extract_extension_from_url, replace_url_ext use crate::utils::size_utils::human_readable_byte_size; use crate::utils::{debug_if_enabled, sys_utils, trace_if_enabled}; use crate::BUILD_TIMESTAMP; +use axum::body::Body; use axum::http::HeaderMap; use axum::response::IntoResponse; use chrono::{DateTime, Utc}; @@ -39,7 +40,6 @@ use std::collections::HashMap; use std::io::BufWriter; use std::path::Path; use std::sync::Arc; -use axum::body::Body; use tokio::sync::Mutex; use url::Url; @@ -181,7 +181,7 @@ fn get_stream_options(app_state: &AppState) -> StreamOptions { // content_length // } -fn get_stream_alternative_url(stream_url: &str, input: &ConfigInput, alias_input: &Arc) -> String { +pub fn get_stream_alternative_url(stream_url: &str, input: &ConfigInput, alias_input: &Arc) -> String { let Some(input_user_info) = input.get_user_info() else { return stream_url.to_owned() }; let Some(alt_input_user_info) = alias_input.get_user_info() else { return stream_url.to_owned() }; @@ -252,39 +252,35 @@ impl StreamDetails { /** * If successfully a provider connection is used, do not forget to release if unsuccessfully */ -async fn get_streaming_options(app_state: &AppState, stream_url: &str, input_opt: Option<&ConfigInput>, force_provider: Option<&str>) +async fn get_streaming_options(app_state: &AppState, stream_url: &str, input: &ConfigInput, force_provider: Option<&str>) -> (Option, StreamingOption, Option>) { - if let Some(input) = input_opt { - let provider_connection_guard = match force_provider { - Some(provider) => app_state.active_provider.force_exact_acquire_connection(provider).await, - None => app_state.active_provider.acquire_connection(&input.name).await - }; - let stream_response_params = match &*provider_connection_guard { - ProviderAllocation::Exhausted => { - let stream = create_provider_connections_exhausted_stream(&app_state.config, &[]); - StreamingOption::Custom(stream) - } - ProviderAllocation::Available(ref provider) - | ProviderAllocation::GracePeriod(ref provider) => { - // force_stream_provider means we keep the url and the provider. - // If force_stream_provider or the input is the same as the config we dont need to get new url - let (provider, url) = if force_provider.is_some() || provider.id == input.id { - (input.name.to_string(), stream_url.to_string()) - } else { - (provider.name.to_string(), get_stream_alternative_url(stream_url, input, provider)) - }; + let provider_connection_guard = match force_provider { + Some(provider) => app_state.active_provider.force_exact_acquire_connection(provider).await, + None => app_state.active_provider.acquire_connection(&input.name).await + }; + let stream_response_params = match &*provider_connection_guard { + ProviderAllocation::Exhausted => { + let stream = create_provider_connections_exhausted_stream(&app_state.config, &[]); + StreamingOption::Custom(stream) + } + ProviderAllocation::Available(ref provider) + | ProviderAllocation::GracePeriod(ref provider) => { + // force_stream_provider means we keep the url and the provider. + // If force_stream_provider or the input is the same as the config we dont need to get new url + let (provider, url) = if force_provider.is_some() || provider.id == input.id { + (input.name.to_string(), stream_url.to_string()) + } else { + (provider.name.to_string(), get_stream_alternative_url(stream_url, input, provider)) + }; - if matches!(&*provider_connection_guard, ProviderAllocation::Available(_)) { - StreamingOption::Available(Some(provider), url) - } else { - StreamingOption::GracePeriod(Some(provider), url) - } + if matches!(&*provider_connection_guard, ProviderAllocation::Available(_)) { + StreamingOption::Available(Some(provider), url) + } else { + StreamingOption::GracePeriod(Some(provider), url) } - }; - (Some(provider_connection_guard), stream_response_params, Some(input.headers.clone())) - } else { - (None, StreamingOption::Available(None, stream_url.to_string()), None) - } + } + }; + (Some(provider_connection_guard), stream_response_params, Some(input.headers.clone())) } @@ -296,13 +292,17 @@ fn get_grace_period_millis(connection_permission: &UserConnectionPermission, str } #[allow(clippy::too_many_arguments)] -async fn create_stream_response_details(app_state: &AppState, stream_options: &StreamOptions, stream_url: &str, - req_headers: &HeaderMap, input_opt: Option<&ConfigInput>, - item_type: PlaylistItemType, share_stream: bool, +async fn create_stream_response_details(app_state: &AppState, + stream_options: &StreamOptions, + stream_url: &str, + req_headers: &HeaderMap, + input: &ConfigInput, + item_type: PlaylistItemType, + share_stream: bool, connection_permission: UserConnectionPermission, force_provider: Option<&str>) -> StreamDetails { let (mut provider_connection_guard, stream_response_params, input_headers) = - get_streaming_options(app_state, stream_url, input_opt, force_provider).await; + get_streaming_options(app_state, stream_url, input, force_provider).await; let config_grace_period_millis = app_state.config.reverse_proxy.as_ref() .and_then(|r| r.stream.as_ref()).map_or_else(default_grace_period_millis, |s| s.grace_period_millis); let grace_period_millis = get_grace_period_millis(&connection_permission, &stream_response_params, config_grace_period_millis); @@ -372,7 +372,7 @@ where pub cluster: XtreamCluster, pub target_type: TargetType, pub target: &'a ConfigTarget, - pub input: Option<&'a ConfigInput>, + pub input: &'a ConfigInput, pub user: &'a ProxyUserCredentials, pub stream_ext: Option<&'a str>, pub req_context: XtreamApiStreamContext, @@ -412,30 +412,23 @@ where if redirect_request || is_dash_request { let redirect_url = if is_hls_request { &replace_url_extension(provider_url, HLS_EXT) } else { provider_url }; let redirect_url = if is_dash_request { &replace_url_extension(redirect_url, DASH_EXT) } else { redirect_url }; - if let Some(input) = params.input { - let redirect_url = get_redirect_alternative_url(app_state, redirect_url, input).await; - debug_if_enabled!("Redirecting stream request to {}", sanitize_sensitive_info(&redirect_url)); - return Some(redirect(&redirect_url).into_response()); - } - return Some(redirect(redirect_url.as_str()).into_response()); + let redirect_url = get_redirect_alternative_url(app_state, redirect_url, params.input).await; + debug_if_enabled!("Redirecting stream request to {}", sanitize_sensitive_info(&redirect_url)); + return Some(redirect(&redirect_url).into_response()); } } else if params.target_type == TargetType::Xtream { let Some(provider_id) = params.provider_id else { return Some(StatusCode::BAD_REQUEST.into_response()); }; - let Some(input) = params.input else { - return Some(StatusCode::BAD_REQUEST.into_response()); - }; - if redirect_request { // handle redirect for series but why ? if params.cluster == XtreamCluster::Series { let ext = params.stream_ext.unwrap_or_default(); - let url = input.url.as_str(); - let username = input.username.as_ref().map_or("", |v| v); - let password = input.password.as_ref().map_or("", |v| v); + let url = params.input.url.as_str(); + let username = params.input.username.as_ref().map_or("", |v| v); + let password = params.input.password.as_ref().map_or("", |v| v); // TODO do i need action_path like for timeshift ? let stream_url = format!("{url}/series/{username}/{password}/{provider_id}{ext}"); debug_if_enabled!("Redirecting stream request to {}", sanitize_sensitive_info(&stream_url)); @@ -444,14 +437,14 @@ where let target_name = params.target.name.as_str(); let virtual_id = params.item.get_virtual_id(); - let stream_url = match get_xtream_player_api_stream_url(input, ¶ms.req_context, ¶ms.get_query_path(provider_id, provider_url), provider_url) { + let stream_url = match get_xtream_player_api_stream_url(params.input, ¶ms.req_context, ¶ms.get_query_path(provider_id, provider_url), provider_url) { None => { error!("Cant find stream url for target {target_name}, context {}, stream_id {virtual_id}", params.req_context); return Some(StatusCode::BAD_REQUEST.into_response()); } Some(url) => { - match app_state.active_provider.get_next_provider(&input.name).await { - Some(provider_cfg) => get_stream_alternative_url(&url, input, &provider_cfg), + match app_state.active_provider.get_next_provider(¶ms.input.name).await { + Some(provider_cfg) => get_stream_alternative_url(&url, params.input, &provider_cfg), None => url, } } @@ -476,65 +469,70 @@ fn is_throttled_stream(item_type: PlaylistItemType, throttle_kbps: usize) -> boo throttle_kbps > 0 && matches!(item_type, PlaylistItemType::Video | PlaylistItemType::Series | PlaylistItemType::SeriesInfo) } -const SESSION_COOKIE: &str = "m3uflt_session="; +const SESSION_COOKIE_NAME: &str = "m3uflt_session="; + +fn create_delete_session_cookie() -> String { + format!("{SESSION_COOKIE_NAME}; Max-Age=0; Path=/; HttpOnly") +} + +fn create_session_cookie(cookie: &str) -> String { + // 3 hours should be enough + format!("{SESSION_COOKIE_NAME}{cookie}; Max-Age=10800; HttpOnly; SameSite=Strict") +} + +pub fn create_token_for_provider(secret: &[u8], virtual_id: u32, provider_name: &str, stream_url: &str) -> Option { + if let Ok(cookie_value) = obfuscate_text(secret, &format!("{virtual_id}:{provider_name}@{stream_url}")) { + return Some(cookie_value); + } + None +} + + +pub fn create_session_cookie_for_provider(secret: &[u8], virtual_id: u32, provider_name: &str, stream_url: &str) -> Option { + if let Some(cookie_value) = create_token_for_provider(secret, virtual_id, provider_name, stream_url) { + return Some(create_session_cookie(&cookie_value)); + } + None +} pub fn read_session_cookie(headers: &HeaderMap) -> Option { if let Some(cookie_header) = headers.get(axum::http::header::COOKIE) { let cookie_value = cookie_header.to_str().unwrap_or_default(); - if let Some(cookie) = cookie_value.split(';').map(str::trim).find(|&part| part.starts_with(SESSION_COOKIE)).and_then(|p| p.strip_prefix(SESSION_COOKIE)) { + if let Some(cookie) = cookie_value.split(';') + .map(str::trim) + .find(|&part| part.starts_with(SESSION_COOKIE_NAME)) + .and_then(|p| { + p.strip_prefix(SESSION_COOKIE_NAME).map(str::trim) + }) { return Some(cookie.to_string()); } } None } -fn get_session_cookie(cookie: &str) -> String { - // 3 hours should be enough - format!("{SESSION_COOKIE}{cookie}; Max-Age=10800; HttpOnly; Secure; SameSite=Strict") -} +pub fn get_stream_info_from_crypted_cookie(secret: &[u8], cookie: &str) -> Option<(u32, String, String)> { + if let Ok(decrypted) = deobfuscate_text(secret, cookie) { + let (virtual_id_and_provider, stream_url) = decrypted.split_once('@')?; + let (virtual_id, provider_name) = virtual_id_and_provider.split_once(':')?; -/// # Panics -pub async fn seek_stream_response(app_state: &AppState, - cookie: &str, - req_headers: &HeaderMap, - input: Option<&ConfigInput>, - item_type: PlaylistItemType, - _target: &ConfigTarget, - user: &ProxyUserCredentials) -> impl axum::response::IntoResponse + Send { - let stream_options = get_stream_options(app_state); - let share_stream = false; - let connection_permission = UserConnectionPermission::Allowed; + if virtual_id.is_empty() || provider_name.is_empty() || stream_url.is_empty() { + return None; + } - if let Ok(provider_and_stream_url) = deobfuscate_text(&app_state.config.t_encrypt_secret, cookie) { - let mut parts = provider_and_stream_url.splitn(2, '@'); - let provider_name = parts.next().unwrap_or(""); - let stream_url = parts.next().unwrap_or(""); - if !provider_name.is_empty() && !stream_url.is_empty() { - let mut stream_details = - create_stream_response_details(app_state, &stream_options, stream_url, req_headers, input, item_type, share_stream, connection_permission.clone(), Some(provider_name)).await; - - if stream_details.has_stream() { - let provider_response = stream_details.stream_info.as_ref().map(|(h, sc)| (h.clone(), *sc)); - let stream = ActiveClientStream::new(stream_details, app_state, user, connection_permission).await; - - let (status_code, header_map) = get_stream_response_with_headers(provider_response); - let mut response = axum::response::Response::builder() - .status(status_code); - for (key, value) in &header_map { - response = response.header(key, value); - } - - response = response.header(axum::http::header::SET_COOKIE, get_session_cookie(cookie)); - - let body_stream = prepare_body_stream(app_state, item_type, stream); - debug_if_enabled!("Streaming partial stream request from {}", sanitize_sensitive_info(stream_url)); - return response.body(body_stream).unwrap().into_response(); - } - drop(stream_details.provider_connection_guard.take()); + if let Ok(vid) = virtual_id.parse::() { + return Some(( + vid, + provider_name.to_string(), + stream_url.to_string(), + )); } } - error!("Cant open partial stream for cookie {cookie}"); - axum::http::StatusCode::BAD_REQUEST.into_response() + None +} + +pub fn get_stream_info_from_cookie(secret: &[u8], headers: &HeaderMap) -> Option<(u32, String, String)> { + let encrypted_cookie = read_session_cookie(headers)?; + get_stream_info_from_crypted_cookie(secret, &encrypted_cookie) } fn prepare_body_stream(app_state: &AppState, item_type: PlaylistItemType, stream: ActiveClientStream) -> Body { @@ -547,13 +545,63 @@ fn prepare_body_stream(app_state: &AppState, item_type: PlaylistItemType, stream body_stream } +pub fn bad_response_with_delete_cookie() -> impl IntoResponse + Send { + error!("Cant open provider forced stream, cookie invalid"); + if let Ok(delete_cookie_header) = axum::http::header::HeaderValue::from_str(&create_delete_session_cookie()) { + ([(axum::http::header::SET_COOKIE, delete_cookie_header)], StatusCode::BAD_REQUEST).into_response() + } else { + StatusCode::BAD_REQUEST.into_response() + } +} + +/// # Panics +pub async fn force_provider_stream_response(app_state: &AppState, + cookie: &str, + virtual_id: u32, + item_type: PlaylistItemType, + req_headers: &HeaderMap, + input: &ConfigInput, + user: &ProxyUserCredentials) -> impl axum::response::IntoResponse + Send { + if let Some((stream_virtual_id, provider_name, stream_url)) = get_stream_info_from_crypted_cookie(&app_state.config.t_encrypt_secret, cookie) { + if stream_virtual_id == virtual_id { + let stream_options = get_stream_options(app_state); + let share_stream = false; + let connection_permission = UserConnectionPermission::Allowed; + + let mut stream_details = + create_stream_response_details(app_state, &stream_options, &stream_url, req_headers, input, item_type, share_stream, connection_permission.clone(), Some(&provider_name)).await; + + if stream_details.has_stream() { + let provider_response = stream_details.stream_info.as_ref().map(|(h, sc)| (h.clone(), *sc)); + let stream = ActiveClientStream::new(stream_details, app_state, user, connection_permission).await; + + let (status_code, header_map) = get_stream_response_with_headers(provider_response); + let mut response = axum::response::Response::builder() + .status(status_code); + for (key, value) in &header_map { + response = response.header(key, value); + } + + response = response.header(axum::http::header::SET_COOKIE, create_session_cookie(cookie)); + + let body_stream = prepare_body_stream(app_state, item_type, stream); + debug_if_enabled!("Streaming provider forced stream request from {}", sanitize_sensitive_info(&stream_url)); + return response.body(body_stream).unwrap().into_response(); + } + drop(stream_details.provider_connection_guard.take()); + } + } + bad_response_with_delete_cookie().into_response() +} + /// # Panics #[allow(clippy::too_many_arguments)] pub async fn stream_response(app_state: &AppState, + virtual_id: u32, + item_type: PlaylistItemType, stream_url: &str, req_headers: &HeaderMap, - input: Option<&ConfigInput>, - item_type: PlaylistItemType, + input: &ConfigInput, target: &ConfigTarget, user: &ProxyUserCredentials, connection_permission: UserConnectionPermission) -> impl axum::response::IntoResponse + Send { @@ -580,6 +628,7 @@ pub async fn stream_response(app_state: &AppState, let stream = ActiveClientStream::new(stream_details, app_state, user, connection_permission).await; let stream_resp = if share_stream { + debug_if_enabled!("Streaming shared stream request from {}", sanitize_sensitive_info(stream_url)); // Shared Stream response let shared_headers = provider_response.as_ref().map_or_else(Vec::new, |(h, _)| h.clone()); SharedStreamManager::subscribe(app_state, stream_url, stream, shared_headers, stream_options.buffer_size).await; @@ -595,16 +644,16 @@ pub async fn stream_response(app_state: &AppState, axum::http::StatusCode::BAD_REQUEST.into_response() } } else { + debug_if_enabled!("Streaming stream request from {}", sanitize_sensitive_info(stream_url)); let (status_code, header_map) = get_stream_response_with_headers(provider_response); - let mut response = axum::response::Response::builder() - .status(status_code); + let mut response = axum::response::Response::builder().status(status_code); for (key, value) in &header_map { response = response.header(key, value); } if let Some(provider) = provider_name { - if let Ok(cookie_value) = obfuscate_text(&app_state.config.t_encrypt_secret, &format!("{provider}@{stream_url}")) { - response = response.header(axum::http::header::SET_COOKIE, get_session_cookie(&cookie_value)); + if let Some(cookie_value) = create_session_cookie_for_provider(&app_state.config.t_encrypt_secret, virtual_id, &provider, stream_url) { + response = response.header(axum::http::header::SET_COOKIE, &cookie_value); } } @@ -629,7 +678,7 @@ fn get_stream_throttle(app_state: &AppState) -> u64 { async fn shared_stream_response(app_state: &AppState, stream_url: &str, user: &ProxyUserCredentials, connect_permission: UserConnectionPermission) -> Option { if let Some(stream) = SharedStreamManager::subscribe_shared_stream(app_state, stream_url).await { - debug_if_enabled!("Using shared channel {}", sanitize_sensitive_info(stream_url)); + debug_if_enabled!("Using shared stream {}", sanitize_sensitive_info(stream_url)); if let Some(headers) = app_state.shared_stream_manager.get_shared_state_headers(stream_url).await { let (status_code, header_map) = get_stream_response_with_headers(Some((headers.clone(), StatusCode::OK))); let stream_details = StreamDetails::from_stream(stream); @@ -771,15 +820,65 @@ pub fn redirect(url: &str) -> impl IntoResponse { .unwrap() } -pub fn is_seek_response(req_headers: &HeaderMap) -> Option { - if let Some(cookie) = read_session_cookie(req_headers) { - if let Some(maybe_value) = req_headers.get("range") { - if let Ok(value) = maybe_value.to_str() { - if !value.starts_with("bytes=0-") { - return Some(cookie); - } - } +pub fn is_seek_response( + cluster: XtreamCluster, + virtual_id: u32, + secret: &[u8], + req_headers: &HeaderMap, +) -> Option { + // seek only for non-live streams + if cluster == XtreamCluster::Live { + return None; + } + + let cookie = read_session_cookie(req_headers)?; + match get_stream_info_from_crypted_cookie(secret, &cookie) { + Some((vid, _, _)) if vid == virtual_id => {} + _ => return None, + } + + // seek requests contain range header + let range = req_headers.get("range")?.to_str().ok()?; + if range.starts_with("bytes=0-") { + return None; + } + + read_session_cookie(req_headers) +} + +pub async fn check_force_provider(app_state: &AppState, virtual_id: u32, req_headers: &HeaderMap, user: &ProxyUserCredentials) -> (Option, UserConnectionPermission) { + + // if you have multi provider setup you need to delegate the same hls requests + // to the same provider. Hls has alternating m3u8 and stream requests. + let mut provider_name = None; + if let Some((stream_virtual_id, stream_provider_name, _stream_url)) = get_stream_info_from_cookie(&app_state.config.t_encrypt_secret, req_headers) { + if stream_virtual_id == virtual_id { + provider_name = Some(stream_provider_name); } } - None + + let connection_permission = match provider_name { + None => user.connection_permission(app_state).await, + Some(_) => UserConnectionPermission::Allowed + }; + + (provider_name, connection_permission) } + + +#[cfg(test)] +mod tests { + use crate::api::api_utils::SESSION_COOKIE_NAME; + + #[test] + fn test_cookie() { + let cookie_value = format!("{SESSION_COOKIE_NAME}bbblegum; sehe=dfd; sdfsdf=sdfsd;"); + let cookie = cookie_value.split(';') + .map(str::trim) + .find(|&part| part.starts_with(SESSION_COOKIE_NAME)) + .and_then(|p| { + p.strip_prefix(SESSION_COOKIE_NAME).map(str::trim) + }).expect("Cookie not found"); + assert_eq!(cookie, "bbblegum"); + } +} \ No newline at end of file diff --git a/src/api/endpoints/hdhomerun_api.rs b/src/api/endpoints/hdhomerun_api.rs index 7eccde995..9b3eabd4a 100644 --- a/src/api/endpoints/hdhomerun_api.rs +++ b/src/api/endpoints/hdhomerun_api.rs @@ -1,18 +1,19 @@ -use std::sync::Arc; -use axum::response::IntoResponse; -use crate::api::model::app_state::{HdHomerunAppState}; -use crate::model::api_proxy::{ProxyUserCredentials}; -use crate::model::config::{Config, TargetType}; +use crate::api::model::app_state::HdHomerunAppState; +use crate::auth::auth_basic::AuthBasic; +use crate::model::api_proxy::ProxyUserCredentials; +use crate::model::config::{Config, ConfigTarget, TargetType}; use crate::model::playlist::{M3uPlaylistItem, PlaylistItemType, XtreamCluster, XtreamPlaylistItem}; use crate::processing::parser::xtream::get_xtream_url; +use crate::repository::m3u_playlist_iterator::M3uPlaylistIterator; +use crate::repository::xtream_playlist_iterator::XtreamPlaylistIterator; use crate::utils::json_utils::get_string_from_serde_value; +use axum::response::IntoResponse; use bytes::Bytes; use futures::{stream, Stream, StreamExt}; use log::{error, warn}; use serde::{Deserialize, Serialize}; -use serde_json::{json}; -use crate::repository::m3u_playlist_iterator::M3uPlaylistIterator; -use crate::repository::xtream_playlist_iterator::{XtreamPlaylistIterator}; +use serde_json::json; +use std::sync::Arc; // https://info.hdhomerun.com/info/http_api @@ -136,7 +137,7 @@ where let lineup = Lineup { guide_number: item.epg_channel_id.unwrap_or(item.name).to_string(), guide_name: item.title.to_string(), - url: (if item.t_stream_url.is_empty() {&item.url} else {&item.t_stream_url}).to_string(), + url: (if item.t_stream_url.is_empty() { &item.url } else { &item.t_stream_url }).to_string(), }; match serde_json::to_string(&lineup) { Ok(content) => { @@ -206,14 +207,14 @@ async fn device_json(axum::extract::State(app_state): axum::extract::State>) -> impl IntoResponse { if let Some(device) = create_device(&app_state).await { - axum::Json(device).into_response() + axum::Json(device).into_response() } else { axum::http::StatusCode::INTERNAL_SERVER_ERROR.into_response() } } async fn lineup_status() -> impl IntoResponse { - axum::Json(json!({ + axum::Json(json!({ "ScanInProgress": 0, "ScanPossible": 0, "Source": "Cable", @@ -221,50 +222,75 @@ async fn lineup_status() -> impl IntoResponse { })) } +async fn lineup(app_state: &Arc, cfg: &Arc, credentials: &Arc, target: &ConfigTarget) -> impl IntoResponse { + let use_output = target.get_hdhomerun_output().as_ref().and_then(|o| o.use_output); + let use_all = use_output.is_none(); + let use_m3u = use_output.as_ref() == Some(&TargetType::M3u); + let use_xtream = use_output.as_ref() == Some(&TargetType::Xtream); + if (use_all || use_m3u) && target.has_output(&TargetType::M3u) { + let iterator = M3uPlaylistIterator::new(cfg, target, credentials).await.ok(); + let stream = m3u_item_to_lineup_stream(iterator); + let body_stream = stream::once(async { Ok(Bytes::from("[")) }) + .chain(stream) + .chain(stream::once(async { Ok(Bytes::from("]")) })); + return axum::response::Response::builder() + .status(axum::http::StatusCode::OK) + .header(axum::http::header::CONTENT_TYPE, mime::APPLICATION_JSON.to_string()) + .body(axum::body::Body::from_stream(body_stream)) + .unwrap().into_response(); + } else if (use_all || use_xtream) && target.has_output(&TargetType::Xtream) { + let server_info = app_state.app_state.config.get_user_server_info(credentials).await; + let base_url = server_info.get_base_url(); + + let base_url_live = if credentials.proxy.is_redirect(PlaylistItemType::Live) || target.is_force_redirect(PlaylistItemType::Live) { None } else { Some(base_url.clone()) }; + let base_url_vod = if credentials.proxy.is_redirect(PlaylistItemType::Video) || target.is_force_redirect(PlaylistItemType::Video) { None } else { Some(base_url) }; + + let live_channels = XtreamPlaylistIterator::new(XtreamCluster::Live, cfg, target, None, credentials).await.ok(); + let vod_channels = XtreamPlaylistIterator::new(XtreamCluster::Video, cfg, target, None, credentials).await.ok(); + // TODO include series when resolved + //let series_channels = xtream_repository::iter_raw_xtream_playlist(cfg, target, XtreamCluster::Series); + let live_stream = xtream_item_to_lineup_stream(Arc::clone(cfg), XtreamCluster::Live, Arc::clone(credentials), base_url_live.clone(), live_channels); + let vod_stream = xtream_item_to_lineup_stream(Arc::clone(cfg), XtreamCluster::Video, Arc::clone(credentials), base_url_vod.clone(), vod_channels); + let live_stream_peek = live_stream.peekable(); + let vod_stream_peek = vod_stream.peekable(); + // helper to decide if a comma is needed + let comma_stream = if live_stream_peek.size_hint().0 > 0 && vod_stream_peek.size_hint().0 > 0 { + stream::once(async { Ok(Bytes::from(",")) }).left_stream() + } else { + stream::empty().right_stream() + }; + let body_stream = stream::once(async { Ok(Bytes::from("[")) }) + .chain(live_stream_peek) + .chain(comma_stream) + .chain(vod_stream_peek) + .chain(stream::once(async { Ok(Bytes::from("]")) })); + return axum::response::Response::builder() + .status(axum::http::StatusCode::OK) + .header(axum::http::header::CONTENT_TYPE, mime::APPLICATION_JSON.to_string()) + .body(axum::body::Body::from_stream(body_stream)) + .unwrap() + .into_response(); + } + axum::http::StatusCode::NOT_FOUND.into_response() +} + +async fn auth_lineup_json(AuthBasic((username, password)): AuthBasic, axum::extract::State(app_state): axum::extract::State>) -> impl IntoResponse { + let cfg = Arc::clone(&app_state.app_state.config); + if let Some((credentials, target)) = cfg.get_target_for_username(&app_state.device.t_username).await { + if !username.eq(&credentials.username) || !password.eq(&credentials.password) { + return axum::http::StatusCode::UNAUTHORIZED.into_response(); + } + let user_credentials = Arc::new(credentials); + return lineup(&app_state, &cfg, &user_credentials, target).await.into_response(); + } + axum::http::StatusCode::NOT_FOUND.into_response() +} + async fn lineup_json(axum::extract::State(app_state): axum::extract::State>) -> impl IntoResponse { let cfg = Arc::clone(&app_state.app_state.config); if let Some((credentials, target)) = cfg.get_target_for_username(&app_state.device.t_username).await { - let use_output = target.get_hdhomerun_output().as_ref().and_then(|o| o.use_output); - let use_all = use_output.is_none(); - let use_m3u = use_output.as_ref() == Some(&TargetType::M3u); - let use_xtream = use_output.as_ref() == Some(&TargetType::Xtream); - if (use_all || use_m3u) && target.has_output(&TargetType::M3u) { - let iterator = M3uPlaylistIterator::new(&cfg,target,&credentials).await.ok(); - let stream = m3u_item_to_lineup_stream(iterator); - let body_stream = stream::once(async { Ok(Bytes::from("[")) }) - .chain(stream) - .chain(stream::once(async { Ok(Bytes::from("]")) })); - return axum::response::Response::builder() - .status(axum::http::StatusCode::OK) - .header(axum::http::header::CONTENT_TYPE, mime::APPLICATION_JSON.to_string()) - .body(axum::body::Body::from_stream(body_stream)) - .unwrap().into_response(); - } else if (use_all || use_xtream) && target.has_output(&TargetType::Xtream) { - let server_info = app_state.app_state.config.get_user_server_info(&credentials).await; - let base_url = server_info.get_base_url(); - - let base_url_live = if credentials.proxy.is_redirect(PlaylistItemType::Live) || target.is_force_redirect(PlaylistItemType::Live) { None } else { Some(base_url.clone()) }; - let base_url_vod = if credentials.proxy.is_redirect(PlaylistItemType::Video) || target.is_force_redirect(PlaylistItemType::Video) { None } else { Some(base_url) }; - - let live_channels = XtreamPlaylistIterator::new(XtreamCluster::Live, &cfg, target, None, &credentials).await.ok(); - let vod_channels = XtreamPlaylistIterator::new(XtreamCluster::Video, &cfg, target, None, &credentials).await.ok(); - // TODO include series when resolved - //let series_channels = xtream_repository::iter_raw_xtream_playlist(cfg, target, XtreamCluster::Series); - let user_credentials = Arc::new(credentials); - let live_stream = xtream_item_to_lineup_stream(Arc::clone(&cfg), XtreamCluster::Live, Arc::clone(&user_credentials), base_url_live.clone(), live_channels); - let vod_stream = xtream_item_to_lineup_stream(Arc::clone(&cfg), XtreamCluster::Video, Arc::clone(&user_credentials), base_url_vod.clone(), vod_channels); - let body_stream = stream::once(async { Ok(Bytes::from("[")) }) - .chain(live_stream) - .chain(stream::once(async { Ok(Bytes::from(",")) })) - .chain(vod_stream) - .chain(stream::once(async { Ok(Bytes::from("]")) })); - return axum::response::Response::builder() - .status(axum::http::StatusCode::OK) - .header(axum::http::header::CONTENT_TYPE, mime::APPLICATION_JSON.to_string()) - .body(axum::body::Body::from_stream(body_stream)) - .unwrap() - .into_response(); - } + let user_credentials = Arc::new(credentials); + return lineup(&app_state, &cfg, &user_credentials, target).await.into_response(); } axum::http::StatusCode::NOT_FOUND.into_response() } @@ -275,17 +301,17 @@ async fn auto_channel(axum::extract::State(_app_state): axum::extract::State axum::Router> { +pub fn hdhr_api_register(basic_auth: bool) -> axum::Router> { axum::Router::new() - .route("/device.xml", axum::routing::get(device_xml)) - .route("/device.json", axum::routing::get(device_json)) - .route("/discover.json", axum::routing::get(discover_json)) - .route("/lineup_status.json", axum::routing::get(lineup_status)) - .route("/lineup.json", axum::routing::get(lineup_json)) - // cfg.service(web::resource("/lineup.xml").route(web::get().to(lineup_xml))); - // cfg.service(web::resource("/lineup.m3u").route(web::get().to(lineup_m3u))); + .route("/device.xml", axum::routing::get(device_xml)) + .route("/device.json", axum::routing::get(device_json)) + .route("/discover.json", axum::routing::get(discover_json)) + .route("/lineup_status.json", axum::routing::get(lineup_status)) + .route("/lineup.json", if basic_auth { axum::routing::get(auth_lineup_json) } else { axum::routing::get(lineup_json) }) + // cfg.service(web::resource("/lineup.xml").route(web::get().to(lineup_xml))); + // cfg.service(web::resource("/lineup.m3u").route(web::get().to(lineup_m3u))); .route("/auto/{channel}", axum::routing::get(auto_channel)) - .route("/tuner{tuner_num}/{channel}", axum::routing::get(auto_channel)) + .route("/tuner{tuner_num}/{channel}", axum::routing::get(auto_channel)) } // fn start_hdhomerum_discovery_handler(ssdp_socket: Arc, server: String, location: String, cache_control: String, usn: String) { diff --git a/src/api/endpoints/hls_api.rs b/src/api/endpoints/hls_api.rs index 9fda6f11e..5ec3eba67 100644 --- a/src/api/endpoints/hls_api.rs +++ b/src/api/endpoints/hls_api.rs @@ -1,19 +1,19 @@ -use crate::api::api_utils::stream_response; -use crate::api::api_utils::try_option_bad_request; +use crate::api::api_utils::{bad_response_with_delete_cookie, check_force_provider, create_session_cookie_for_provider, force_provider_stream_response, get_stream_alternative_url}; +use crate::api::api_utils::{get_stream_info_from_crypted_cookie, try_option_bad_request}; +use crate::api::model::active_provider_manager::ProviderConfig; use crate::api::model::app_state::AppState; use crate::api::model::streams::provider_stream::{create_custom_video_stream_response, CustomVideoStreamType}; use crate::model::api_proxy::{ProxyUserCredentials, UserConnectionPermission}; -use crate::model::config::{ConfigInput}; +use crate::model::config::ConfigInput; use crate::model::playlist::{PlaylistItemType, XtreamCluster}; use crate::processing::parser::hls::{rewrite_hls, RewriteHlsProps}; +use crate::utils::constants::HLS_EXT; use crate::utils::network::request; use crate::utils::network::request::{is_hls_url, replace_url_extension, sanitize_sensitive_info}; use axum::response::IntoResponse; use log::{debug, error}; use serde::Deserialize; use std::sync::Arc; -use crate::utils::constants::HLS_EXT; -use crate::utils::crypto_utils; #[derive(Debug, Deserialize)] struct HlsApiPathParams { @@ -24,24 +24,50 @@ struct HlsApiPathParams { token: String, } -fn hls_response(hls_content: String) -> impl IntoResponse + Send { - axum::response::Response::builder() - .status(axum::http::StatusCode::OK) - .header(axum::http::header::CONTENT_TYPE, "application/x-mpegurl") - .body(hls_content) - .unwrap() - .into_response() +fn hls_response(hls_content: String, cookie: Option) -> impl IntoResponse + Send { + let mut builder = axum::response::Response::builder() + .status(axum::http::StatusCode::OK) + .header(axum::http::header::CONTENT_TYPE, "application/x-mpegurl"); + if let Some(cookie) = cookie { + builder = builder.header(axum::http::header::SET_COOKIE, cookie); + } + builder.body(hls_content) + .unwrap() + .into_response() } pub(in crate::api) async fn handle_hls_stream_request(app_state: &Arc, user: &ProxyUserCredentials, + provider_name: Option, hls_url: &str, virtual_id: u32, input: &ConfigInput) -> impl IntoResponse + Send { let url = replace_url_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 { + let create_stream_and_cookie = |provider_cfg: &Arc| { + let stream_url = get_stream_alternative_url(&url, input, provider_cfg); + let cookie = create_session_cookie_for_provider( + &app_state.config.t_encrypt_secret, + virtual_id, + &provider_cfg.name, + &stream_url, + ); + (stream_url, Some(provider_cfg.name.to_string()), cookie) + }; + + let (request_url, provider, cookie) = match provider_name { + None => match app_state.active_provider.get_next_provider(&input.name).await { + Some(provider_cfg) => create_stream_and_cookie(&provider_cfg), + None => (url, None, None), + }, + Some(provider) => match app_state.active_provider.force_exact_acquire_connection(&provider).await.get_provider_config() { + Some(provider_cfg) => create_stream_and_cookie(&provider_cfg), + None => (url, None, None), + }, + }; + + match request::download_text_content(Arc::clone(&app_state.http_client), input, &request_url, None).await { Ok((content, response_url)) => { let rewrite_hls_props = RewriteHlsProps { secret: &app_state.config.t_encrypt_secret, @@ -50,9 +76,10 @@ pub(in crate::api) async fn handle_hls_stream_request(app_state: &Arc, hls_url: response_url, virtual_id, input_id: input.id, + provider_name: provider.unwrap_or_default(), // this should not happen }; let hls_content = rewrite_hls(user, &rewrite_hls_props); - hls_response(hls_content).into_response() + hls_response(hls_content, cookie).into_response() } Err(err) => { error!("Failed to download m3u8 {}", sanitize_sensitive_info(err.to_string().as_str())); @@ -72,22 +99,37 @@ async fn hls_api_stream( if user.permission_denied(&app_state) { return axum::http::StatusCode::FORBIDDEN.into_response(); } - let connection_permission = user.connection_permission(&app_state).await; - if connection_permission == UserConnectionPermission::Exhausted { - return create_custom_video_stream_response(&app_state.config, &CustomVideoStreamType::UserConnectionsExhausted).into_response(); - } - - let Ok(hls_url) = crypto_utils::deobfuscate_text(&app_state.config.t_encrypt_secret, ¶ms.token) else { return axum::http::StatusCode::BAD_REQUEST.into_response(); }; let target_name = &target.name; let virtual_id = params.stream_id; let input = try_option_bad_request!(app_state.config.get_input_by_id(params.input_id), true, format!("Cant find input for target {target_name}, context {}, stream_id {virtual_id}", XtreamCluster::Live)); - if is_hls_url(&hls_url) { - return handle_hls_stream_request(&app_state, &user, &hls_url, virtual_id, input).await.into_response(); + let (_provider_name, connection_permission) = check_force_provider(&app_state, virtual_id, &req_headers, &user).await; + if connection_permission == UserConnectionPermission::Exhausted { + return create_custom_video_stream_response(&app_state.config, &CustomVideoStreamType::UserConnectionsExhausted).into_response(); } - stream_response(&app_state, &hls_url, &req_headers, Some(input), PlaylistItemType::LiveHls, target, &user, connection_permission).await.into_response() + let Some((stream_virtual_id, stream_provider_name, hls_url)) = get_stream_info_from_crypted_cookie(&app_state.config.t_encrypt_secret, ¶ms.token) + else { + return bad_response_with_delete_cookie().into_response(); + }; + + if stream_virtual_id != virtual_id { + return bad_response_with_delete_cookie().into_response(); + } + + let provider_name = Some(stream_provider_name); + + if is_hls_url(&hls_url) { + return handle_hls_stream_request(&app_state, &user, provider_name, &hls_url, virtual_id, input).await.into_response(); + } + + // if provider_name.is_some() { + // TODO we decode twice the cookie, one time to check for connection permission and one time in force_provider_stream_response + force_provider_stream_response(&app_state, ¶ms.token, virtual_id, PlaylistItemType::LiveHls, &req_headers, input, &user).await.into_response() + // } else { + // stream_response(&app_state, virtual_id, PlaylistItemType::LiveHls, &hls_url, &req_headers, input, target, &user, connection_permission).await.into_response() + // } } pub fn hls_api_register() -> axum::Router> { @@ -96,13 +138,3 @@ pub fn hls_api_register() -> axum::Router> { //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))); } - -#[cfg(test)] -mod tests { - - #[test] - fn test_hls_api_register() { - - } - -} \ No newline at end of file diff --git a/src/api/endpoints/m3u_api.rs b/src/api/endpoints/m3u_api.rs index ec6b8b55c..c5614603e 100644 --- a/src/api/endpoints/m3u_api.rs +++ b/src/api/endpoints/m3u_api.rs @@ -1,4 +1,4 @@ -use crate::api::api_utils::{get_user_target, get_user_target_by_credentials, is_seek_response, redirect, redirect_response, resource_response, seek_stream_response, separate_number_and_remainder, stream_response, try_option_bad_request, try_result_bad_request, RedirectParams}; +use crate::api::api_utils::{get_user_target, get_user_target_by_credentials, is_seek_response, redirect, redirect_response, resource_response, force_provider_stream_response, separate_number_and_remainder, stream_response, try_option_bad_request, try_result_bad_request, RedirectParams, check_force_provider}; use crate::api::endpoints::hls_api::handle_hls_stream_request; use crate::api::model::app_state::AppState; use crate::api::model::request::UserApiRequest; @@ -6,12 +6,13 @@ use crate::model::api_proxy::{UserConnectionPermission}; use crate::model::config::{TargetType}; use crate::model::playlist::{FieldGetAccessor, PlaylistEntry, PlaylistItemType, XtreamCluster}; use crate::repository::m3u_repository::{m3u_get_item_for_stream_id, m3u_load_rewrite_playlist}; -use crate::utils::network::request::{sanitize_sensitive_info}; +use crate::utils::network::request::{extract_extension_from_url, sanitize_sensitive_info}; use axum::response::IntoResponse; use bytes::Bytes; use futures::stream; use log::{debug, error}; use std::sync::Arc; +use axum::http::StatusCode; use crate::api::endpoints::xtream_api::XtreamApiStreamContext; use crate::api::model::streams::provider_stream::{create_custom_video_stream_response, CustomVideoStreamType}; use crate::repository::storage_const; @@ -66,45 +67,39 @@ async fn m3u_api_stream( axum::extract::Path((username, password, stream_id)): axum::extract::Path<(String, String, String)>, 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(&username, &password, &api_req, &app_state).await, false, format!("Could not find any user {username}")); + if user.permission_denied(&app_state) { + return StatusCode::FORBIDDEN.into_response(); + } + + let target_name = &target.name; + if !target.has_output(&TargetType::M3u) { + debug!("Target has no m3u output {target_name}"); + return StatusCode::BAD_REQUEST.into_response(); + } + 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).await - else { return axum::http::StatusCode::BAD_REQUEST.into_response() }; - if user.permission_denied(&app_state) { - return axum::http::StatusCode::FORBIDDEN.into_response(); - } + let pli = try_result_bad_request!(m3u_get_item_for_stream_id(virtual_id, &app_state.config, target).await, true, format!("Failed to read m3u 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}, stream_id {virtual_id}")); - if !target.has_output(&TargetType::M3u) { - return axum::http::StatusCode::BAD_REQUEST.into_response(); - } + let cluster = XtreamCluster::try_from(pli.item_type).unwrap_or(XtreamCluster::Live); - 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())); - return axum::http::StatusCode::BAD_REQUEST.into_response(); - } - }; - - let input = app_state.config.get_input_by_name(m3u_item.input_name.as_str()); - - - if let Some(cookie) = is_seek_response(&req_headers) { + if let Some(cookie) = is_seek_response(cluster, pli.virtual_id, &app_state.config.t_encrypt_secret, &req_headers) { // partial request means we are in reverse proxy mode, seek happened - return seek_stream_response(&app_state, &cookie, &req_headers, input, m3u_item.item_type, target, &user).await.into_response() + return force_provider_stream_response(&app_state, &cookie, pli.virtual_id, pli.item_type, &req_headers, input, &user).await.into_response() } - let connection_permission = user.connection_permission(&app_state).await; + let (provider_name, connection_permission) = check_force_provider(&app_state, virtual_id, &req_headers, &user).await; if connection_permission == UserConnectionPermission::Exhausted { return create_custom_video_stream_response(&app_state.config, &CustomVideoStreamType::UserConnectionsExhausted).into_response(); } - let cluster = XtreamCluster::try_from(m3u_item.item_type).unwrap_or(XtreamCluster::Live); let context = XtreamApiStreamContext::try_from(cluster).unwrap_or(XtreamApiStreamContext::Live); let redirect_params = RedirectParams { - item: &m3u_item, - provider_id: m3u_item.get_provider_id(), + item: &pli, + provider_id: pli.get_provider_id(), cluster, target_type: TargetType::Xtream, target, @@ -119,15 +114,16 @@ async fn m3u_api_stream( return response.into_response(); } - let is_hls_request = m3u_item.item_type == PlaylistItemType::LiveHls || stream_ext.as_deref() == Some(HLS_EXT); + let extension = stream_ext.unwrap_or_else( + || extract_extension_from_url(&pli.url).map_or_else(String::new, std::string::ToString::to_string)); + + let is_hls_request = pli.item_type == PlaylistItemType::LiveHls || pli.item_type == PlaylistItemType::LiveDash || extension == HLS_EXT; // 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.url, m3u_item.virtual_id, hls_input).await.into_response(); + return handle_hls_stream_request(&app_state, &user, provider_name, &pli.url, pli.virtual_id, input).await.into_response(); } - stream_response(&app_state, m3u_item.url.as_str(), &req_headers, input, m3u_item.item_type, target, &user, connection_permission).await.into_response() + stream_response(&app_state, pli.virtual_id, pli.item_type, pli.url.as_str(), &req_headers, input, target, &user, connection_permission).await.into_response() } async fn m3u_api_resource( diff --git a/src/api/endpoints/xtream_api.rs b/src/api/endpoints/xtream_api.rs index 4faabc3b0..db3735359 100644 --- a/src/api/endpoints/xtream_api.rs +++ b/src/api/endpoints/xtream_api.rs @@ -1,7 +1,7 @@ // https://github.com/tellytv/go.xtream-codes/blob/master/structs.go use crate::api::api_utils; -use crate::api::api_utils::{get_user_target, get_user_target_by_credentials, is_seek_response, redirect_response, resource_response, seek_stream_response, separate_number_and_remainder, serve_file, stream_response, RedirectParams}; +use crate::api::api_utils::{get_user_target, get_user_target_by_credentials, is_seek_response, redirect_response, resource_response, force_provider_stream_response, separate_number_and_remainder, serve_file, stream_response, RedirectParams, check_force_provider}; use crate::api::api_utils::{redirect, try_option_bad_request, try_result_bad_request}; use crate::api::endpoints::hls_api::handle_hls_stream_request; use crate::api::endpoints::xmltv_api::get_empty_epg_response; @@ -21,7 +21,6 @@ use crate::repository::playlist_repository::get_target_id_mapping; use crate::repository::storage::{get_target_storage_path, hex_encode}; use crate::repository::{storage_const, user_repository, xtream_repository}; use crate::utils::constants::HLS_EXT; -use crate::utils::debug_if_enabled; use crate::utils::hash_utils::generate_playlist_uuid; use crate::utils::json_utils; use crate::utils::json_utils::get_u32_from_serde_value; @@ -192,35 +191,34 @@ async fn xtream_player_api_stream( return StatusCode::BAD_REQUEST.into_response(); } - 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), 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)); + let cluster = pli.xtream_cluster; - if let Some(cookie) = is_seek_response(req_headers) { + if let Some(cookie) = is_seek_response(cluster, pli.virtual_id, &app_state.config.t_encrypt_secret, req_headers) { // partial request means we are in reverse proxy mode, seek happened - return seek_stream_response(app_state, &cookie, req_headers, Some(input), pli.item_type, target, &user).await.into_response() - + return force_provider_stream_response(app_state, &cookie, pli.virtual_id, pli.item_type, req_headers, input, &user).await.into_response() } - let connection_permission = user.connection_permission(app_state).await; + let (provider_name, connection_permission) = check_force_provider(app_state, virtual_id, req_headers, &user).await; if connection_permission == UserConnectionPermission::Exhausted { return create_custom_video_stream_response(&app_state.config, &CustomVideoStreamType::UserConnectionsExhausted).into_response(); } + let context = stream_req.context.clone(); let redirect_params = RedirectParams { item: &pli, provider_id: Some(mapping.provider_id), - cluster: pli.xtream_cluster, + cluster, target_type: TargetType::Xtream, target, - input: Some(input), + input, user: &user, stream_ext: stream_ext.as_deref(), - req_context: stream_req.context.clone(), + req_context: context, action_path: stream_req.action_path, }; if let Some(response) = redirect_response(app_state, &redirect_params).await { @@ -239,19 +237,15 @@ async fn xtream_player_api_stream( let stream_url = try_option_bad_request!(get_xtream_player_api_stream_url(input, &stream_req.context, &query_path, &pli.url), true, format!("Cant find stream url for target {target_name}, context {}, stream_id {virtual_id}", stream_req.context)); - - let is_hls_request = pli.item_type == PlaylistItemType::LiveHls || extension == HLS_EXT; - + let is_hls_request = pli.item_type == PlaylistItemType::LiveHls || pli.item_type == PlaylistItemType::LiveDash || extension == HLS_EXT; // Reverse proxy mode if is_hls_request { - return handle_hls_stream_request(app_state, &user, &stream_url, pli.virtual_id, input).await.into_response(); + return handle_hls_stream_request(app_state, &user, provider_name, &stream_url, pli.virtual_id, input).await.into_response(); } - debug_if_enabled!("Streaming stream request from {}", sanitize_sensitive_info(&stream_url)); - stream_response(app_state, &stream_url, req_headers, Some(input), pli.item_type, target, &user, connection_permission).await.into_response() + stream_response(app_state, pli.virtual_id, pli.item_type, &stream_url, req_headers, input, target, &user, connection_permission).await.into_response() } - async fn xtream_player_api_stream_with_token( req_headers: &HeaderMap, app_state: &Arc, @@ -287,9 +281,11 @@ async fn xtream_player_api_stream_with_token( ui_enabled: false, }; + // TODO how should we use fixed provider for hls in multi provider config? + // Reverse proxy mode if is_hls_request { - return handle_hls_stream_request(app_state, &user, &pli.url, pli.virtual_id, input).await.into_response(); + return handle_hls_stream_request(app_state, &user, None, &pli.url, pli.virtual_id, input).await.into_response(); } let extension = stream_ext.unwrap_or_else( @@ -307,7 +303,7 @@ async fn xtream_player_api_stream_with_token( stream_req.context)); trace_if_enabled!("Streaming stream request from {}", sanitize_sensitive_info(&stream_url)); - stream_response(app_state, &stream_url, req_headers, Some(input), pli.item_type, target, &user, UserConnectionPermission::Allowed).await.into_response() + stream_response(app_state, pli.virtual_id, pli.item_type, &stream_url, req_headers, input, target, &user, UserConnectionPermission::Allowed).await.into_response() } else { axum::http::StatusCode::BAD_REQUEST.into_response() } diff --git a/src/api/main_api.rs b/src/api/main_api.rs index cd629da5e..4d7b98299 100644 --- a/src/api/main_api.rs +++ b/src/api/main_api.rs @@ -177,13 +177,14 @@ fn start_hdhomerun(cfg: &Arc, app_state: &Arc, infos: &mut Vec let app_host = host.clone(); let port = device.port; let device_clone = Arc::new(device.clone()); + let basic_auth = hdhomerun.auth; infos.push(format!("HdHomeRun Server '{}' running: http://{host}:{port}", device.name)); tokio::spawn(async move { let router = axum::Router::>::new() .layer(create_cors_layer()) .layer(create_compression_layer()) // .layer(TraceLayer::new_for_http()) // `Logger::default()` - .merge(hdhr_api_register()); + .merge(hdhr_api_register(basic_auth)); let router: axum::Router<()> = router.with_state(Arc::new(HdHomerunAppState { app_state: Arc::clone(&app_data), diff --git a/src/api/model/active_provider_manager.rs b/src/api/model/active_provider_manager.rs index ed539c408..f77ecfee7 100644 --- a/src/api/model/active_provider_manager.rs +++ b/src/api/model/active_provider_manager.rs @@ -21,6 +21,15 @@ impl ProviderConnectionGuard { } } } + pub fn get_provider_config(&self) -> Option> { + match self.allocation { + ProviderAllocation::Exhausted => None, + ProviderAllocation::Available(ref cfg) | + ProviderAllocation::GracePeriod(ref cfg) => { + Some(Arc::clone(cfg)) + } + } + } } impl Deref for ProviderConnectionGuard { @@ -251,7 +260,7 @@ impl SingleProviderLineup { } fn get_next(&self) -> Option> { - self.provider.get_next(true) + self.provider.get_next(false) } fn acquire(&self) -> ProviderAllocation { @@ -564,8 +573,10 @@ impl ActiveProviderManager { ProviderLineup::Multi(multi) => { for group in &multi.providers { match group { - ProviderPriorityGroup::SingleProviderGroup(config) => { - return Some((lineup, config)); + ProviderPriorityGroup::SingleProviderGroup(single) => { + if single.name == name { + return Some((lineup, single)); + } } ProviderPriorityGroup::MultiProviderGroup(_, configs) => { for config in configs { @@ -582,9 +593,9 @@ impl ActiveProviderManager { None } - pub async fn force_exact_acquire_connection(&self, input_name: &str) -> ProviderConnectionGuard { + pub async fn force_exact_acquire_connection(&self, provider_name: &str) -> ProviderConnectionGuard { let providers = self.providers.read().await; - let allocation = match Self::get_provider_config(input_name, &providers) { + let allocation = match Self::get_provider_config(provider_name, &providers) { None => ProviderAllocation::Exhausted, // No Name matched, we don't have this provider Some((_lineup, config)) => config.force_allocate(), }; @@ -736,7 +747,7 @@ mod tests { id, name: name.to_string(), url: "http://example.com".to_string(), - epg: Default::default(), + epg: Option::default(), username: None, password: None, persist: None, @@ -747,7 +758,7 @@ mod tests { max_connections, priority, aliases: None, - headers: Default::default(), + headers: HashMap::default(), options: None, method: InputFetchMethod::default(), t_base_url: String::default(), diff --git a/src/api/model/streams/chunked_buffer.rs b/src/api/model/streams/chunked_buffer.rs index 52885b133..ce4e9bc87 100644 --- a/src/api/model/streams/chunked_buffer.rs +++ b/src/api/model/streams/chunked_buffer.rs @@ -59,7 +59,7 @@ mod tests { let mut index:usize = 0; while let Some(chunk) = chunked_buffer.next_chunk() { - for &byte in chunk.iter() { + for byte in chunk { let expected_value = buffer[index % buffer.len()]; assert_eq!(byte, expected_value, "Wrong value {byte} != {expected_value} at index {index} detected!"); index+=1; diff --git a/src/api/model/streams/readonly_ring_buffer.rs b/src/api/model/streams/readonly_ring_buffer.rs index 5c7356af6..12e222180 100644 --- a/src/api/model/streams/readonly_ring_buffer.rs +++ b/src/api/model/streams/readonly_ring_buffer.rs @@ -68,7 +68,7 @@ mod tests { let mut index:usize = 0; while let Some(chunk) = ring_buffer.next_chunk() { - for &byte in chunk.iter() { + for byte in chunk { let expected_value = buffer[index % buffer.len()]; assert_eq!(byte, expected_value, "Wrong value {byte} != {expected_value} at index {index} detected!"); index+=1; diff --git a/src/api/scheduler.rs b/src/api/scheduler.rs index bac5c69ae..e869ced2e 100644 --- a/src/api/scheduler.rs +++ b/src/api/scheduler.rs @@ -72,7 +72,7 @@ mod tests { } } Err(_) => {} - }; + } let duration = start.elapsed(); assert!(runs.load(Ordering::SeqCst) == 6, "Failed to run"); diff --git a/src/auth/auth_basic.rs b/src/auth/auth_basic.rs new file mode 100644 index 000000000..75605bbf0 --- /dev/null +++ b/src/auth/auth_basic.rs @@ -0,0 +1,59 @@ +use axum::extract::FromRequestParts; +use axum::http::request::Parts; +use axum::http::StatusCode; +use base64::Engine; +use base64::engine::general_purpose; + +pub type Rejection = (StatusCode, &'static str); +#[derive(Debug, PartialEq, Eq, Clone)] +pub struct AuthBasic(pub (String, String)); + +impl FromRequestParts for AuthBasic +where + B: Send + Sync, +{ + type Rejection = Rejection; + + async fn from_request_parts(req: &mut Parts, _: &B) -> Result { + Self::decode_request_parts(req) + } +} + +impl AuthBasic { + fn from_header(contents: (String, String)) -> Self { + Self(contents) + } + + fn decode_request_parts(req: &mut Parts) -> Result { + let authorization = req + .headers + .get(axum::http::header::AUTHORIZATION) + .ok_or((StatusCode::BAD_REQUEST, "Authorization header is missing"))? + .to_str() + .map_err(|_| (StatusCode::BAD_REQUEST, "Authorization header contains invalid characters"))?; + + let split = authorization.split_once(' '); + match split { + Some(("Basic", contents)) => { + let decoded = decode(contents)?; + Ok(Self::from_header(decoded)) + }, + _ => Err((StatusCode::BAD_REQUEST, "`Authorization` header must be a basic auth")), + } + } +} + +/// Decodes the two parts of basic auth using the colon +fn decode(input: &str) -> Result<(String, String), Rejection> { + // Decode from base64 into a string + let decoded = general_purpose::STANDARD.decode(input).map_err(|_| (StatusCode::BAD_REQUEST, "Authorization header contains invalid characters"))?; + let decoded = String::from_utf8(decoded).map_err(|_| (StatusCode::BAD_REQUEST, "Authorization header contains invalid characters"))?; + + + // Return depending on if password is present + if let Some((username, password)) = decoded.split_once(':') { + Ok((username.trim().to_string(), password.trim().to_string())) + } else { + Err((StatusCode::BAD_REQUEST, "Authorization header contains no password")) + } +} \ No newline at end of file diff --git a/src/foundation/filter.pest b/src/foundation/filter.pest index 1c0ab389b..40556b6d9 100644 --- a/src/foundation/filter.pest +++ b/src/foundation/filter.pest @@ -1,5 +1,5 @@ WHITESPACE = _{ " " | "\t" | "\r" | "\n"} -field = { ^"group" | ^"title" | ^"name" | ^"url" | ^"input"} +field = { ^"group" | ^"title" | ^"name" | ^"url" | ^"input" | ^"caption"} and = { ^"and" } or = { ^"or" } not = { ^"not" } diff --git a/src/foundation/filter.rs b/src/foundation/filter.rs index 13f5cf0df..3ed10713c 100644 --- a/src/foundation/filter.rs +++ b/src/foundation/filter.rs @@ -710,7 +710,7 @@ mod tests { match get_filter(flt, None) { Ok(filter) => { - let result = CONSTANTS.re_whitespace.replace_all(&flt, " "); + let result = CONSTANTS.re_whitespace.replace_all(flt, " "); assert_eq!(format!("{filter}"), result.trim()); } Err(e) => { diff --git a/src/model/config.rs b/src/model/config.rs index 87d7f7e84..945647de9 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -38,13 +38,13 @@ use crate::utils::sys_utils::exit; const DEFAULT_USER_AGENT: &str = "Mozilla/5.0 (AppleTV; U; CPU OS 14_2 like Mac OS X; en-us) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/14.0.1 Safari/605.1.15"; pub const MAPPER_ATTRIBUTE_FIELDS: &[&str] = &[ - "name", "title", "group", "id", "chno", "logo", + "name", "title", "caption", "group", "id", "chno", "logo", "logo_small", "parent_code", "audio_track", "time_shift", "rec", "url", "epg_channel_id", "epg_id" ]; -pub const AFFIX_FIELDS: &[&str] = &["name", "title", "group"]; -pub const COUNTER_FIELDS: &[&str] = &["name", "title", "chno"]; +pub const AFFIX_FIELDS: &[&str] = &["name", "title", "caption", "group"]; +pub const COUNTER_FIELDS: &[&str] = &["name", "title", "caption", "chno"]; const STREAM_QUEUE_SIZE: usize = 1024; // mpsc channel holding messages. with 8192byte chunks and 2Mbit/s approx 8MB @@ -166,7 +166,7 @@ impl Display for ProcessingOrder { } } -#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Sequence)] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Sequence, Eq, PartialEq)] pub enum ItemField { #[serde(rename = "group")] Group, @@ -180,6 +180,8 @@ pub enum ItemField { Input, #[serde(rename = "type")] Type, + #[serde(rename = "caption")] + Caption, } impl ItemField { @@ -189,6 +191,7 @@ impl ItemField { const URL: &'static str = "Url"; const INPUT: &'static str = "Input"; const TYPE: &'static str = "Type"; + const CAPTION: &'static str = "Caption"; } impl Display for ItemField { @@ -200,6 +203,7 @@ impl Display for ItemField { Self::Url => Self::URL, Self::Input => Self::INPUT, Self::Type => Self::TYPE, + Self::Caption => Self::CAPTION, }) } } diff --git a/src/model/hdhomerun_config.rs b/src/model/hdhomerun_config.rs index 23e9e1230..fc397b916 100644 --- a/src/model/hdhomerun_config.rs +++ b/src/model/hdhomerun_config.rs @@ -66,6 +66,8 @@ impl HdHomeRunDeviceConfig { pub struct HdHomeRunConfig { #[serde(default)] pub enabled: bool, + #[serde(default)] + pub auth: bool, pub devices: Vec, } diff --git a/src/model/mapping.rs b/src/model/mapping.rs index 491bc8fc8..9dce5c502 100644 --- a/src/model/mapping.rs +++ b/src/model/mapping.rs @@ -1,6 +1,6 @@ use enum_iterator::Sequence; use log::{debug, error, trace}; -use regex::{Regex}; +use regex::Regex; use std::borrow::Cow; use std::collections::HashMap; use std::fmt::Display; @@ -295,7 +295,7 @@ impl MappingValueProcessor<'_> { } } - fn apply_tags(&mut self, value: &str, captures: &HashMap<&str, &str>) -> Option { + fn apply_tags(&self, value: &str, captures: &HashMap<&str, &str>) -> Option { let mut new_value = String::from(value); let tag_captures = CONSTANTS.re_template_tag.captures_iter(value) .filter(|caps| caps.len() > 1) @@ -336,12 +336,12 @@ impl MappingValueProcessor<'_> { fn apply_suffix(&mut self, captures: &HashMap<&str, &str>) { let mapper = self.mapper; - let suffix = &mapper.suffix; + let suffixes = &mapper.suffix; - for (key, value) in suffix { + for (key, value) in suffixes { if let Some(suffix) = self.apply_tags(value, captures) { if let Some(old_value) = self.get_property(key) { - let new_value = format!("{}{}", &old_value, suffix); + let new_value = format!("{old_value}{suffix}"); self.set_property(key, &new_value); } } diff --git a/src/model/playlist.rs b/src/model/playlist.rs index 072c257c1..b2c51d8ea 100644 --- a/src/model/playlist.rs +++ b/src/model/playlist.rs @@ -259,6 +259,7 @@ macro_rules! generate_field_accessor_impl_for_playlist_item_header { $( stringify!($prop) => Some(self.$prop.clone()), )* + "caption" => Some(if self.title.is_empty() { self.name.clone() } else { self.title.clone() }), "epg_channel_id" | "epg_id" => self.epg_channel_id.clone(), _ => None, } @@ -274,6 +275,11 @@ macro_rules! generate_field_accessor_impl_for_playlist_item_header { true } )* + "caption" => { + self.title = val.clone(); + self.name = val; + true + } "epg_channel_id" | "epg_id" => { self.epg_channel_id = Some(value.to_owned()); true diff --git a/src/processing/parser/hls.rs b/src/processing/parser/hls.rs index 955816f81..801bef759 100644 --- a/src/processing/parser/hls.rs +++ b/src/processing/parser/hls.rs @@ -1,6 +1,6 @@ use crate::model::api_proxy::ProxyUserCredentials; -use crate::utils::crypto_utils::{obfuscate_text}; use std::str; +use crate::api::api_utils::{create_token_for_provider}; use crate::utils::constants::{CONSTANTS, HLS_PREFIX}; pub struct RewriteHlsProps<'a> { @@ -10,6 +10,7 @@ pub struct RewriteHlsProps<'a> { pub hls_url: String, pub virtual_id: u32, pub input_id: u16, + pub provider_name: String, } fn rewrite_hls_url(input: &str, replacement: &str) -> String { @@ -29,10 +30,11 @@ fn rewrite_hls_url(input: &str, replacement: &str) -> String { } fn rewrite_uri_attrib(line: &str, props: &RewriteHlsProps) -> String { - if let Some(caps) = CONSTANTS.re_memory_usage.captures(line) { + if let Some(caps) = CONSTANTS.re_hls_uri.captures(line) { let uri = &caps[1]; - if let Ok(encrypted_uri) = obfuscate_text(props.secret, &rewrite_hls_url(&props.hls_url, uri)) { - return CONSTANTS.re_hls_uri.replace(line, format!(r#"URI="{encrypted_uri}""#)).to_string(); + let target_url = &rewrite_hls_url(&props.hls_url, uri); + if let Some(token) = create_token_for_provider(props.secret, props.virtual_id, &props.provider_name, target_url) { + return CONSTANTS.re_hls_uri.replace(line, format!(r#"URI="{token}""#)).to_string(); } } line.to_string() @@ -43,14 +45,30 @@ pub fn rewrite_hls(user: &ProxyUserCredentials, props: &RewriteHlsProps) -> Stri let password = &user.password; let mut result = Vec::new(); for line in props.content.lines() { + // skip comments if line.starts_with('#') { - result.push(rewrite_uri_attrib(line, props)); - } else if let Ok(token) = if line.starts_with("http") { - obfuscate_text(props.secret, line) + let rewritten = rewrite_uri_attrib(line, props); + result.push(rewritten); + continue; + } + + // target url + let target_url = if line.starts_with("http") { + line.to_string() } else { - obfuscate_text(props.secret, &rewrite_hls_url(&props.hls_url, line)) - } { - result.push(format!("{}/{HLS_PREFIX}/{username}/{password}/{}/{}/{token}", props.base_url, props.input_id, props.virtual_id)); + rewrite_hls_url(&props.hls_url, line) + }; + if let Some(token) = create_token_for_provider(props.secret, props.virtual_id, &props.provider_name, &target_url) { + let url = format!( + "{}/{HLS_PREFIX}/{}/{}/{}/{}/{}", + props.base_url, + username, + password, + props.input_id, + props.virtual_id, + token + ); + result.push(url); } } result.join("\r\n") diff --git a/src/processing/parser/xtream.rs b/src/processing/parser/xtream.rs index f540dfd42..61512f88c 100644 --- a/src/processing/parser/xtream.rs +++ b/src/processing/parser/xtream.rs @@ -121,7 +121,7 @@ pub fn parse_xtream(input: &ConfigInput, match map_to_xtream_streams(xtream_cluster, streams) { Ok(mut xtream_streams) => { - let mut group_map: HashMap:: = + let mut group_map: HashMap = xtream_categories.into_iter().map(|category| (category.category_id.to_string(), category) ).collect(); diff --git a/src/processing/processor/epg.rs b/src/processing/processor/epg.rs index 6764e590c..d8206af81 100644 --- a/src/processing/processor/epg.rs +++ b/src/processing/processor/epg.rs @@ -37,7 +37,7 @@ impl EpgIdCache<'_> { fn normalize_and_store(&mut self, name: &str, epg_id: Option<&String>) { let normalized_name = self.normalize(name); - self.normalized.insert(normalized_name, epg_id.map(|v| v.to_string())); + self.normalized.insert(normalized_name, epg_id.map(std::string::ToString::to_string)); } fn normalize(&self, name: &str) -> String { @@ -107,7 +107,7 @@ fn assign_channel_epg(new_epg: &mut Vec, fp: &mut FetchedPlaylist, id_cache if not_processed { let normalized = id_cache.normalize(&chan.header.name); if let Some(epg_id) = id_cache.normalized.get(&normalized) { - chan.header.epg_channel_id = epg_id.clone(); + chan.header.epg_channel_id.clone_from(epg_id); } } } diff --git a/src/repository/bplustree.rs b/src/repository/bplustree.rs index 20465df24..f725b315f 100644 --- a/src/repository/bplustree.rs +++ b/src/repository/bplustree.rs @@ -425,7 +425,7 @@ where } } - const fn new_with_root(root: BPlusTreeNode::) -> Self { + const fn new_with_root(root: BPlusTreeNode) -> Self { let (inner_order, leaf_order) = calc_order::(); Self { root, @@ -822,22 +822,22 @@ mod tests { // Query the tree for i in 0u32..=test_size { let found = tree.query(&i); - assert!(found.is_some(), "{content} {} not found", i); + assert!(found.is_some(), "{content} {i} not found"); assert!(found.unwrap().eq(&Record { id: i, data: format!("{content} {i}"), - }), "{content} {} not found", i); + }), "{content} {i} not found"); } let mut tree_query: BPlusTreeQuery = BPlusTreeQuery::try_new(&filepath)?; for i in 0u32..=test_size { let found = tree_query.query(&i); - assert!(found.is_some(), "{content} {} not found", i); + assert!(found.is_some(), "{content} {i} not found"); let entry = found.unwrap(); assert!(entry.eq(&Record { id: i, data: format!("{content} {i}"), - }), "{content} {} not found", i); + }), "{content} {i} not found"); } let mut tree_update: BPlusTreeUpdate = BPlusTreeUpdate::try_new(&filepath)?; @@ -850,7 +850,7 @@ mod tests { }; tree_update.update(&i, new_record)?; } else { - assert!(false, "{content} {} not found", i); + assert!(false, "{content} {i} not found"); } } @@ -858,13 +858,13 @@ mod tests { for i in 0u32..=test_size { let found = tree_query.query(&i); - assert!(found.is_some(), "{content} {} not found", i); + assert!(found.is_some(), "{content} {i} not found"); let entry = found.unwrap(); let expected = Record { id: i, data: format!("{content} {}", i + 9000), }; - assert!(entry.eq(&expected), "Entry not equal {:?} != {:?}", entry, expected); + assert!(entry.eq(&expected), "Entry not equal {entry:?} != {expected:?}"); } Ok(()) @@ -890,7 +890,7 @@ mod tests { tree.traverse(|keys, values| { keys.iter().zip(values.iter()).for_each(|(k, v)| { - assert!(format!("{content} {}", k + 1).eq(&v.data), "Wrong entry") + assert!(format!("{content} {}", k + 1).eq(&v.data), "Wrong entry"); }); }); @@ -914,8 +914,8 @@ mod tests { let tree: BPlusTree = BPlusTree::load(&filepath)?; // Traverse the tree - for (key, value) in tree.iter() { - assert!(format!("Entry {}", key).eq(&value.data), "Wrong entry"); + for (key, value) in &tree { + assert!(format!("Entry {key}").eq(&value.data), "Wrong entry"); entry_set.remove(key); } assert!(entry_set.is_empty()); diff --git a/src/repository/indexed_document.rs b/src/repository/indexed_document.rs index c264f4ab1..1d62b98c2 100644 --- a/src/repository/indexed_document.rs +++ b/src/repository/indexed_document.rs @@ -544,7 +544,7 @@ mod tests { for i in 0u32..=500 { idw.write_doc(i, &Record { id: i, - data: format!("E {}", i), + data: format!("E {i}"), })?; } diff --git a/src/tools/directed_graph.rs b/src/tools/directed_graph.rs index 3a3a4eb74..f6cbd1ca5 100644 --- a/src/tools/directed_graph.rs +++ b/src/tools/directed_graph.rs @@ -196,6 +196,14 @@ where } } +impl Default for DirectedGraph +where + K: Eq + std::hash::Hash + Clone + Display + Debug, +{ + fn default() -> Self { + Self::new() + } +} #[cfg(test)] mod tests { @@ -203,8 +211,8 @@ mod tests { use std::collections::HashSet; fn are_vecs_equal(vec1: &Vec<&str>, vec2: Vec<&str>) -> bool { - let set1: HashSet = vec1.into_iter().map(|s| s.to_string()).collect(); - let set2: HashSet = vec2.into_iter().map(|s| s.to_string()).collect(); + let set1: HashSet = vec1.iter().map(|s| s.to_string()).collect(); + let set2: HashSet = vec2.iter().map(|s| s.to_string()).collect(); set1 == set2 } @@ -326,12 +334,3 @@ mod tests { } } - -impl Default for DirectedGraph -where - K: Eq + std::hash::Hash + Clone + Display + Debug, -{ - fn default() -> Self { - Self::new() - } -} diff --git a/src/utils/constants.rs b/src/utils/constants.rs index 1126690c1..4e44fd350 100644 --- a/src/utils/constants.rs +++ b/src/utils/constants.rs @@ -31,7 +31,8 @@ pub const FILENAME_TRIM_PATTERNS: &[char] = &['.', '-', '_']; pub const MEDIA_STREAM_HEADERS: &[&str] = &["accept", "content-type", "content-length", "connection", "accept-ranges", "content-range", "vary", "transfer-encoding", "access-control-allow-origin", - "access-control-allow-credentials", "icy-metadata"]; + "access-control-allow-credentials", "icy-metadata", "cache-control", "referer", "last-modified", + "etag", "expires"]; pub struct KodiStyle { pub year: Regex, @@ -69,7 +70,7 @@ pub static CONSTANTS: LazyLock = LazyLock::new(|| re_username: Regex::new(r"(username=)[^&]*").unwrap(), re_password: Regex::new(r"(password=)[^&]*").unwrap(), re_token: Regex::new(r"(token=)[^&]*").unwrap(), - re_stream_url: Regex::new(r"(.*://).*/(live|video|movie|series|m3u-stream|resource)/\w+/\w+").unwrap(), + re_stream_url: Regex::new(r"(?i)^(?Phttps?://)[^/]+/(?Plive|video|movie|series|m3u-stream|resource)/[^/]+/[^/]+/").unwrap(), re_url: Regex::new(r"(.*://).*?/(.*)").unwrap(), re_base_href: Regex::new(r#"(href|src)="/([^"]*)""#).unwrap(), re_env_var: Regex::new(r"\$\{env:(?P[a-zA-Z_][a-zA-Z0-9_]*)}").unwrap(), diff --git a/src/utils/crypto_utils.rs b/src/utils/crypto_utils.rs index b1d231437..de7d74774 100644 --- a/src/utils/crypto_utils.rs +++ b/src/utils/crypto_utils.rs @@ -69,7 +69,7 @@ mod tests { fn test_encrypt() { let secret: [u8; 16] = rand::rng().random(); // Random IV (AES-CBC 16 Bytes) let plain = "hello world"; - let encrypted = encrypt_text(&secret, &plain); + let encrypted = encrypt_text(&secret, plain); let decrypted = decrypt_text(&secret, &encrypted.unwrap()).unwrap(); assert_eq!(decrypted, plain); diff --git a/src/utils/file/config_reader.rs b/src/utils/file/config_reader.rs index a82d679ac..7d7fc5b20 100644 --- a/src/utils/file/config_reader.rs +++ b/src/utils/file/config_reader.rs @@ -104,15 +104,15 @@ pub fn read_mapping(mapping_file: &str, resolve_var: bool) -> Result = serde_yaml::from_reader(config_file_reader(file, resolve_var)); - match mapping { + return match mapping { Ok(mut result) => { handle_m3u_filter_error_result!(M3uFilterErrorKind::Info, result.prepare()); - return Ok(Some(result)); + Ok(Some(result)) } Err(err) => { - return Err(info_err!(err.to_string())); + Err(info_err!(err.to_string())) } - } + }; } warn!("cant read mapping file: {}", mapping_file.to_str().unwrap_or("?")); Ok(None) @@ -339,10 +339,10 @@ fn get_csv_file_path(file_uri: &String) -> Result { if file_uri.contains("://") { if let Ok(url) = file_uri.parse::() { if url.scheme() == "file" { - match url.to_file_path() { - Ok(path) => return Ok(path), - Err(()) => return Err(str_to_io_error(&format!("Could not open {file_uri}"))), - } + return match url.to_file_path() { + Ok(path) => Ok(path), + Err(()) => Err(str_to_io_error(&format!("Could not open {file_uri}"))), + }; } } Err(str_to_io_error(&format!("Only file:// is supported {file_uri}"))) diff --git a/src/utils/network/request.rs b/src/utils/network/request.rs index f4d9ff99e..7ecbb663e 100644 --- a/src/utils/network/request.rs +++ b/src/utils/network/request.rs @@ -180,9 +180,9 @@ pub fn get_local_file_content(file_path: &PathBuf) -> Result { if content.len() >= 2 && is_gzip(&content[0..2]) { let mut decoder = GzDecoder::new(&content[..]); let mut decode_buffer = String::new(); - match decoder.read_to_string(&mut decode_buffer) { - Ok(_) => return Ok(decode_buffer), - Err(err) => return Err(str_to_io_error(&format!("failed to decode gzip content {err}"))) + return match decoder.read_to_string(&mut decode_buffer) { + Ok(_) => Ok(decode_buffer), + Err(err) => Err(str_to_io_error(&format!("failed to decode gzip content {err}"))) }; } return Ok(String::from_utf8_lossy(&content).parse().unwrap()); @@ -504,8 +504,8 @@ mod tests { fn test_url_mask() { // Replace with "***" let query = "https://bubblegum.tv/live/username/password/2344"; - let masked = sanitize_sensitive_info(&query); - println!("{masked}") + let masked = sanitize_sensitive_info(query); + println!("{masked}"); } #[test] diff --git a/src/utils/step_measure.rs b/src/utils/step_measure.rs index aafc3e1a1..f2aac7295 100644 --- a/src/utils/step_measure.rs +++ b/src/utils/step_measure.rs @@ -9,11 +9,11 @@ fn format_duration(duration: Duration) -> String { let millis_rem = duration.subsec_millis(); if millis < 1_000 { - format!("{} ms", millis) + format!("{millis} ms") } else if secs < 60 { - format!("{}.{:03} s", secs, millis_rem) + format!("{secs}.{millis_rem:03} s") } else { - format!("{}:{:02}.{:03} min", mins, secs_rem, millis_rem) + format!("{mins}:{secs_rem:02}.{millis_rem:03} min") } }