From 9f3c2bca69f6577be4bf2e5b5798b4a1651f331e Mon Sep 17 00:00:00 2001 From: euzu Date: Thu, 24 Apr 2025 18:12:29 +0200 Subject: [PATCH] multi-provider hls stream fix --- CHANGELOG.md | 2 +- src/api/api_utils.rs | 325 +++++++++++++++-------- src/api/endpoints/hls_api.rs | 98 ++++--- src/api/endpoints/m3u_api.rs | 60 ++--- src/api/endpoints/xtream_api.rs | 36 ++- src/api/model/active_provider_manager.rs | 17 +- src/processing/parser/hls.rs | 36 ++- src/utils/constants.rs | 3 +- 8 files changed, 363 insertions(+), 214 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 6d91c9bc0..db97a98a0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -86,7 +86,7 @@ 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 http://user:password@ip:port/lineup.json +- hdhomerun supports now basic auth like you need to enable auth in config ```yaml hdhomerun: diff --git a/src/api/api_utils.rs b/src/api/api_utils.rs index 3bd95de3d..5f38fc55f 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,63 @@ 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, + req_headers: &HeaderMap, +) -> Option { + // seek only for non-live streams + if cluster == XtreamCluster::Live { + return None; + } + + // cookie needs to have virtual_id to identify that we have the same item + let cookie = read_session_cookie(req_headers)?; + if !cookie.starts_with(&format!("{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; + } + + Some(cookie) +} + +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;"); + 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) + }) {} + } +} \ No newline at end of file diff --git a/src/api/endpoints/hls_api.rs b/src/api/endpoints/hls_api.rs index 9fda6f11e..18caaa720 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::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..84f76fbbe 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, &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..6a3000e6c 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, 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/model/active_provider_manager.rs b/src/api/model/active_provider_manager.rs index ed539c408..7949f4d28 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<&Arc> { + match self.allocation { + ProviderAllocation::Exhausted => None, + ProviderAllocation::Available(ref cfg) | + ProviderAllocation::GracePeriod(ref cfg) => { + Some(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 { @@ -479,7 +488,7 @@ impl MultiProviderLineup { let allocation = { let config = Self::get_next_provider_from_group(priority_group, false); if config.is_none() { - Self::get_next_provider_from_group(priority_group, true) + Self::get_next_provider_from_group(priority_group, false) } else { config } @@ -582,9 +591,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(), }; diff --git a/src/processing/parser/hls.rs b/src/processing/parser/hls.rs index 955816f81..cf1db2918 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 { @@ -31,8 +32,9 @@ 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) { 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/utils/constants.rs b/src/utils/constants.rs index 8b17f8f2e..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,