diff --git a/Cargo.lock b/Cargo.lock index 06a2181d3..af0780d03 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1096,7 +1096,7 @@ dependencies = [ [[package]] name = "frontend" -version = "3.2.13" +version = "3.2.14" dependencies = [ "anyhow", "base64", @@ -3793,7 +3793,7 @@ dependencies = [ [[package]] name = "shared" -version = "3.2.13" +version = "3.2.14" dependencies = [ "base64", "bitflags 2.10.0", @@ -4356,7 +4356,7 @@ checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b" [[package]] name = "tuliprox" -version = "3.2.13" +version = "3.2.14" dependencies = [ "arc-swap", "async-compression", diff --git a/backend/Cargo.toml b/backend/Cargo.toml index 0f72bca30..9a606172d 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "tuliprox" -version = "3.2.13" +version = "3.2.14" edition = "2021" # See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html diff --git a/backend/src/api/api_utils.rs b/backend/src/api/api_utils.rs index 83606dcfd..34a4c4a8e 100644 --- a/backend/src/api/api_utils.rs +++ b/backend/src/api/api_utils.rs @@ -1,5 +1,5 @@ use crate::api::endpoints::xtream_api::{get_xtream_player_api_stream_url, ApiStreamContext}; -use crate::api::model::UserSession; +use crate::api::model::{UserSession}; use crate::api::model::{ create_channel_unavailable_stream, create_custom_video_stream_response, create_provider_connections_exhausted_stream, create_provider_stream, @@ -17,12 +17,12 @@ use crate::utils::request; use crate::utils::{debug_if_enabled, trace_if_enabled}; use crate::BUILD_TIMESTAMP; use arc_swap::ArcSwapOption; -use axum::http::HeaderMap; +use axum::http::{HeaderMap}; use axum::response::IntoResponse; use chrono::{DateTime, Utc}; use futures::{StreamExt, TryStreamExt}; use jsonwebtoken::{decode, Algorithm, DecodingKey, Validation}; -use log::{debug, error, log_enabled, trace, warn}; +use log::{debug, error, log_enabled, trace, warn}; use reqwest::header::RETRY_AFTER; use serde::Serialize; use shared::model::{Claims, InputFetchMethod, PlaylistEntry, PlaylistItemType, StreamChannel, TargetType, UserConnectionPermission, XtreamCluster}; @@ -373,16 +373,13 @@ async fn resolve_streaming_strategy( ) }; - let available = match allocation { - ProviderAllocation::Exhausted - | ProviderAllocation::Available(_) => true, - ProviderAllocation::GracePeriod(_) => false, - }; - - if available { - ProviderStreamState::Available(Some(provider), url) - } else { - ProviderStreamState::GracePeriod(Some(provider), url) + match allocation { + ProviderAllocation::Exhausted => { + let stream = create_provider_connections_exhausted_stream(&app_state.app_config, &[]); + ProviderStreamState::Custom(stream) + }, + ProviderAllocation::Available(_) => ProviderStreamState::Available(Some(provider), url), + ProviderAllocation::GracePeriod(_) => ProviderStreamState::GracePeriod(Some(provider), url), } } } @@ -418,7 +415,7 @@ fn get_grace_period_millis( } } -#[allow(clippy::too_many_arguments)] +#[allow(clippy::too_many_arguments, clippy::too_many_lines)] async fn create_stream_response_details( app_state: &Arc, stream_options: &StreamOptions, @@ -450,6 +447,22 @@ async fn create_stream_response_details( .as_ref() .and_then(|guard| guard.allocation.get_provider_name()); + if let ProviderStreamState::GracePeriod(ref provider_grace_check, ref stream_url) = streaming_strategy.provider_stream_state { + tokio::time::sleep(tokio::time::Duration::from_millis(grace_period_millis)).await; + if let Some(provider_name) = provider_grace_check.as_ref() { + if app_state.active_provider.is_over_limit(provider_name).await { + app_state.connection_manager.update_stream_detail(&fingerprint.addr, CustomVideoStreamType::ProviderConnectionsExhausted).await; + debug!("Provider connections exhausted after grace period for provider: {provider_name}"); + let stream = create_provider_connections_exhausted_stream(&app_state.app_config, &[]); + streaming_strategy.provider_stream_state = ProviderStreamState::Custom(stream); + let provider_handle = streaming_strategy.provider_handle.take(); + app_state.connection_manager.release_provider_handle(provider_handle).await; + } + } else { + streaming_strategy.provider_stream_state = ProviderStreamState::Available(provider_grace_check.clone(), stream_url.clone()); + } + } + match streaming_strategy.provider_stream_state { // custom stream means we display our own stream like connection exhausted, channel-unavailable... ProviderStreamState::Custom(provider_stream) => { @@ -837,8 +850,7 @@ pub async fn stream_response( share_stream, connection_permission, None, - ) - .await; + ).await; if stream_details.has_stream() { // let content_length = get_stream_content_length(provider_response.as_ref()); @@ -853,15 +865,20 @@ pub async fn stream_response( } else { None }; - stream_channel.shared = share_stream; + + let mut is_stream_shared = share_stream; + if let Some((_header, _status_code, _url, Some(_custom_video))) = stream_details.stream_info.as_ref() { + if stream_details.stream.is_some() { + is_stream_shared = false; + } + } + + stream_channel.shared = is_stream_shared; let stream = ActiveClientStream::new(stream_details, app_state, user, connection_permission, fingerprint, stream_channel, Some(session_token), req_headers) .await; - let stream_resp = if share_stream { - debug_if_enabled!( - "Streaming shared stream request from {}", - sanitize_sensitive_info(stream_url) - ); + let stream_resp = if is_stream_shared { + debug_if_enabled!("Streaming shared stream request from {}",sanitize_sensitive_info(stream_url)); // Shared Stream response let shared_headers = provider_response .as_ref() @@ -898,9 +915,7 @@ pub async fn stream_response( ); if log_enabled!(log::Level::Debug) { if session_url.eq(&stream_url) { - debug!( - "Streaming stream request from {}", - sanitize_sensitive_info(stream_url) + debug!("Streaming stream request from {}", sanitize_sensitive_info(stream_url) ); } else { debug!( @@ -977,9 +992,7 @@ async fn try_shared_stream_response_if_any( if let Some((stream, provider)) = SharedStreamManager::subscribe_shared_stream(app_state, stream_url, &fingerprint.addr).await { - debug_if_enabled!( - "Using shared stream {}", - 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 @@ -991,6 +1004,7 @@ async fn try_shared_stream_response_if_any( axum::http::StatusCode::OK, ))); let mut stream_details = StreamDetails::from_stream(stream); + stream_details.provider_name = provider; stream_channel.shared = true; let stream = diff --git a/backend/src/api/model/streams/shared_stream_manager.rs b/backend/src/api/model/streams/shared_stream_manager.rs index 514657836..db1d77698 100644 --- a/backend/src/api/model/streams/shared_stream_manager.rs +++ b/backend/src/api/model/streams/shared_stream_manager.rs @@ -210,7 +210,7 @@ impl SharedStreamState { result = broadcast_rx.recv() => { match result { Ok(data) => { - // Wenn der Client pausiert, einfach skippen oder warten + // If the client press pause, skip if client_tx_clone.is_closed() { continue; } diff --git a/frontend/Cargo.toml b/frontend/Cargo.toml index a1d1c1a3f..c206be646 100644 --- a/frontend/Cargo.toml +++ b/frontend/Cargo.toml @@ -1,10 +1,10 @@ [package] name = "frontend" -version = "3.2.13" +version = "3.2.14" edition = "2021" [dependencies] -shared = { version = "3.2.13", path = "../shared" } +shared = { version = "3.2.14", path = "../shared" } chrono = "0" yew = "0.21" yew-router = "0.18" diff --git a/frontend/src/app/components/source_editor/input_form.rs b/frontend/src/app/components/source_editor/input_form.rs index b4a190e1e..8b76bf660 100644 --- a/frontend/src/app/components/source_editor/input_form.rs +++ b/frontend/src/app/components/source_editor/input_form.rs @@ -174,40 +174,27 @@ pub fn ConfigInputView(props: &ConfigInputViewProps) -> Html { use_effect_with(config_input, move |cfg| { if let Some(input) = cfg { input_form_state.dispatch(ConfigInputFormAction::SetAll(input.as_ref().clone())); + input_options_state.dispatch(ConfigInputOptionsFormAction::SetAll( - input - .options - .as_ref() - .map_or_else(ConfigInputOptionsDto::default, |d| d.clone()), + input.options.as_ref().map_or_else(ConfigInputOptionsDto::default, |d| d.clone()), )); + staged_input_state.dispatch(StagedInputFormAction::SetAll( - input - .staged - .as_ref() - .map_or_else(StagedInputDto::default, |c| c.clone()), + input.staged.as_ref().map_or_else(StagedInputDto::default, |c| c.clone()), )); // Load headers headers_state.set(input.headers.clone()); // Load EPG sources - epg_sources_state.set( - input - .epg - .as_ref() - .and_then(|epg| epg.sources.clone()) - .unwrap_or_default() - ); + epg_sources_state.set(input.epg.as_ref().and_then(|epg| epg.sources.clone()).unwrap_or_default()); // Load aliases aliases_state.set(input.aliases.clone().unwrap_or_default()); } else { input_form_state.dispatch(ConfigInputFormAction::SetAll(ConfigInputDto::default())); - input_options_state.dispatch(ConfigInputOptionsFormAction::SetAll( - ConfigInputOptionsDto::default(), - )); - staged_input_state - .dispatch(StagedInputFormAction::SetAll(StagedInputDto::default())); + input_options_state.dispatch(ConfigInputOptionsFormAction::SetAll(ConfigInputOptionsDto::default())); + staged_input_state.dispatch(StagedInputFormAction::SetAll(StagedInputDto::default())); headers_state.set(HashMap::new()); epg_sources_state.set(Vec::new()); aliases_state.set(Vec::new()); diff --git a/shared/Cargo.toml b/shared/Cargo.toml index 034f45819..fa0a309e1 100644 --- a/shared/Cargo.toml +++ b/shared/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "shared" -version = "3.2.13" +version = "3.2.14" edition = "2021" [dependencies]