From e2601bcfdc8a90e04379470881f08d677773f432 Mon Sep 17 00:00:00 2001 From: euzu Date: Wed, 30 Apr 2025 09:55:08 +0200 Subject: [PATCH] User Session handling for fast frequent requests --- Cargo.lock | 32 ++-- Cargo.toml | 2 +- src/api/api_utils.rs | 197 ++++++----------------- src/api/endpoints/hls_api.rs | 125 ++++++++------ src/api/endpoints/m3u_api.rs | 38 +++-- src/api/endpoints/xtream_api.rs | 36 +++-- src/api/model/active_provider_manager.rs | 2 +- src/api/model/active_user_manager.rs | 176 +++++++++++++------- src/api/model/provider_config.rs | 24 ++- src/api/model/streams/provider_stream.rs | 4 +- src/processing/parser/hls.rs | 31 +++- src/utils/default_utils.rs | 2 +- src/utils/hash_utils.rs | 21 +++ src/utils/network/xtream.rs | 9 +- 14 files changed, 379 insertions(+), 320 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index b789a7172..0f12841a0 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -362,9 +362,9 @@ checksum = "a2698f953def977c68f935bb0dfa959375ad4638570e969e2f1e9f433cbf1af6" [[package]] name = "cc" -version = "1.2.19" +version = "1.2.20" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8e3a13707ac958681c13b39b458c073d0d9bc8a22cb1b2f4c8e55eb72c13f362" +checksum = "04da6a0d40b948dfc4fa8f5bbf402b0fc1a64a28dbf7d12ffd683550f2c1b63a" dependencies = [ "jobserver", "libc", @@ -385,9 +385,9 @@ checksum = "613afe47fcd5fac7ccf1db93babcb082c5994d996f20b8b159f2ad1658eb5724" [[package]] name = "chrono" -version = "0.4.40" +version = "0.4.41" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1a7964611d71df112cb1730f2ee67324fcf4d0fc6606acbbe9bfe06df124637c" +checksum = "c469d952047f47f91b68d1cba3f10d63c11d73e4636f24f08daf0278abf01c4d" dependencies = [ "android-tzdata", "iana-time-zone", @@ -660,9 +660,9 @@ dependencies = [ [[package]] name = "deunicode" -version = "1.6.1" +version = "1.6.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "dc55fe0d1f6c107595572ec8b107c0999bb1a2e0b75e37429a4fb0d6474a0e7d" +checksum = "abd57806937c9cc163efc8ea3910e00a62e2aeb0b8119f1793a978088f8f6b04" [[package]] name = "digest" @@ -2046,9 +2046,9 @@ dependencies = [ [[package]] name = "quick-xml" -version = "0.37.4" +version = "0.37.5" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a4ce8c88de324ff838700f36fb6ab86c96df0e3c4ab6ef3a9b2044465cce1369" +checksum = "331e97a1af0bf59823e6eadffe373d7b27f485be8748f71471c662c1f269b7fb" dependencies = [ "memchr", ] @@ -2632,9 +2632,9 @@ checksum = "13c2bddecc57b384dee18652358fb23172facb8a2c51ccc10d74c157bdea3292" [[package]] name = "syn" -version = "2.0.100" +version = "2.0.101" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b09a44accad81e1ba1cd74a32461ba89dee89095ba17b32f5d03683b1b1fc2a0" +checksum = "8ce2b7fc941b3a24138a0a7cf8e858bfc6a992e7978a068a5c760deb0ed43caf" dependencies = [ "proc-macro2", "quote", @@ -3241,9 +3241,9 @@ dependencies = [ [[package]] name = "webpki-roots" -version = "0.26.8" +version = "0.26.9" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2210b291f7ea53617fbafcc4939f10914214ec15aace5ba62293a668f322c5c9" +checksum = "29aad86cec885cafd03e8305fd727c418e970a521322c91688414d5b8efba16b" dependencies = [ "rustls-pki-types", ] @@ -3617,18 +3617,18 @@ dependencies = [ [[package]] name = "zerocopy" -version = "0.8.24" +version = "0.8.25" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2586fea28e186957ef732a5f8b3be2da217d65c5969d4b1e17f973ebbe876879" +checksum = "a1702d9583232ddb9174e01bb7c15a2ab8fb1bc6f227aa1233858c351a3ba0cb" dependencies = [ "zerocopy-derive", ] [[package]] name = "zerocopy-derive" -version = "0.8.24" +version = "0.8.25" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a996a8f63c5c4448cd959ac1bab0aaa3306ccfd060472f85943ee0750f0169be" +checksum = "28a6e20d751156648aa063f3800b706ee209a32c0b4d9f24be3d980b01be55ef" dependencies = [ "proc-macro2", "quote", diff --git a/Cargo.toml b/Cargo.toml index 6a474c402..964cf86a2 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -36,7 +36,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" +deunicode = "1.6.2" mime = "0.3" log = "0.4" env_logger = "0.11" diff --git a/src/api/api_utils.rs b/src/api/api_utils.rs index 1bc8a3214..4e903d147 100644 --- a/src/api/api_utils.rs +++ b/src/api/api_utils.rs @@ -19,13 +19,12 @@ use crate::model::playlist::{PlaylistEntry, PlaylistItemType, XtreamCluster}; use crate::tools::atomic_once_flag::AtomicOnceFlag; use crate::tools::lru_cache::LRUResourceCache; use crate::utils::constants::{DASH_EXT, HLS_EXT}; -use crate::utils::crypto_utils::{deobfuscate_text, obfuscate_text}; use crate::utils::default_utils::default_grace_period_millis; use crate::utils::file::file_utils::create_new_file_for_write; use crate::utils::network::request; use crate::utils::network::request::{extract_extension_from_url, replace_url_extension, sanitize_sensitive_info}; use crate::utils::size_utils::human_readable_byte_size; -use crate::utils::{debug_if_enabled, sys_utils, trace_if_enabled}; +use crate::utils::{debug_if_enabled, hash_utils, sys_utils, trace_if_enabled}; use crate::BUILD_TIMESTAMP; use axum::body::Body; use axum::http::HeaderMap; @@ -83,7 +82,9 @@ macro_rules! try_result_bad_request { pub use try_option_bad_request; pub use try_result_bad_request; +use crate::api::model::active_user_manager::UserSession; use crate::api::model::provider_config::ProviderConfig; +use crate::utils::hash_utils::base64_to_u32; pub fn get_server_time() -> String { chrono::offset::Local::now().with_timezone(&chrono::Local).format("%Y-%m-%d %H:%M:%S %Z").to_string() @@ -473,31 +474,16 @@ fn is_throttled_stream(item_type: PlaylistItemType, throttle_kbps: usize) -> boo const SESSION_COOKIE_NAME: &str = "m3uflt_session="; -fn create_delete_session_cookie() -> String { - format!("{SESSION_COOKIE_NAME}; Max-Age=0; Path=/; HttpOnly") +// fn create_delete_session_cookie() -> String { +// format!("{SESSION_COOKIE_NAME}; Max-Age=0; Path=/; HttpOnly") +// } + +pub fn create_session_cookie(token: u32) -> String { + let cookie = hash_utils::u32_to_base64(token); + format!("{SESSION_COOKIE_NAME}{cookie}; Path=/; HttpOnly; SameSite=Lax") } -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], token: &str, virtual_id: u32, provider_name: &str, stream_url: &str) -> Option { - if let Ok(cookie_value) = obfuscate_text(secret, &format!("{virtual_id}:{token}:{provider_name}@{stream_url}")) { - return Some(cookie_value); - } - None -} - - -pub fn create_session_cookie_for_provider(secret: &[u8], token: &str, virtual_id: u32, provider_name: &str, stream_url: &str) -> Option { - if let Some(cookie_value) = create_token_for_provider(secret, token, virtual_id, provider_name, stream_url) { - return Some(create_session_cookie(&cookie_value)); - } - None -} - -pub fn read_session_cookie(headers: &HeaderMap) -> Option { +pub fn read_session_token(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(';') @@ -506,40 +492,12 @@ pub fn read_session_cookie(headers: &HeaderMap) -> Option { .and_then(|p| { p.strip_prefix(SESSION_COOKIE_NAME).map(str::trim) }) { - return Some(cookie.to_string()); + return base64_to_u32(cookie); } } None } -/// # Panics -pub fn get_stream_info_from_crypted_cookie(secret: &[u8], cookie: &str) -> Option<(String, u32, String, String)> { - if let Ok(decrypted) = deobfuscate_text(secret, cookie) { - let (virtual_id_and_provider, stream_url) = decrypted.split_once('@')?; - let mut items: Vec = virtual_id_and_provider.split(':').filter(|s| !s.is_empty()).map(ToString::to_string).collect(); - if items.len() != 3 { - return None; - } - items.reverse(); - let virtual_id = items.pop().unwrap(); - let vid = virtual_id.parse::().ok()?; - let token = items.pop().unwrap(); - let provider_name = items.pop().unwrap(); - return Some(( - token, - vid, - provider_name.to_string(), - stream_url.to_string(), - )); - } - None -} - -pub fn get_stream_info_from_cookie(secret: &[u8], headers: &HeaderMap) -> Option<(String, 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 { let throttle_kbps = usize::try_from(get_stream_throttle(app_state)).unwrap_or_default(); let body_stream = if is_throttled_stream(item_type, throttle_kbps) { @@ -550,53 +508,38 @@ 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, + user_session: &UserSession, item_type: PlaylistItemType, req_headers: &HeaderMap, input: &ConfigInput, user: &ProxyUserCredentials) -> impl axum::response::IntoResponse + Send { - if let Some((stream_token, 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 && app_state.active_users.has_token(&user.username, &stream_token).await { - let stream_options = get_stream_options(app_state); - let share_stream = false; - let connection_permission = UserConnectionPermission::Allowed; + 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; + let mut stream_details = + create_stream_response_details(app_state, &stream_options, &user_session.stream_url, req_headers, input, item_type, share_stream, connection_permission.clone(), Some(&user_session.provider)).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; + 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()); + 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); } + + let body_stream = prepare_body_stream(app_state, item_type, stream); + debug_if_enabled!("Streaming provider forced stream request from {}", sanitize_sensitive_info(&user_session.stream_url)); + return response.body(body_stream).unwrap().into_response(); } - bad_response_with_delete_cookie().into_response() + drop(stream_details.provider_connection_guard.take()); + StatusCode::BAD_REQUEST.into_response() } /// # Panics @@ -631,7 +574,7 @@ pub async fn stream_response(app_state: &AppState, let provider_response = stream_details.stream_info.as_ref().map(|(h, sc)| (h.clone(), *sc)); let provider_name = stream_details.provider_connection_guard.as_ref().and_then(ProviderConnectionGuard::get_provider_name); - let stream = ActiveClientStream::new(stream_details, app_state, user, connection_permission).await; + let stream = ActiveClientStream::new(stream_details, app_state, user, connection_permission.clone()).await; let stream_resp = if share_stream { debug_if_enabled!("Streaming shared stream request from {}", sanitize_sensitive_info(stream_url)); // Shared Stream response @@ -658,10 +601,8 @@ pub async fn stream_response(app_state: &AppState, if let Some(provider) = provider_name { if matches!(item_type, PlaylistItemType::LiveHls | PlaylistItemType::LiveDash | PlaylistItemType::Video | PlaylistItemType::Series) { - if let Some(grace_token) = app_state.active_users.get_or_create_token(&user.username).await { - if let Some(cookie_value) = create_session_cookie_for_provider(&app_state.config.t_encrypt_secret, &grace_token, virtual_id, &provider, stream_url) { - response = response.header(axum::http::header::SET_COOKIE, &cookie_value); - } + if let Some(token) = app_state.active_users.create_user_session(&user.username, virtual_id, &provider, &stream_url, connection_permission).await { + response = response.header(axum::http::header::SET_COOKIE, &create_session_cookie(token)); } } } @@ -828,68 +769,30 @@ pub fn redirect(url: &str) -> impl IntoResponse { .unwrap() } -pub async fn is_seek_response( - app_state: &AppState, +pub async fn is_seek_request<'a>( cluster: XtreamCluster, - virtual_id: u32, req_headers: &HeaderMap, - username: &str, -) -> Option { +) -> bool { // seek only for non-live streams if cluster == XtreamCluster::Live { - return None; + return false; } - let cookie = read_session_cookie(req_headers)?; - match get_stream_info_from_crypted_cookie(&app_state.config.t_encrypt_secret, &cookie) { - Some((token, vid, _, _)) if vid == virtual_id => { - if !app_state.active_users.has_token(username, &token).await { - return None; - } + // seek requests contains range header + let range = req_headers + .get("range") + .and_then(|h| h.to_str().ok()) + .map(ToString::to_string); + + if let Some(range) = range { + if range.starts_with("bytes=0-") { + return false; } - _ => 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, item_type: PlaylistItemType, req_headers: &HeaderMap, user: &ProxyUserCredentials) -> (Option, UserConnectionPermission) { - if !matches!(item_type, PlaylistItemType::LiveHls | PlaylistItemType::LiveDash | PlaylistItemType::Series | PlaylistItemType::Video) { - return (None, user.connection_permission(app_state).await); - } - - // 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_token, 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 && app_state.active_users.has_token(&user.username, &stream_token).await { - provider_name = Some(stream_provider_name); + if range.starts_with("bytes=") { + return true; } } - - let connection_permission = if provider_name.is_some() { UserConnectionPermission::Allowed } else { - let permission = user.connection_permission(app_state).await; - match permission { - UserConnectionPermission::GracePeriod => { - if app_state.active_users.get_token(&user.username).await.is_some() { - UserConnectionPermission::Exhausted - } else { - UserConnectionPermission::GracePeriod - } - } - _ => permission, - } - }; - - (provider_name, connection_permission) -} - + false} #[cfg(test)] mod tests { diff --git a/src/api/endpoints/hls_api.rs b/src/api/endpoints/hls_api.rs index 1f07acded..bbaf670a3 100644 --- a/src/api/endpoints/hls_api.rs +++ b/src/api/endpoints/hls_api.rs @@ -1,11 +1,11 @@ -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::api_utils::{create_session_cookie, force_provider_stream_response, get_stream_alternative_url, is_seek_request, read_session_token}; +use crate::api::api_utils::{try_option_bad_request}; 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::playlist::{PlaylistItemType, XtreamCluster}; -use crate::processing::parser::hls::{rewrite_hls, RewriteHlsProps}; +use crate::processing::parser::hls::{get_hls_session_token_and_url_from_token, 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}; @@ -13,7 +13,7 @@ use axum::response::IntoResponse; use log::{debug, error}; use serde::Deserialize; use std::sync::Arc; -use crate::api::model::provider_config::ProviderConfig; +use crate::api::model::active_user_manager::UserSession; #[derive(Debug, Deserialize)] struct HlsApiPathParams { @@ -38,35 +38,34 @@ fn hls_response(hls_content: String, cookie: Option) -> impl IntoRespons pub(in crate::api) async fn handle_hls_stream_request(app_state: &Arc, user: &ProxyUserCredentials, - provider_name: Option, + user_session: Option<&UserSession>, hls_url: &str, virtual_id: u32, - input: &ConfigInput) -> impl IntoResponse + Send { + input: &ConfigInput, + connection_permission: UserConnectionPermission) -> impl IntoResponse + Send { let url = replace_url_extension(hls_url, HLS_EXT); let server_info = app_state.config.get_user_server_info(user).await; - let grace_token = app_state.active_users.get_or_create_token(&user.username).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, - &grace_token.clone().unwrap_or_default(), - 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), + let (request_url, session_token) = match user_session { + Some(session) => { + match app_state.active_provider.force_exact_acquire_connection(&session.provider).await.get_provider_config() { + Some(provider_cfg) => { + let stream_url = get_stream_alternative_url(&url, input, &provider_cfg); + (stream_url, Some(session.token)) + }, + None => (url, None), + } }, + None => { + match app_state.active_provider.get_next_provider(&input.name).await { + Some(provider_cfg) => { + let stream_url = get_stream_alternative_url(&url, input, &provider_cfg); + let session_token= app_state.active_users.create_user_session(&user.username, virtual_id, &provider_cfg.name, &stream_url, connection_permission).await; + (stream_url, session_token) + }, + None => (url, None), + } + } }; match request::download_text_content(Arc::clone(&app_state.http_client), input, &request_url, None).await { @@ -78,11 +77,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 - user_token: grace_token.unwrap_or_default().to_string(), + user_token: session_token.unwrap_or_default(), }; let hls_content = rewrite_hls(user, &rewrite_hls_props); - hls_response(hls_content, cookie).into_response() + hls_response(hls_content, session_token.map(|token| create_session_cookie(token))).into_response() } Err(err) => { error!("Failed to download m3u8 {}", sanitize_sensitive_info(err.to_string().as_str())); @@ -107,32 +105,53 @@ async fn hls_api_stream( 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)); - let (_provider_name, connection_permission) = check_force_provider(&app_state, virtual_id, PlaylistItemType::LiveHls, &req_headers, &user).await; - if connection_permission == UserConnectionPermission::Exhausted { - return create_custom_video_stream_response(&app_state.config, &CustomVideoStreamType::UserConnectionsExhausted).into_response(); - } - - let Some((stream_token, 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(); + let mut user_session = match read_session_token(&req_headers) { + None => None, + Some(token) => app_state.active_users.get_user_session(&user.username, token).await, }; - if stream_virtual_id != virtual_id || app_state.active_users.get_token(&user.username).await.is_some_and(|t| ! t.eq(&stream_token)) { - return bad_response_with_delete_cookie().into_response(); + if let Some(session) = &mut user_session { + if session.permission == UserConnectionPermission::Exhausted { + return create_custom_video_stream_response(&app_state.config, &CustomVideoStreamType::UserConnectionsExhausted).into_response(); + } + + if app_state.active_provider.is_over_limit(&session.provider).await { + return create_custom_video_stream_response(&app_state.config, &CustomVideoStreamType::ProviderConnectionsExhausted).into_response(); + } + + let hls_url = if let Some((session_token_opt, hls_url)) = get_hls_session_token_and_url_from_token(&app_state.config.t_encrypt_secret, ¶ms.token) { + if let Some(session_token) = session_token_opt { + if session.token != session_token { + return axum::http::StatusCode::BAD_REQUEST.into_response(); + } + } + hls_url + } else { + return axum::http::StatusCode::BAD_REQUEST.into_response(); + }; + session.stream_url = hls_url; + if session.virtual_id == virtual_id { + if is_seek_request(XtreamCluster::Live, &req_headers).await { + // partial request means we are in reverse proxy mode, seek happened + return force_provider_stream_response(&app_state, session, PlaylistItemType::LiveHls, &req_headers, input, &user).await.into_response() + } + } else { + return axum::http::StatusCode::BAD_REQUEST.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(); + } + + if is_hls_url(&session.stream_url) { + return handle_hls_stream_request(&app_state, &user, Some(&session), &session.stream_url, virtual_id, input, connection_permission).await.into_response(); + } + + force_provider_stream_response(&app_state, session, PlaylistItemType::LiveHls, &req_headers, input, &user).await.into_response() + } else { + axum::http::StatusCode::BAD_REQUEST.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> { diff --git a/src/api/endpoints/m3u_api.rs b/src/api/endpoints/m3u_api.rs index a87e89703..1a0937b7d 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, 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::api_utils::{get_user_target, get_user_target_by_credentials, is_seek_request, redirect, redirect_response, resource_response, force_provider_stream_response, separate_number_and_remainder, stream_response, try_option_bad_request, try_result_bad_request, RedirectParams, read_session_token}; use crate::api::endpoints::hls_api::handle_hls_stream_request; use crate::api::model::app_state::AppState; use crate::api::model::request::UserApiRequest; @@ -74,7 +74,7 @@ async fn m3u_api_stream( let target_name = &target.name; if !target.has_output(&TargetType::M3u) { - debug!("Target has no m3u output {target_name}"); + debug!("Target has no m3u playlist {target_name}"); return StatusCode::BAD_REQUEST.into_response(); } @@ -85,12 +85,26 @@ async fn m3u_api_stream( let cluster = XtreamCluster::try_from(pli.item_type).unwrap_or(XtreamCluster::Live); - if let Some(cookie) = is_seek_response(&app_state, cluster, pli.virtual_id, &req_headers, &user.username).await { - // partial request means we are in reverse proxy mode, seek happened - return force_provider_stream_response(&app_state, &cookie, pli.virtual_id, pli.item_type, &req_headers, input, &user).await.into_response() - } + let user_session = match read_session_token(&req_headers) { + None => None, + Some(token) => app_state.active_users.get_user_session(&user.username, token).await, + }; - let (provider_name, connection_permission) = check_force_provider(&app_state, virtual_id, pli.item_type, &req_headers, &user).await; + if let Some(session) = &user_session { + if session.permission == UserConnectionPermission::Exhausted { + return create_custom_video_stream_response(&app_state.config, &CustomVideoStreamType::UserConnectionsExhausted).into_response(); + } + + if app_state.active_provider.is_over_limit(&session.provider).await { + return create_custom_video_stream_response(&app_state.config, &CustomVideoStreamType::ProviderConnectionsExhausted).into_response(); + } + if session.virtual_id == virtual_id && is_seek_request(cluster, &req_headers).await { + // partial request means we are in reverse proxy mode, seek happened + return force_provider_stream_response(&app_state, session, pli.item_type, &req_headers, input, &user).await.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(); } @@ -120,7 +134,7 @@ async fn m3u_api_stream( 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, provider_name, &pli.url, pli.virtual_id, input).await.into_response(); + return handle_hls_stream_request(&app_state, &user, user_session.as_ref(), &pli.url, pli.virtual_id, input, 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() @@ -134,13 +148,15 @@ async fn m3u_api_resource( ) -> impl axum::response::IntoResponse + Send { let Ok(m3u_stream_id) = stream_id.parse::() else { return axum::http::StatusCode::BAD_REQUEST.into_response() }; 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() }; + else { return StatusCode::BAD_REQUEST.into_response() }; if user.permission_denied(&app_state) { - return axum::http::StatusCode::FORBIDDEN.into_response(); + return StatusCode::FORBIDDEN.into_response(); } + let target_name = &target.name; if !target.has_output(&TargetType::M3u) { - return axum::http::StatusCode::BAD_REQUEST.into_response(); + debug!("Target has no m3u playlist {target_name}"); + return StatusCode::BAD_REQUEST.into_response(); } let m3u_item = match m3u_get_item_for_stream_id(m3u_stream_id, &app_state.config, target).await { Ok(item) => item, diff --git a/src/api/endpoints/xtream_api.rs b/src/api/endpoints/xtream_api.rs index 40e806f24..a36f3bb95 100644 --- a/src/api/endpoints/xtream_api.rs +++ b/src/api/endpoints/xtream_api.rs @@ -1,13 +1,14 @@ // 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, force_provider_stream_response, separate_number_and_remainder, serve_file, stream_response, RedirectParams, check_force_provider}; +use crate::api::api_utils::{get_user_target, get_user_target_by_credentials, is_seek_request, redirect_response, resource_response, force_provider_stream_response, separate_number_and_remainder, serve_file, stream_response, RedirectParams, read_session_token}; 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; use crate::api::model::app_state::AppState; use crate::api::model::request::UserApiRequest; -use crate::api::model::streams::provider_stream::{create_custom_video_stream_response, CustomVideoStreamType}; +use crate::api::model::streams::provider_stream::{create_custom_video_stream_response, CustomVideoStreamType} +; use crate::api::model::xtream::XtreamAuthorizationResponse; use crate::m3u_filter_error::create_m3u_filter_error_result; use crate::m3u_filter_error::info_err; @@ -187,7 +188,7 @@ async fn xtream_player_api_stream( let target_name = &target.name; if !target.has_output(&TargetType::Xtream) { - debug!("Target has no xtream output {target_name}"); + debug!("Target has no xtream codes playlist {target_name}"); return StatusCode::BAD_REQUEST.into_response(); } @@ -197,12 +198,27 @@ async fn xtream_player_api_stream( 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(app_state, cluster, pli.virtual_id, req_headers, &user.username).await { - // partial request means we are in reverse proxy mode, seek happened - return force_provider_stream_response(app_state, &cookie, pli.virtual_id, pli.item_type, req_headers, input, &user).await.into_response() - } + let user_session = match read_session_token(req_headers) { + None => None, + Some(token) => app_state.active_users.get_user_session(&user.username, token).await, + }; - let (provider_name, connection_permission) = check_force_provider(app_state, virtual_id, pli.item_type, req_headers, &user).await; + if let Some(session) = &user_session { + if session.permission == UserConnectionPermission::Exhausted { + return create_custom_video_stream_response(&app_state.config, &CustomVideoStreamType::UserConnectionsExhausted).into_response(); + } + + if app_state.active_provider.is_over_limit(&session.provider).await { + return create_custom_video_stream_response(&app_state.config, &CustomVideoStreamType::ProviderConnectionsExhausted).into_response(); + } + + if session.virtual_id == virtual_id && is_seek_request(cluster, req_headers).await { + // partial request means we are in reverse proxy mode, seek happened + return force_provider_stream_response(app_state, session, pli.item_type, req_headers, input, &user).await.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(); } @@ -240,7 +256,7 @@ async fn xtream_player_api_stream( 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, provider_name, &stream_url, pli.virtual_id, input).await.into_response(); + return handle_hls_stream_request(app_state, &user, user_session.as_ref(), &stream_url, pli.virtual_id, input, 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() @@ -285,7 +301,7 @@ async fn xtream_player_api_stream_with_token( // Reverse proxy mode if is_hls_request { - return handle_hls_stream_request(app_state, &user, None, &pli.url, pli.virtual_id, input).await.into_response(); + return handle_hls_stream_request(app_state, &user, None, &pli.url, pli.virtual_id, input, UserConnectionPermission::Allowed).await.into_response(); } let extension = stream_ext.unwrap_or_else( diff --git a/src/api/model/active_provider_manager.rs b/src/api/model/active_provider_manager.rs index 332573b23..f8dc61213 100644 --- a/src/api/model/active_provider_manager.rs +++ b/src/api/model/active_provider_manager.rs @@ -563,7 +563,7 @@ impl ActiveProviderManager { pub async fn is_over_limit(&self, provider_name: &str) -> bool { let providers = self.providers.read().await; if let Some((_, config)) = Self::get_provider_config(provider_name, &providers) { - config.is_over_limit().await + config.is_over_limit(self.grace_period_timeout_secs).await } else { false } diff --git a/src/api/model/active_user_manager.rs b/src/api/model/active_user_manager.rs index 4b6cb1ee1..3a3beaaa9 100644 --- a/src/api/model/active_user_manager.rs +++ b/src/api/model/active_user_manager.rs @@ -1,11 +1,13 @@ use crate::model::api_proxy::UserConnectionPermission; +use crate::model::config::Config; +use crate::utils::default_utils::{default_grace_period_millis, default_grace_period_timeout_secs}; +use crate::utils::time_utils::current_time_secs; use jsonwebtoken::get_current_timestamp; use log::{debug, info}; +use rand::RngCore; use std::collections::HashMap; use std::sync::Arc; use tokio::sync::RwLock; -use crate::model::config::Config; -use crate::utils::default_utils::{default_grace_period_millis, default_grace_period_timeout_secs}; pub struct UserConnectionGuard { manager: Arc, @@ -21,22 +23,32 @@ impl Drop for UserConnectionGuard { } } +#[derive(Clone, Debug)] +pub struct UserSession { + pub token: u32, + pub virtual_id: u32, + pub provider: String, + pub stream_url: String, + pub ts: u64, + pub permission: UserConnectionPermission, +} + struct UserConnectionData { max_connections: u32, connections: u32, granted_grace: bool, grace_ts: u64, - token: Option, + sessions: Vec, } impl UserConnectionData { fn new(max_connections: u32) -> Self { Self { + max_connections, connections: 1, granted_grace: false, grace_ts: 0, - token: None, - max_connections, + sessions: Vec::new(), } } } @@ -79,6 +91,42 @@ impl ActiveUserManager { 0 } + async fn check_connection_permission(&self, username: &str, connection_data: &mut UserConnectionData) -> UserConnectionPermission { + let current_connections = connection_data.connections; + + if current_connections < connection_data.max_connections { + // Reset grace period because user is back under max_connections + connection_data.granted_grace = false; + connection_data.grace_ts = 0; + return UserConnectionPermission::Allowed; + } + + let now = get_current_timestamp(); + // Check if user already used grace period + if connection_data.granted_grace { + if now - connection_data.grace_ts <= self.grace_period_timeout_secs { + // Grace timeout still active, deny connection + debug!("User access denied, grace exhausted, too many connections: {username}"); + return UserConnectionPermission::Exhausted; + } + // Grace timeout expired, reset grace counters + connection_data.granted_grace = false; + connection_data.grace_ts = 0; + } + + if self.grace_period_millis > 0 && current_connections == connection_data.max_connections { + // Allow grace period once + connection_data.granted_grace = true; + connection_data.grace_ts = now; + debug!("Granted grace period for user access: {username}"); + return UserConnectionPermission::GracePeriod; + } + + // Too many connections, no grace allowed + debug!("User access denied, too many connections: {username}"); + UserConnectionPermission::Exhausted + } + pub async fn connection_permission( &self, username: &str, @@ -86,42 +134,9 @@ impl ActiveUserManager { ) -> UserConnectionPermission { if max_connections > 0 { if let Some(connection_data) = self.user.write().await.get_mut(username) { - let current_connections = connection_data.connections; - - if current_connections < max_connections { - // Reset grace period because user is back under max_connections - connection_data.granted_grace = false; - connection_data.grace_ts = 0; - return UserConnectionPermission::Allowed; - } - - let now = get_current_timestamp(); - // Check if user already used grace period - if connection_data.granted_grace { - if now - connection_data.grace_ts <= self.grace_period_timeout_secs { - // Grace timeout still active, deny connection - debug!("User access denied, grace exhausted, too many connections: {username}"); - return UserConnectionPermission::Exhausted; - } - // Grace timeout expired, reset grace counters - connection_data.granted_grace = false; - connection_data.grace_ts = 0; - } - - if self.grace_period_millis > 0 && current_connections == max_connections { - // Allow grace period once - connection_data.granted_grace = true; - connection_data.grace_ts = now; - debug!("Granted grace period for user access: {username}"); - return UserConnectionPermission::GracePeriod; - } - - // Too many connections, no grace allowed - debug!("User access denied, too many connections: {username}"); - return UserConnectionPermission::Exhausted; + return self.check_connection_permission(username, connection_data).await; } } - UserConnectionPermission::Allowed } @@ -158,50 +173,87 @@ impl ActiveUserManager { if connection_data.connections > 0 { connection_data.connections -= 1; } - // DO NOT reset granted_grace or grace_ts here! // We must preserve the grace period state until connection_permission() checks it. - if connection_data.connections == 0 { - lock.remove(username); - } else if connection_data.connections < connection_data.max_connections { - connection_data.token = None; - } + // if connection_data.connections == 0 { + // lock.remove(username); + // } else + // if connection_data.connections < connection_data.max_connections { + // connection_data.token = None; + // } } drop(lock); self.log_active_user().await; } - pub async fn get_or_create_token(&self, username: &str) -> Option { - let token = crate::utils::string_utils::generate_random_string(6); - let mut result = None; + fn find_user_session(token: u32, sessions: &[UserSession]) -> Option<&UserSession> { + for session in sessions { + if session.token == token { + return Some(session); + } + } + None + } + + pub async fn create_user_session(&self, username: &str, virtual_id: u32, provider: &str, stream_url: &str, connection_permission: UserConnectionPermission) -> Option { let mut lock = self.user.write().await; if let Some(connection_data) = lock.get_mut(username) { - result = if connection_data.token.is_some() { - connection_data.token.clone() - } else { - connection_data.token = Some(token.to_string()); - Some(token) + let session_token = rand::rng().next_u32(); + let session = UserSession { + token: session_token, + virtual_id, + provider: provider.to_string(), + stream_url: stream_url.to_string(), + ts: current_time_secs(), + permission: connection_permission, }; + connection_data.sessions.push(session); + return Some(session_token); } drop(lock); - result + None } - pub async fn get_token(&self, username: &str) -> Option { + pub async fn get_user_session(&self, username: &str, token: u32) -> Option { + self.update_user_session(username, token).await + // let mut lock = self.user.write().await; + // lock.get_mut(username) + // .and_then(|conn| Self::find_user_session(token, &conn.sessions)) + // .cloned() // owned copy + } + + async fn update_user_session(&self, username: &str, token: u32) -> Option { let mut lock = self.user.write().await; if let Some(connection_data) = lock.get_mut(username) { - connection_data.token.clone() - } else { - None + if connection_data.max_connections == 0 { + return Self::find_user_session(token, &connection_data.sessions).cloned(); + } + + // Separate mutable borrow of the session + let mut found_session_index = None; + for (i, session) in connection_data.sessions.iter().enumerate() { + if session.token == token { + found_session_index = Some(i); + break; + } + } + + if let Some(index) = found_session_index { + let session_permission = connection_data.sessions[index].permission.clone(); + if session_permission == UserConnectionPermission::GracePeriod { + let new_permission = self.check_connection_permission(username, connection_data).await; + connection_data.sessions[index].permission = new_permission; + } + return Some(connection_data.sessions[index].clone()); + } } - } - pub async fn has_token(&self, username: &str, token: &str) -> bool { - self.get_token(username).await.is_some_and(|t| token == t) + None } - async fn log_active_user(&self) { + + async fn log_active_user(&self) { if self.log_active_user { let user_count = self.active_users().await; let user_connection_count = self.active_connections().await; diff --git a/src/api/model/provider_config.rs b/src/api/model/provider_config.rs index 33700caf8..a40997e8f 100644 --- a/src/api/model/provider_config.rs +++ b/src/api/model/provider_config.rs @@ -1,9 +1,9 @@ use crate::api::model::active_provider_manager::ProviderAllocation; use crate::model::config::{ConfigInput, ConfigInputAlias, InputType, InputUserInfo}; -use std::ops::Deref; -use std::sync::Arc; use jsonwebtoken::get_current_timestamp; use log::debug; +use std::ops::Deref; +use std::sync::Arc; use tokio::sync::RwLock; #[derive(Debug)] @@ -83,12 +83,26 @@ impl ProviderConfig { } #[inline] - pub async fn is_over_limit(&self) -> bool { + pub async fn is_over_limit(&self, grace_period_timeout_secs: u64) -> bool { let max = self.max_connections; if max == 0 { return false; } - self.connection.read().await.current_connections > max + let mut guard = self.connection.write().await; + if guard.current_connections < self.max_connections { + guard.granted_grace = false; + guard.grace_ts = 0; + } + + if guard.granted_grace { + let now = get_current_timestamp(); + if now - guard.grace_ts <= grace_period_timeout_secs { + // Grace timeout still active, deny connection + debug!("Provider access denied, grace exhausted, too many connections: {}", self.name); + return true; + } + } + guard.current_connections > max } // @@ -132,7 +146,7 @@ impl ProviderConfig { guard.granted_grace = true; guard.grace_ts = now; guard.current_connections += 1; - return ProviderConfigAllocation::GracePeriod + return ProviderConfigAllocation::GracePeriod; } ProviderConfigAllocation::Exhausted } diff --git a/src/api/model/streams/provider_stream.rs b/src/api/model/streams/provider_stream.rs index d81adca9c..5287f35c6 100644 --- a/src/api/model/streams/provider_stream.rs +++ b/src/api/model/streams/provider_stream.rs @@ -21,7 +21,7 @@ use crate::api::model::stream::ProviderStreamResponse; pub enum CustomVideoStreamType { ChannelUnavailable, UserConnectionsExhausted, - // ProviderConnectionsExhausted, + ProviderConnectionsExhausted, } fn create_video_stream(video: Option<&Arc>>, headers: &[(String, String)], log_message: &str) -> ProviderStreamResponse { @@ -53,7 +53,7 @@ pub fn create_custom_video_stream_response(config: &Config, video_response: &Cus if let (Some(stream), Some((headers, status_code))) = match video_response { CustomVideoStreamType::ChannelUnavailable => create_channel_unavailable_stream(config, &[], StatusCode::BAD_REQUEST), CustomVideoStreamType::UserConnectionsExhausted => create_user_connections_exhausted_stream(config, &[]), - // CustomVideoStreamType::ProviderConnectionsExhausted => create_provider_connections_exhausted_stream(config, &[]), + CustomVideoStreamType::ProviderConnectionsExhausted => create_provider_connections_exhausted_stream(config, &[]), } { let mut builder = axum::response::Response::builder() .status(status_code); diff --git a/src/processing/parser/hls.rs b/src/processing/parser/hls.rs index d501d2f0e..9a2c479e8 100644 --- a/src/processing/parser/hls.rs +++ b/src/processing/parser/hls.rs @@ -1,17 +1,34 @@ use crate::model::api_proxy::ProxyUserCredentials; -use std::str; -use crate::api::api_utils::{create_token_for_provider}; use crate::utils::constants::{CONSTANTS, HLS_PREFIX}; +use crate::utils::crypto_utils::{deobfuscate_text, obfuscate_text}; +use crate::utils::hash_utils; +use std::str; + +fn create_hls_session_token_and_url(secret: &[u8], session_token: u32, stream_url: &str) -> Option { + let token = hash_utils::u32_to_base64(session_token); + if let Ok(cookie_value) = obfuscate_text(secret, &format!("{token}@{stream_url}")) { + return Some(cookie_value); + } + None +} + +pub fn get_hls_session_token_and_url_from_token(secret: &[u8], token: &str) -> Option<(Option, String)> { + if let Ok(decrypted) = deobfuscate_text(secret, token) { + let (session_token, stream_url) = decrypted.split_once('@')?; + return Some((hash_utils::base64_to_u32(session_token), stream_url.to_owned())); + } + None +} + pub struct RewriteHlsProps<'a> { - pub secret: &'a [u8;16], + pub secret: &'a [u8; 16], pub base_url: &'a str, pub content: &'a str, pub hls_url: String, pub virtual_id: u32, pub input_id: u16, - pub provider_name: String, - pub user_token: String, + pub user_token: u32, } fn rewrite_hls_url(input: &str, replacement: &str) -> String { @@ -34,7 +51,7 @@ fn rewrite_uri_attrib(line: &str, props: &RewriteHlsProps) -> String { if let Some(caps) = CONSTANTS.re_hls_uri.captures(line) { let uri = &caps[1]; let target_url = &rewrite_hls_url(&props.hls_url, uri); - if let Some(token) = create_token_for_provider(props.secret, &props.user_token, props.virtual_id, &props.provider_name, target_url) { + if let Some(token) = create_hls_session_token_and_url(props.secret, props.user_token, target_url) { return CONSTANTS.re_hls_uri.replace(line, format!(r#"URI="{token}""#)).to_string(); } } @@ -59,7 +76,7 @@ pub fn rewrite_hls(user: &ProxyUserCredentials, props: &RewriteHlsProps) -> Stri } else { rewrite_hls_url(&props.hls_url, line) }; - if let Some(token) = create_token_for_provider(props.secret, &props.user_token, props.virtual_id, &props.provider_name, &target_url) { + if let Some(token) = create_hls_session_token_and_url(props.secret, props.user_token, &target_url) { let url = format!( "{}/{HLS_PREFIX}/{}/{}/{}/{}/{}", props.base_url, diff --git a/src/utils/default_utils.rs b/src/utils/default_utils.rs index 3127c94a0..80063427b 100644 --- a/src/utils/default_utils.rs +++ b/src/utils/default_utils.rs @@ -9,5 +9,5 @@ pub const fn default_resolve_delay_secs() -> u16 { 2 } // Default grace values to accommodate rapid channel changes and seek requests, // helping avoid triggering hard max_connection enforcement. pub const fn default_grace_period_millis() -> u64 { 500 } -pub const fn default_grace_period_timeout_secs() -> u64 { 10 } +pub const fn default_grace_period_timeout_secs() -> u64 { 2 } pub const fn default_connect_timeout_secs() -> u32 { 6 } \ No newline at end of file diff --git a/src/utils/hash_utils.rs b/src/utils/hash_utils.rs index 7adf72410..4b660022e 100644 --- a/src/utils/hash_utils.rs +++ b/src/utils/hash_utils.rs @@ -1,3 +1,5 @@ +use base64::Engine; +use base64::engine::general_purpose; use crate::model::playlist::{PlaylistItemType, UUIDType}; use crate::repository::storage::hash_string; @@ -22,3 +24,22 @@ pub fn generate_playlist_uuid(key: &str, provider_id: &str, item_type: PlaylistI } hash_string(url) } + +pub fn u32_to_base64(value: u32) -> String { + // big-endian is safer and more portable when you care about consistent ordering or cross-platform data + let bytes = value.to_be_bytes(); + general_purpose::STANDARD.encode(&bytes) +} + +pub fn base64_to_u32(encoded: &str) -> Option { + let decoded = general_purpose::STANDARD.decode(encoded).ok()?; + + if decoded.len() != 4 { + return None; + } + + let arr: [u8; 4] = decoded + .as_slice() + .try_into().ok()?; + Some(u32::from_be_bytes(arr)) +} \ No newline at end of file diff --git a/src/utils/network/xtream.rs b/src/utils/network/xtream.rs index 6efde6811..4612750dc 100644 --- a/src/utils/network/xtream.rs +++ b/src/utils/network/xtream.rs @@ -137,12 +137,13 @@ pub async fn get_xtream_playlist(client: Arc, input: &ConfigInp let base_url = get_xtream_stream_url_base(&input.url, username, password); - if let Err(err) = request::get_input_json_content(Arc::clone(&client), input, base_url.as_str(), None).await { - warn!("Failed to login xtream account {username} {err}"); - return (Vec::with_capacity(0), vec![err]); + if let Err(_err) = request::get_input_json_content(Arc::clone(&client), input, base_url.as_str(), None).await { + if let Err(err) = request::get_input_json_content(Arc::clone(&client), input, &format!("{base_url}&action=get_account_info"), None).await { + warn!("Failed to login xtream account {username} {err}"); + return (Vec::with_capacity(0), vec![err]); + } } - let mut playlist_groups: Vec = Vec::with_capacity(128); let skip_cluster = get_skip_cluster(input);