Files
tuliprox/backend/src/api/api_utils.rs
T

1846 lines
71 KiB
Rust
Raw Normal View History

use crate::{
api::{
endpoints::xtream_api::{get_xtream_player_api_stream_url, ApiStreamContext},
model::{
create_active_client_stream, create_channel_unavailable_stream, create_custom_video_stream_response,
create_provider_connections_exhausted_stream, create_provider_stream, get_stream_response_with_headers,
tee_stream, AppState, CustomVideoStreamType, ProviderAllocation, ProviderConfig,
ProviderStreamFactoryOptions, ProviderStreamState, SharedStreamManager, StreamDetails, StreamError,
StreamingStrategy, ThrottledStream, UserApiRequest, UserSession,
},
},
auth::Fingerprint,
model::{ConfigInput, ConfigTarget, ProxyUserCredentials},
tools::lru_cache::LRUResourceCache,
utils::{
async_file_reader, async_file_writer, create_new_file_for_write, debug_if_enabled, get_file_extension, request,
request::{content_type_from_ext, parse_range, send_with_retry_and_provider},
trace_if_enabled,
},
BUILD_TIMESTAMP,
2026-02-14 17:52:35 +01:00
};
use arc_swap::ArcSwapOption;
use axum::{
body::Body,
http::{header, HeaderMap, HeaderValue, Response, StatusCode},
response::IntoResponse,
};
2026-01-20 10:09:58 +01:00
use bytes::Bytes;
2025-04-22 17:08:31 +02:00
use chrono::{DateTime, Utc};
use futures::{stream, Stream, StreamExt, TryStreamExt};
2025-03-20 11:57:32 +01:00
use jsonwebtoken::{decode, Algorithm, DecodingKey, Validation};
2026-01-02 13:55:19 +01:00
use log::{debug, error, info, log_enabled, trace, warn};
use serde::Serialize;
use shared::{
concat_string,
model::{
Claims, InputFetchMethod, PlaylistEntry, PlaylistItemType, ProxyType, StreamChannel, TargetType,
UserConnectionPermission, VirtualId, XtreamCluster,
},
utils::{
bin_serialize, extract_extension_from_url, human_readable_kbps, is_sanitize_sensitive_info_enabled,
replace_url_extension, sanitize_sensitive_info,
trim_slash, Internable, CONTENT_TYPE_CBOR, CONTENT_TYPE_JSON, DASH_EXT, HLS_EXT,
},
};
use std::{
borrow::Cow,
collections::HashMap,
convert::Infallible,
io::SeekFrom,
path::{Path, PathBuf},
sync::Arc,
};
use tokio::{
io::{AsyncReadExt, AsyncSeekExt},
sync::Mutex,
2025-07-24 14:27:57 +02:00
};
2025-12-16 16:23:15 +01:00
use tokio_util::io::ReaderStream;
2024-12-05 19:54:34 +01:00
use url::Url;
2026-03-10 13:49:52 +01:00
pub(crate) fn resolve_request_url_for_logging<'a>(input: &ConfigInput, stream_url: &'a str) -> Cow<'a, str> {
if is_sanitize_sensitive_info_enabled() {
return Cow::Borrowed(stream_url);
}
2026-03-10 13:49:52 +01:00
let provider = input.get_resolve_provider(stream_url);
if let Ok(url) = Url::parse(stream_url) {
return Cow::Owned(request::preview_request_target_for_logging(&url, provider.as_ref()));
}
input
.resolve_url(stream_url)
.ok()
.and_then(|resolved| {
Url::parse(resolved.as_ref())
.ok()
.map(|url| Cow::Owned(request::preview_request_target_for_logging(&url, provider.as_ref())))
})
.unwrap_or(Cow::Borrowed(stream_url))
}
2025-01-31 23:08:23 +01:00
#[macro_export]
macro_rules! try_option_bad_request {
($option:expr, $msg_is_error:expr, $msg:expr) => {
match $option {
Some(value) => value,
None => {
2025-07-24 14:27:57 +02:00
if $msg_is_error {
error!("{}", $msg);
} else {
debug!("{}", $msg);
}
2025-03-11 14:56:10 +01:00
return axum::http::StatusCode::BAD_REQUEST.into_response();
2025-01-31 23:08:23 +01:00
}
}
};
($option:expr) => {
match $option {
Some(value) => value,
2025-03-11 14:56:10 +01:00
None => return axum::http::StatusCode::BAD_REQUEST.into_response(),
2025-01-31 23:08:23 +01:00
}
};
}
2026-01-20 10:09:58 +01:00
#[macro_export]
macro_rules! internal_server_error {
() => {
axum::http::StatusCode::INTERNAL_SERVER_ERROR.into_response()
};
}
2025-01-31 23:08:23 +01:00
#[macro_export]
2025-10-16 11:58:40 +02:00
macro_rules! try_unwrap_body {
($body:expr) => {
2026-02-14 17:52:35 +01:00
$body
.map_or_else(|_| axum::http::StatusCode::INTERNAL_SERVER_ERROR.into_response(), |resp| resp.into_response())
2025-10-16 11:58:40 +02:00
};
}
#[macro_export]
macro_rules! try_result_or_status {
($option:expr, $status:expr, $msg_is_error:expr, $msg:expr) => {
2025-01-31 23:08:23 +01:00
match $option {
Ok(value) => value,
Err(_) => {
2025-07-24 14:27:57 +02:00
if $msg_is_error {
error!("{}", $msg);
} else {
debug!("{}", $msg);
}
2025-10-16 11:58:40 +02:00
return $status.into_response();
2025-01-31 23:08:23 +01:00
}
}
};
2025-10-16 11:58:40 +02:00
($option:expr, $status:expr) => {
2025-01-31 23:08:23 +01:00
match $option {
Ok(value) => value,
2025-10-16 11:58:40 +02:00
Err(_) => return $status.into_response(),
2025-01-31 23:08:23 +01:00
}
};
}
#[macro_export]
2025-10-16 11:58:40 +02:00
macro_rules! try_result_bad_request {
($option:expr, $msg_is_error:expr, $msg:expr) => {
2026-02-14 17:52:35 +01:00
$crate::api::api_utils::try_result_or_status!($option, axum::http::StatusCode::BAD_REQUEST, $msg_is_error, $msg)
2025-10-16 11:58:40 +02:00
};
($option:expr) => {
2026-02-14 17:52:35 +01:00
$crate::api::api_utils::try_result_or_status!($option, axum::http::StatusCode::BAD_REQUEST)
2025-10-16 11:58:40 +02:00
};
}
#[macro_export]
macro_rules! try_result_not_found {
($option:expr, $msg_is_error:expr, $msg:expr) => {
2026-02-14 17:52:35 +01:00
$crate::api::api_utils::try_result_or_status!($option, axum::http::StatusCode::NOT_FOUND, $msg_is_error, $msg)
2025-10-16 11:58:40 +02:00
};
($option:expr) => {
2026-02-14 17:52:35 +01:00
$crate::api::api_utils::try_result_or_status!($option, axum::http::StatusCode::NOT_FOUND)
};
}
2026-02-14 17:52:35 +01:00
use crate::api::panel_api::{can_provision_on_exhausted, create_panel_api_provisioning_stream_details};
pub use internal_server_error;
use shared::error::TuliproxError;
2025-02-05 20:01:49 +01:00
pub use try_option_bad_request;
pub use try_result_bad_request;
2025-10-16 11:58:40 +02:00
pub use try_result_not_found;
2026-02-14 17:52:35 +01:00
pub use try_result_or_status;
pub use try_unwrap_body;
2025-04-22 00:21:45 +02:00
pub fn get_server_time() -> String {
2026-02-14 17:52:35 +01:00
chrono::offset::Local::now().with_timezone(&chrono::Local).format("%Y-%m-%d %H:%M:%S %Z").to_string()
2025-04-22 00:21:45 +02:00
}
pub fn get_build_time() -> Option<String> {
2025-07-24 14:27:57 +02:00
BUILD_TIMESTAMP
.to_string()
.parse::<DateTime<Utc>>()
.ok()
.map(|datetime| datetime.format("%Y-%m-%d %H:%M:%S %Z").to_string())
2025-04-22 00:21:45 +02:00
}
2025-04-10 00:06:48 +02:00
#[allow(clippy::missing_panics_doc)]
pub async fn serve_file(file_path: &Path, mime_type: String, cache_control: Option<&str>) -> impl IntoResponse + Send {
2025-11-20 14:22:26 +01:00
match tokio::fs::try_exists(file_path).await {
Ok(exists) => {
2025-12-16 16:23:15 +01:00
if !exists {
2025-11-20 14:22:26 +01:00
return axum::http::StatusCode::NOT_FOUND.into_response();
2025-03-11 14:56:10 +01:00
}
2025-11-20 14:22:26 +01:00
}
Err(err) => {
error!("Failed to open file {}, {err:?}", file_path.display());
2025-11-20 14:22:26 +01:00
return axum::http::StatusCode::NOT_FOUND.into_response();
}
}
match tokio::fs::File::open(file_path).await {
Ok(file) => {
2026-02-14 17:52:35 +01:00
let last_modified = file.metadata().await.ok().and_then(|m| m.modified().ok()).map(|m| {
let dt: DateTime<Utc> = m.into();
dt.format("%a, %d %b %Y %H:%M:%S GMT").to_string()
});
let reader = async_file_reader(file);
2025-11-20 14:22:26 +01:00
let stream = tokio_util::io::ReaderStream::new(reader);
let body = axum::body::Body::from_stream(stream);
let mut builder = axum::response::Response::builder()
2025-11-20 14:22:26 +01:00
.status(axum::http::StatusCode::OK)
2025-12-25 12:33:28 +01:00
.header(axum::http::header::CONTENT_TYPE, mime_type)
2026-02-14 17:52:35 +01:00
.header(axum::http::header::CACHE_CONTROL, cache_control.unwrap_or("no-cache"));
if let Some(lm) = last_modified {
builder = builder.header(axum::http::header::LAST_MODIFIED, lm);
}
try_unwrap_body!(builder.body(body))
2025-11-20 14:22:26 +01:00
}
2026-01-20 10:09:58 +01:00
Err(_) => internal_server_error!(),
}
}
2025-07-24 14:27:57 +02:00
pub fn get_user_target_by_username(
username: &str,
2025-12-19 21:13:12 +01:00
app_state: &Arc<AppState>,
2025-07-24 14:27:57 +02:00
) -> Option<(ProxyUserCredentials, Arc<ConfigTarget>)> {
if !username.is_empty() {
2025-06-24 12:31:38 +02:00
return app_state.app_config.get_target_for_username(username);
}
None
}
2025-07-24 14:27:57 +02:00
pub fn get_user_target_by_credentials<'a>(
username: &str,
password: &str,
api_req: &'a UserApiRequest,
app_state: &'a AppState,
) -> Option<(ProxyUserCredentials, Arc<ConfigTarget>)> {
if !username.is_empty() && !password.is_empty() {
2025-06-24 12:31:38 +02:00
app_state.app_config.get_target_for_user(username, password)
} else {
2023-10-25 21:07:38 +02:00
let token = api_req.token.as_str().trim();
2024-05-10 12:01:48 +02:00
if token.is_empty() {
None
2024-05-10 12:01:48 +02:00
} else {
2025-06-24 12:31:38 +02:00
app_state.app_config.get_target_for_user_by_token(token)
}
}
2023-11-22 21:06:58 +01:00
}
2025-07-24 14:27:57 +02:00
pub fn get_user_target<'a>(
api_req: &'a UserApiRequest,
app_state: &'a AppState,
) -> Option<(ProxyUserCredentials, Arc<ConfigTarget>)> {
2023-11-22 21:06:58 +01:00
let username = api_req.username.as_str().trim();
let password = api_req.password.as_str().trim();
get_user_target_by_credentials(username, password, api_req, app_state)
}
2025-03-16 17:08:00 +01:00
pub struct StreamOptions {
pub stream_retry: bool,
pub buffer_enabled: bool,
pub buffer_size: usize,
2025-03-20 11:57:32 +01:00
pub pipe_provider_stream: bool,
2025-03-16 17:08:00 +01:00
}
2025-05-09 17:13:37 +02:00
/// Constructs a `StreamOptions` object based on the application's reverse proxy configuration.
///
/// This function retrieves streaming-related settings from the `AppState`:
/// - `stream_retry`: whether retrying the stream is enabled,
/// - `buffer_enabled`: whether stream buffering is enabled,
/// - `buffer_size`: the size of the stream buffer.
///
/// If the reverse proxy or stream settings are not defined, default values are used:
/// - retry: `false`
/// - buffering: `false`
/// - buffer size: `0`
///
/// Additionally, it computes `pipe_provider_stream`, which is `true` only if
/// both retry and buffering are disabled—indicating that the stream can be piped directly
/// from the provider without additional handling.
///
/// Returns a `StreamOptions` instance with the resolved configuration.
2025-12-19 21:13:12 +01:00
fn get_stream_options(app_state: &Arc<AppState>) -> StreamOptions {
2025-12-13 13:54:01 +01:00
let (stream_retry, buffer_enabled, buffer_size) = app_state
2025-07-24 14:27:57 +02:00
.app_config
.config
.load()
2025-01-14 10:47:37 +01:00
.reverse_proxy
.as_ref()
.and_then(|reverse_proxy| reverse_proxy.stream.as_ref())
2025-12-13 13:54:01 +01:00
.map_or((false, false, 0), |stream| {
2026-02-14 17:52:35 +01:00
let (buffer_enabled, buffer_size) =
stream.buffer.as_ref().map_or((false, 0), |buffer| (buffer.enabled, buffer.size));
(stream.retry, buffer_enabled, buffer_size)
2025-01-14 10:47:37 +01:00
});
2025-01-27 15:58:57 +01:00
let pipe_provider_stream = !stream_retry && !buffer_enabled;
2026-02-14 17:52:35 +01:00
StreamOptions { stream_retry, buffer_enabled, buffer_size, pipe_provider_stream }
2025-01-27 15:58:57 +01:00
}
2026-02-14 17:52:35 +01:00
pub fn get_stream_alternative_url(stream_url: &str, input: &ConfigInput, alias_input: &Arc<ProviderConfig>) -> String {
2025-07-24 14:27:57 +02:00
let Some(input_user_info) = input.get_user_info() else {
2026-01-20 10:09:58 +01:00
return stream_url.to_string();
2025-07-24 14:27:57 +02:00
};
let Some(alt_input_user_info) = alias_input.get_user_info() else {
2026-01-20 10:09:58 +01:00
return stream_url.to_string();
2025-07-24 14:27:57 +02:00
};
2025-03-19 09:54:59 +01:00
2025-10-24 20:03:47 +02:00
let modified = stream_url.replacen(&input_user_info.base_url, &alt_input_user_info.base_url, 1);
let modified = modified.replacen(&input_user_info.username, &alt_input_user_info.username, 1);
modified.replacen(&input_user_info.password, &alt_input_user_info.password, 1)
2025-03-19 09:54:59 +01:00
}
2026-01-20 10:09:58 +01:00
async fn get_redirect_alternative_url(
2025-12-19 21:13:12 +01:00
app_state: &Arc<AppState>,
2026-01-20 10:09:58 +01:00
redirect_url: &Arc<str>,
2025-07-24 14:27:57 +02:00
input: &ConfigInput,
2026-01-20 10:09:58 +01:00
) -> Arc<str> {
if let Some((base_url, username, password)) = input.get_matched_config_by_url(redirect_url) {
2026-02-14 17:52:35 +01:00
if let Some(provider_cfg) = app_state.active_provider.get_next_provider(&input.name).await {
let mut new_url = redirect_url.replacen(base_url, provider_cfg.url.as_str(), 1);
if let (Some(old_username), Some(old_password)) = (username, password) {
2026-02-14 17:52:35 +01:00
if let (Some(new_username), Some(new_password)) =
(provider_cfg.username.as_ref(), provider_cfg.password.as_ref())
{
new_url = new_url.replacen(old_username, new_username, 1);
new_url = new_url.replacen(old_password, new_password, 1);
2026-01-20 10:09:58 +01:00
return new_url.into();
}
// one has credentials the other not, something not right
2026-01-20 10:09:58 +01:00
return redirect_url.clone();
}
2026-01-20 10:09:58 +01:00
return new_url.into();
}
}
2026-01-20 10:09:58 +01:00
redirect_url.clone()
}
2025-05-09 17:13:37 +02:00
/// Determines the appropriate streaming strategy for the given input and stream URL.
///
/// This function attempts to acquire a connection to a streaming provider, either using a forced provider
/// (if specified), or based on the input name. It then selects a corresponding `StreamingOption`:
///
/// - If no connections are available (`Exhausted`), it returns a custom stream indicating exhaustion.
/// - If a connection is available or in a grace period, it constructs a streaming URL accordingly:
/// - If the provider was forced or matches the input, the original URL is reused.
/// - Otherwise, an alternative URL is generated based on the provider and input.
///
/// The function returns:
/// - an optional `ProviderConnectionGuard` to manage the connection's lifecycle,
/// - a `ProviderStreamState` describing how the stream state is,
/// - and optional HTTP headers to include in the request.
///
/// This logic helps abstract the decision-making behind provider selection and stream URL resolution.
2025-07-24 14:27:57 +02:00
async fn resolve_streaming_strategy(
2025-12-19 18:30:17 +01:00
app_state: &Arc<AppState>,
2025-07-24 14:27:57 +02:00
stream_url: &str,
2025-10-30 13:00:25 +01:00
fingerprint: &Fingerprint,
2025-07-24 14:27:57 +02:00
input: &ConfigInput,
2026-01-20 10:09:58 +01:00
force_provider: Option<&Arc<str>>,
2026-02-24 12:51:51 +01:00
allow_provider_grace: bool,
2026-03-09 17:12:46 +01:00
user_priority: i8,
2025-07-24 14:27:57 +02:00
) -> StreamingStrategy {
2025-05-09 17:13:37 +02:00
// allocate a provider connection
2026-02-24 16:03:39 +01:00
let mut forced_provider_allocated = false;
2026-01-20 18:51:46 +01:00
let provider_connection_handle = match force_provider {
2026-02-24 16:03:39 +01:00
Some(provider) => {
// First try to stay on the exact pinned provider account without over-allocating.
// If that account is no longer available, fall back to any available account in the same lineup.
if let Some(handle) = app_state
.active_provider
2026-03-09 17:12:46 +01:00
.acquire_exact_connection_with_grace(provider, &fingerprint.addr, allow_provider_grace, user_priority)
2026-02-24 16:03:39 +01:00
.await
{
forced_provider_allocated = true;
Some(handle)
} else {
debug_if_enabled!(
"Pinned provider {} unavailable for {}; falling back to lineup allocation",
sanitize_sensitive_info(provider),
sanitize_sensitive_info(&fingerprint.addr.to_string())
);
app_state
.active_provider
2026-03-09 17:12:46 +01:00
.acquire_connection_with_grace(&input.name, &fingerprint.addr, allow_provider_grace, user_priority)
2026-02-24 16:03:39 +01:00
.await
}
}
2026-02-24 12:51:51 +01:00
None => {
app_state
.active_provider
2026-03-09 17:12:46 +01:00
.acquire_connection_with_grace(&input.name, &fingerprint.addr, allow_provider_grace, user_priority)
2026-02-24 12:51:51 +01:00
.await
}
2025-04-24 18:12:29 +02:00
};
2025-07-21 22:17:57 +02:00
2026-01-20 18:51:46 +01:00
// panel_api provisioning/loading is handled later in the stream creation flow
2025-12-16 23:03:24 +01:00
2026-02-14 17:52:35 +01:00
let stream_response_params = if let Some(allocation) = provider_connection_handle.as_ref().map(|ph| &ph.allocation)
{
match allocation {
ProviderAllocation::Exhausted => {
debug!("Provider {} is exhausted. No connections allowed.", input.name);
let stream = create_provider_connections_exhausted_stream(&app_state.app_config, &[]);
ProviderStreamState::Custom(stream)
}
ProviderAllocation::Available(ref provider_cfg) | ProviderAllocation::GracePeriod(ref provider_cfg) => {
// 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 don't need to get new url
2026-02-24 16:03:39 +01:00
let (selected_provider_name, url) = if forced_provider_allocated || provider_cfg.id == input.id {
2026-02-14 17:52:35 +01:00
(input.name.clone(), stream_url.to_string())
} else {
(provider_cfg.name.clone(), get_stream_alternative_url(stream_url, input, provider_cfg))
};
2026-02-14 17:52:35 +01:00
debug_if_enabled!(
"provider session: input={} provider_cfg={} user={} allocation={} stream_url={}",
sanitize_sensitive_info(&input.name),
sanitize_sensitive_info(&provider_cfg.name),
sanitize_sensitive_info(
provider_cfg.get_user_info().as_ref().map_or_else(|| "?", |u| u.username.as_str())
),
allocation.short_key(),
2026-03-10 13:49:52 +01:00
sanitize_sensitive_info(resolve_request_url_for_logging(input, &url).as_ref())
2026-02-14 17:52:35 +01:00
);
if matches!(allocation, ProviderAllocation::Available(_)) {
ProviderStreamState::Available(Some(selected_provider_name.intern()), url.intern())
} else {
ProviderStreamState::GracePeriod(Some(selected_provider_name.intern()), url.intern())
2025-11-04 18:01:44 +01:00
}
2025-03-19 09:54:59 +01:00
}
2026-02-14 17:52:35 +01:00
}
} else {
debug!("Provider {} is exhausted. No connections allowed.", input.name);
let stream = create_provider_connections_exhausted_stream(&app_state.app_config, &[]);
ProviderStreamState::Custom(stream)
};
2025-11-04 18:01:44 +01:00
2025-05-09 17:13:37 +02:00
StreamingStrategy {
2025-11-04 18:01:44 +01:00
provider_handle: provider_connection_handle,
2025-05-09 17:13:37 +02:00
provider_stream_state: stream_response_params,
input_headers: Some(input.headers.clone()),
2025-05-09 17:13:37 +02:00
}
2025-03-19 09:54:59 +01:00
}
2025-07-24 14:27:57 +02:00
fn get_grace_period_millis(
connection_permission: UserConnectionPermission,
stream_response_params: &ProviderStreamState,
config_grace_period_millis: u64,
) -> u64 {
if config_grace_period_millis > 0
&& (
2026-02-14 17:52:35 +01:00
matches!(stream_response_params, ProviderStreamState::GracePeriod(_, _)) // provider grace period
2025-07-24 14:27:57 +02:00
|| connection_permission == UserConnectionPermission::GracePeriod
2026-02-14 17:52:35 +01:00
// user grace period
)
2025-07-24 14:27:57 +02:00
{
config_grace_period_millis
} else {
0
}
2025-04-17 12:11:51 +02:00
}
#[allow(clippy::too_many_arguments, clippy::too_many_lines)]
2025-07-24 14:27:57 +02:00
async fn create_stream_response_details(
2025-11-10 01:04:13 +01:00
app_state: &Arc<AppState>,
2025-07-24 14:27:57 +02:00
stream_options: &StreamOptions,
stream_url: &str,
2025-10-30 13:00:25 +01:00
fingerprint: &Fingerprint,
2025-07-24 14:27:57 +02:00
req_headers: &HeaderMap,
2026-02-10 11:53:42 +01:00
input: &Arc<ConfigInput>,
2025-07-24 14:27:57 +02:00
item_type: PlaylistItemType,
share_stream: bool,
connection_permission: UserConnectionPermission,
2026-01-20 10:09:58 +01:00
force_provider: Option<&Arc<str>>,
2026-02-24 12:51:51 +01:00
allow_provider_grace: bool,
2026-01-20 18:51:46 +01:00
virtual_id: VirtualId,
2026-03-09 17:12:46 +01:00
user_priority: i8,
2026-02-10 11:53:42 +01:00
) -> Result<StreamDetails, TuliproxError> {
2026-02-14 17:52:35 +01:00
let mut streaming_strategy =
2026-03-09 17:12:46 +01:00
resolve_streaming_strategy(app_state, stream_url, fingerprint, input, force_provider, allow_provider_grace, user_priority)
.await;
let mut grace_period_options = app_state.get_grace_options();
grace_period_options.period_millis = get_grace_period_millis(
2025-12-18 11:07:30 +01:00
connection_permission,
&streaming_strategy.provider_stream_state,
grace_period_options.period_millis,
2025-12-18 11:07:30 +01:00
);
2026-03-11 12:25:54 +01:00
let provider_grace_active = matches!(
streaming_strategy.provider_stream_state,
ProviderStreamState::GracePeriod(_, _)
);
2025-12-18 11:07:30 +01:00
2026-02-14 17:52:35 +01:00
let guard_provider_name =
streaming_strategy.provider_handle.as_ref().and_then(|guard| guard.allocation.get_provider_name());
2025-12-18 11:07:30 +01:00
2026-01-20 18:51:46 +01:00
if matches!(streaming_strategy.provider_stream_state, ProviderStreamState::Custom(_))
&& can_provision_on_exhausted(app_state, input)
{
if let Some(handle) = streaming_strategy.provider_handle.take() {
2026-02-14 17:52:35 +01:00
app_state.connection_manager.release_provider_handle(Some(handle)).await;
2026-01-20 18:51:46 +01:00
}
debug_if_enabled!(
"panel_api: provider connections exhausted; sending provisioning stream for input {}",
sanitize_sensitive_info(&input.name)
);
2026-02-10 11:53:42 +01:00
return Ok(create_panel_api_provisioning_stream_details(
2026-01-20 18:51:46 +01:00
app_state,
input,
guard_provider_name.clone(),
&grace_period_options,
2026-01-20 18:51:46 +01:00
fingerprint.addr,
virtual_id,
2026-02-10 11:53:42 +01:00
));
2026-01-20 18:51:46 +01:00
}
2025-05-09 17:13:37 +02:00
match streaming_strategy.provider_stream_state {
// custom stream means we display our own stream like connection exhausted, channel-unavailable...
2025-05-09 17:13:37 +02:00
ProviderStreamState::Custom(provider_stream) => {
2025-03-20 11:57:32 +01:00
let (stream, stream_info) = provider_stream;
2026-02-10 11:53:42 +01:00
Ok(StreamDetails {
2025-03-20 11:57:32 +01:00
stream,
stream_info,
2025-10-29 16:40:17 +01:00
provider_name: guard_provider_name.clone(),
2026-02-23 21:16:45 +01:00
request_url: None,
grace_period: grace_period_options,
2026-03-11 12:25:54 +01:00
provider_grace_active: false,
2026-01-20 18:51:46 +01:00
disable_provider_grace: false,
2025-03-20 11:57:32 +01:00
reconnect_flag: None,
2025-11-04 18:01:44 +01:00
provider_handle: streaming_strategy.provider_handle.clone(),
2026-02-10 11:53:42 +01:00
})
2025-03-20 11:57:32 +01:00
}
2025-10-29 16:40:17 +01:00
ProviderStreamState::Available(_provider_name, request_url)
| ProviderStreamState::GracePeriod(_provider_name, request_url) => {
2026-02-23 21:16:45 +01:00
debug_if_enabled!(
"Provider stream selection: allocated_provider={} actual_request_url={}",
sanitize_sensitive_info(guard_provider_name.as_deref().unwrap_or("?")),
2026-03-10 13:49:52 +01:00
sanitize_sensitive_info(resolve_request_url_for_logging(input, request_url.as_ref()).as_ref())
2026-02-23 21:16:45 +01:00
);
2026-03-11 12:25:54 +01:00
let defer_provider_stream_until_grace_check = if provider_grace_active && grace_period_options.hold_stream {
if let Some(provider_name) = guard_provider_name.as_ref() {
app_state.active_provider.is_over_limit(provider_name).await
} else {
false
}
} else {
false
};
let (stream, stream_info, reconnect_flag) = if defer_provider_stream_until_grace_check {
debug_if_enabled!(
"Deferring provider stream open until grace check completes for {}",
sanitize_sensitive_info(resolve_request_url_for_logging(input, request_url.as_ref()).as_ref())
2025-07-24 14:27:57 +02:00
);
2026-03-11 12:25:54 +01:00
(None, None, None)
} else {
let parsed_url = Url::parse(&request_url);
let ((stream, stream_info), reconnect_flag) = if let Ok(url) = parsed_url {
let default_user_agent = app_state.app_config.config.load().default_user_agent.clone();
let disabled_headers = app_state.get_disabled_headers();
let mut provider_stream_factory_options = ProviderStreamFactoryOptions::new(
fingerprint.addr,
item_type,
share_stream,
stream_options,
&url,
req_headers,
streaming_strategy.input_headers.as_ref(),
disabled_headers.as_ref(),
default_user_agent.as_deref(),
);
2026-02-10 11:53:42 +01:00
2026-03-11 12:25:54 +01:00
let provider_config = input.get_resolve_provider(url.as_ref());
provider_stream_factory_options.set_provider(provider_config);
2026-02-10 11:53:42 +01:00
2026-03-11 12:25:54 +01:00
let reconnect_flag = provider_stream_factory_options.get_reconnect_flag_clone();
let provider_stream = match create_provider_stream(
app_state,
&app_state.http_client.load(),
provider_stream_factory_options,
)
.await
{
None => (None, None),
Some((stream, info)) => (Some(stream), info),
};
(provider_stream, Some(reconnect_flag))
} else {
((None, None), None)
2025-05-09 17:13:37 +02:00
};
2026-03-11 12:25:54 +01:00
(stream, stream_info, reconnect_flag)
2025-03-20 11:57:32 +01:00
};
2025-03-19 09:54:59 +01:00
2025-03-20 11:57:32 +01:00
if log_enabled!(log::Level::Debug) {
2025-11-10 01:04:13 +01:00
if let Some((headers, status_code, response_url, _custom_video_type)) = stream_info.as_ref() {
2025-03-20 11:57:32 +01:00
debug!(
"Responding stream request {} with status {}, headers {:?}",
2026-02-14 17:52:35 +01:00
sanitize_sensitive_info(response_url.as_ref().map_or(stream_url, |s| s.as_str())),
2025-03-20 11:57:32 +01:00
status_code,
headers
);
}
}
2025-03-19 09:54:59 +01:00
2026-03-11 12:25:54 +01:00
// if we have no stream, we should release the provider unless we intentionally defer
// opening it for provider-grace hold_stream behavior.
let provider_handle = if stream.is_none() && !defer_provider_stream_until_grace_check {
2025-11-10 01:04:13 +01:00
let provider_handle = streaming_strategy.provider_handle.take();
app_state.connection_manager.release_provider_handle(provider_handle).await;
error!("Can't open stream {}", sanitize_sensitive_info(&request_url));
2025-11-10 01:04:13 +01:00
None
} else {
streaming_strategy.provider_handle.take()
};
2025-10-25 20:40:58 +02:00
2026-02-10 11:53:42 +01:00
Ok(StreamDetails {
2025-03-20 11:57:32 +01:00
stream,
stream_info,
2025-10-29 16:40:17 +01:00
provider_name: guard_provider_name.clone(),
2026-02-23 21:16:45 +01:00
request_url: Some(request_url.clone()),
grace_period: grace_period_options,
2026-03-11 12:25:54 +01:00
provider_grace_active,
2026-01-20 18:51:46 +01:00
disable_provider_grace: false,
2025-03-20 11:57:32 +01:00
reconnect_flag,
2025-11-10 01:30:23 +01:00
provider_handle,
2026-02-10 11:53:42 +01:00
})
2025-03-19 09:54:59 +01:00
}
}
}
2025-04-19 13:50:41 +02:00
pub struct RedirectParams<'a, P>
where
2025-04-19 15:01:23 +02:00
P: PlaylistEntry,
2025-04-19 13:50:41 +02:00
{
pub item: &'a P,
pub provider_id: Option<u32>,
pub cluster: XtreamCluster,
pub target_type: TargetType,
pub target: &'a ConfigTarget,
2025-04-24 18:12:29 +02:00
pub input: &'a ConfigInput,
2025-04-19 13:50:41 +02:00
pub user: &'a ProxyUserCredentials,
pub stream_ext: Option<&'a str>,
2025-06-03 18:25:15 +02:00
pub req_context: ApiStreamContext,
2025-04-19 13:50:41 +02:00
pub action_path: &'a str,
}
2025-04-19 15:01:23 +02:00
impl<P> RedirectParams<'_, P>
2025-04-19 13:50:41 +02:00
where
2025-04-19 15:01:23 +02:00
P: PlaylistEntry,
2025-04-19 13:50:41 +02:00
{
2025-04-19 15:01:23 +02:00
pub fn get_query_path(&self, provider_id: u32, url: &str) -> String {
2026-02-14 17:52:35 +01:00
let extension =
self.stream_ext.map_or_else(|| extract_extension_from_url(url).unwrap_or_default(), ToString::to_string);
2025-04-19 15:01:23 +02:00
// if there is an action_path (like for timeshift duration/start), it will be added in front of the stream_id
2025-04-19 15:01:23 +02:00
if self.action_path.is_empty() {
concat_string!(&provider_id.to_string(), &extension)
2025-04-19 15:01:23 +02:00
} else {
concat_string!(&trim_slash(self.action_path), "/", &provider_id.to_string(), &extension)
2025-04-19 15:01:23 +02:00
}
}
}
2025-04-19 13:50:41 +02:00
2025-07-24 14:27:57 +02:00
pub async fn redirect_response<'a, P>(
2025-12-19 21:13:12 +01:00
app_state: &Arc<AppState>,
2025-07-24 14:27:57 +02:00
params: &'a RedirectParams<'a, P>,
) -> Option<impl IntoResponse + Send>
2025-04-19 15:01:23 +02:00
where
P: PlaylistEntry,
{
2025-04-19 13:50:41 +02:00
let item_type = params.item.get_item_type();
2025-12-26 16:03:03 +01:00
let provider_url = params.item.get_provider_url();
2025-04-19 13:50:41 +02:00
2026-02-14 17:52:35 +01:00
let redirect_request = params.user.proxy.is_redirect(item_type) || params.target.is_force_redirect(item_type);
let is_hls_request = item_type == PlaylistItemType::LiveHls || params.stream_ext == Some(HLS_EXT);
let is_dash_request =
(!is_hls_request && item_type == PlaylistItemType::LiveDash) || params.stream_ext == Some(DASH_EXT);
2025-04-19 13:50:41 +02:00
if params.target_type == TargetType::M3u {
if redirect_request || is_dash_request {
2026-01-20 10:09:58 +01:00
let redirect_url: Arc<str> = if is_hls_request {
replace_url_extension(&provider_url, HLS_EXT).into()
2025-07-24 14:27:57 +02:00
} else {
2026-01-20 10:09:58 +01:00
provider_url.clone()
2025-07-24 14:27:57 +02:00
};
let redirect_url =
2026-02-14 17:52:35 +01:00
if is_dash_request { replace_url_extension(&redirect_url, DASH_EXT).into() } else { redirect_url };
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));
2025-04-24 18:12:29 +02:00
return Some(redirect(&redirect_url).into_response());
2025-04-19 13:50:41 +02:00
}
} else if params.target_type == TargetType::Xtream {
let Some(provider_id) = params.provider_id else {
2026-01-20 10:09:58 +01:00
return Some(StatusCode::BAD_REQUEST.into_response());
2025-04-19 13:50:41 +02:00
};
2025-04-19 15:01:23 +02:00
if redirect_request {
// handle redirect for series but why?
2025-04-19 15:01:23 +02:00
if params.cluster == XtreamCluster::Series {
let ext = params.stream_ext.unwrap_or_default();
2025-04-24 18:12:29 +02:00
let url = params.input.url.as_str();
2026-02-14 17:52:35 +01:00
let (username, password) =
(params.input.username.as_deref().unwrap_or(""), params.input.password.as_deref().unwrap_or(""));
2025-04-19 15:01:23 +02:00
// 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 {}",
2026-03-10 13:49:52 +01:00
sanitize_sensitive_info(resolve_request_url_for_logging(params.input, &stream_url).as_ref())
);
2025-04-19 15:01:23 +02:00
return Some(redirect(&stream_url).into_response());
}
2025-04-19 13:50:41 +02:00
2025-04-19 15:01:23 +02:00
let target_name = params.target.name.as_str();
let virtual_id = params.item.get_virtual_id();
2025-07-24 14:27:57 +02:00
let stream_url = match get_xtream_player_api_stream_url(
params.input,
params.req_context,
2025-12-26 16:03:03 +01:00
&params.get_query_path(provider_id, &provider_url),
&provider_url,
2025-07-24 14:27:57 +02:00
) {
2025-04-19 15:01:23 +02:00
None => {
2026-02-14 17:52:35 +01:00
error!(
"Can't find stream url for target {target_name}, context {}, stream_id {virtual_id}",
params.req_context
);
2026-01-20 10:09:58 +01:00
return Some(StatusCode::BAD_REQUEST.into_response());
2025-04-19 15:01:23 +02:00
}
2026-02-14 17:52:35 +01:00
Some(url) => match app_state.active_provider.get_next_provider(&params.input.name).await {
Some(provider_cfg) => get_stream_alternative_url(&url, params.input, &provider_cfg),
None => url.to_string(),
},
2025-04-19 15:01:23 +02:00
};
2025-04-19 13:50:41 +02:00
2025-04-19 15:01:23 +02:00
// hls or dash redirect
if is_dash_request {
2025-07-24 14:27:57 +02:00
let redirect_url = if is_hls_request {
&replace_url_extension(&stream_url, HLS_EXT)
} else {
&replace_url_extension(&stream_url, DASH_EXT)
};
debug_if_enabled!(
"Redirecting stream request to {}",
2026-03-10 13:49:52 +01:00
sanitize_sensitive_info(resolve_request_url_for_logging(params.input, redirect_url).as_ref())
);
2025-04-19 15:01:23 +02:00
return Some(redirect(redirect_url).into_response());
2025-04-19 13:50:41 +02:00
}
debug_if_enabled!(
"Redirecting stream request to {}",
2026-03-10 13:49:52 +01:00
sanitize_sensitive_info(resolve_request_url_for_logging(params.input, &stream_url).as_ref())
);
2025-04-19 13:50:41 +02:00
return Some(redirect(&stream_url).into_response());
}
}
None
}
2025-04-22 17:08:31 +02:00
fn is_throttled_stream(item_type: PlaylistItemType, throttle_kbps: usize) -> bool {
2025-07-24 14:27:57 +02:00
throttle_kbps > 0
&& matches!(
item_type,
PlaylistItemType::Video
| PlaylistItemType::Series
| PlaylistItemType::SeriesInfo
| PlaylistItemType::Catchup
2025-12-16 19:59:53 +01:00
| PlaylistItemType::LocalVideo
| PlaylistItemType::LocalSeries
| PlaylistItemType::LocalSeriesInfo
2025-07-24 14:27:57 +02:00
)
2025-04-22 17:08:31 +02:00
}
2026-02-14 17:52:35 +01:00
fn prepare_body_stream<S>(app_state: &Arc<AppState>, item_type: PlaylistItemType, stream: S) -> axum::body::Body
2025-12-18 15:55:54 +01:00
where
2026-02-14 17:52:35 +01:00
S: futures::Stream<Item = Result<bytes::Bytes, StreamError>> + Send + 'static,
2025-12-18 15:55:54 +01:00
{
2025-04-24 18:12:29 +02:00
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) {
2026-01-02 13:55:19 +01:00
info!("Stream throttling active: {}", human_readable_kbps(u64::try_from(throttle_kbps).unwrap_or_default()));
2025-04-24 18:12:29 +02:00
axum::body::Body::from_stream(ThrottledStream::new(stream.boxed(), throttle_kbps))
} else {
axum::body::Body::from_stream(stream)
};
body_stream
}
2025-04-22 17:08:31 +02:00
/// # Panics
2025-07-24 14:27:57 +02:00
pub async fn force_provider_stream_response(
2025-10-30 13:00:25 +01:00
fingerprint: &Fingerprint,
2025-11-10 01:04:13 +01:00
app_state: &Arc<AppState>,
2025-07-24 14:27:57 +02:00
user_session: &UserSession,
2025-10-24 19:47:39 +02:00
mut stream_channel: StreamChannel,
2025-07-24 14:27:57 +02:00
req_headers: &HeaderMap,
2026-02-10 11:53:42 +01:00
input: &Arc<ConfigInput>,
2025-07-24 14:27:57 +02:00
user: &ProxyUserCredentials,
) -> impl IntoResponse + Send {
let stream_options = get_stream_options(app_state);
let share_stream = false;
let connection_permission = UserConnectionPermission::Allowed;
2025-10-24 18:36:13 +02:00
let item_type = stream_channel.item_type;
2025-04-22 17:08:31 +02:00
// Release the existing provider connection for this session before acquiring a new one.
// This is critical for users with a connection limit of 1 to avoid "Provider exhausted" or provider-side 502/509 errors during seeking.
app_state.connection_manager.release_provider_connection(&user_session.addr).await;
2026-02-24 16:03:39 +01:00
// Keep seek/range reconnects on the same provider account whenever that exact account is still available.
// If it is not available anymore, resolve_streaming_strategy will fall back to another free account.
let preferred_provider = Some(&user_session.provider);
// Never allow provider-side grace for forced seek/session reacquire.
// Over-allocation here would break provider-side one-connection limits.
let allow_provider_grace = false;
2026-02-24 12:51:51 +01:00
2026-02-10 11:53:42 +01:00
let stream_details = match create_stream_response_details(
2025-07-24 14:27:57 +02:00
app_state,
&stream_options,
&user_session.stream_url,
2025-10-30 13:00:25 +01:00
fingerprint,
2025-07-24 14:27:57 +02:00
req_headers,
input,
item_type,
share_stream,
connection_permission,
2026-02-24 12:51:51 +01:00
preferred_provider,
allow_provider_grace,
2026-01-20 18:51:46 +01:00
stream_channel.virtual_id,
2026-03-09 17:12:46 +01:00
user.priority,
2025-07-24 14:27:57 +02:00
)
2026-02-14 17:52:35 +01:00
.await
{
2026-02-10 11:53:42 +01:00
Ok(stream_details) => stream_details,
Err(err) => {
error!("Failed to stream: {err}");
return StatusCode::INTERNAL_SERVER_ERROR.into_response();
}
};
2025-04-22 17:08:31 +02:00
2026-03-11 12:25:54 +01:00
let deferred_grace_hold_stream =
!stream_details.has_stream() && stream_details.provider_grace_active && stream_details.grace_period.hold_stream;
if stream_details.has_stream() || deferred_grace_hold_stream {
2026-02-14 17:52:35 +01:00
let provider_response =
stream_details.stream_info.as_ref().map(|(h, sc, url, cvt)| (h.clone(), *sc, url.clone(), *cvt));
app_state.active_users.update_session_addr(&user.username, &user_session.token, &fingerprint.addr).await;
2025-10-24 19:47:39 +02:00
stream_channel.shared = share_stream;
2026-02-14 17:52:35 +01:00
let stream = create_active_client_stream(
stream_details,
app_state,
user,
connection_permission,
fingerprint,
stream_channel,
Some(&user_session.token),
req_headers,
)
.await;
2025-07-24 14:27:57 +02:00
2026-02-14 17:52:35 +01:00
let (status_code, header_map) = get_stream_response_with_headers(provider_response.map(|(h, s, _, _)| (h, s)));
2025-05-09 17:13:37 +02:00
let mut response = axum::response::Response::builder().status(status_code);
for (key, value) in &header_map {
response = response.header(key, value);
2025-04-22 17:08:31 +02:00
}
2026-01-20 18:51:46 +01:00
let body_stream = prepare_body_stream(app_state, item_type, stream);
2025-07-24 14:27:57 +02:00
debug_if_enabled!(
"Streaming provider forced stream request from {}",
2026-03-10 13:49:52 +01:00
sanitize_sensitive_info(resolve_request_url_for_logging(input, user_session.stream_url.as_ref()).as_ref())
2025-07-24 14:27:57 +02:00
);
return try_unwrap_body!(response.body(body_stream));
2025-04-22 17:08:31 +02:00
}
2025-11-10 01:04:13 +01:00
app_state.connection_manager.release_provider_handle(stream_details.provider_handle).await;
2026-02-14 17:52:35 +01:00
if let (Some(stream), _stream_info) =
create_channel_unavailable_stream(&app_state.app_config, &[], StatusCode::SERVICE_UNAVAILABLE)
{
app_state
.connection_manager
.update_stream_detail(&fingerprint.addr, CustomVideoStreamType::ChannelUnavailable)
.await;
2025-05-09 17:13:37 +02:00
debug!("Streaming custom stream");
2025-07-24 14:27:57 +02:00
try_unwrap_body!(axum::response::Response::builder()
2026-01-20 10:09:58 +01:00
.status(StatusCode::OK)
2025-07-24 14:27:57 +02:00
.body(axum::body::Body::from_stream(stream)))
2025-05-09 17:13:37 +02:00
} else {
2026-01-20 10:09:58 +01:00
StatusCode::BAD_REQUEST.into_response()
2025-05-09 17:13:37 +02:00
}
2025-04-22 17:08:31 +02:00
}
2025-04-10 10:43:14 +02:00
/// # Panics
2025-07-24 14:27:57 +02:00
#[allow(clippy::too_many_arguments, clippy::too_many_lines)]
pub async fn stream_response(
2025-10-30 13:00:25 +01:00
fingerprint: &Fingerprint,
2025-11-10 01:04:13 +01:00
app_state: &Arc<AppState>,
2025-07-24 14:27:57 +02:00
session_token: &str,
2025-10-24 19:47:39 +02:00
mut stream_channel: StreamChannel,
2025-07-24 14:27:57 +02:00
stream_url: &str,
req_headers: &HeaderMap,
2026-02-10 11:53:42 +01:00
input: &Arc<ConfigInput>,
target: &Arc<ConfigTarget>,
2025-07-24 14:27:57 +02:00
user: &ProxyUserCredentials,
connection_permission: UserConnectionPermission,
) -> impl IntoResponse + Send {
2026-03-10 13:49:52 +01:00
let request_log_stream_url = resolve_request_url_for_logging(input, stream_url);
2025-07-24 14:27:57 +02:00
if log_enabled!(log::Level::Trace) {
2026-03-10 13:49:52 +01:00
trace!("Try to open stream {}", sanitize_sensitive_info(request_log_stream_url.as_ref()));
2025-07-24 14:27:57 +02:00
}
2025-01-27 15:58:57 +01:00
2025-04-16 18:04:57 +02:00
if connection_permission == UserConnectionPermission::Exhausted {
2025-07-24 14:27:57 +02:00
return create_custom_video_stream_response(
2025-11-10 01:04:13 +01:00
app_state,
&fingerprint.addr,
2025-07-24 14:27:57 +02:00
CustomVideoStreamType::UserConnectionsExhausted,
2026-02-14 17:52:35 +01:00
)
.into_response();
2025-04-16 18:04:57 +02:00
}
2025-10-24 18:36:13 +02:00
let virtual_id = stream_channel.virtual_id;
let item_type = stream_channel.item_type;
2026-03-11 12:25:54 +01:00
// Grace candidates are intentionally not shared:
// they are transient and may be terminated after the grace check.
// Keeping them out of the shared-stream lifecycle avoids cross-effects on
// concurrently running streams of the same user.
let share_stream =
is_stream_share_enabled(item_type, target) && connection_permission != UserConnectionPermission::GracePeriod;
2025-11-30 11:51:45 +01:00
let _shared_lock = if share_stream {
let write_lock = app_state.app_config.file_locks.write_lock_str(stream_url).await;
2026-02-14 17:52:35 +01:00
if let Some(value) = try_shared_stream_response_if_any(
app_state,
stream_url,
fingerprint,
user,
connection_permission,
stream_channel.clone(),
session_token,
req_headers,
)
.await
2025-07-24 14:27:57 +02:00
{
2025-03-11 14:56:10 +01:00
return value.into_response();
2025-01-27 15:58:57 +01:00
}
2025-11-30 11:51:45 +01:00
Some(write_lock)
} else {
2026-02-12 19:18:16 +01:00
// Opportunistic cross-target sharing: if another target already runs a shared stream
// for the same provider URL, subscribe to it instead of opening a separate connection.
if item_type == PlaylistItemType::Live {
2026-02-14 17:52:35 +01:00
if let Some(value) = try_shared_stream_response_if_any(
app_state,
stream_url,
fingerprint,
user,
connection_permission,
stream_channel.clone(),
session_token,
req_headers,
)
.await
2026-02-12 19:18:16 +01:00
{
debug_if_enabled!(
"Opportunistic shared stream reuse for {}",
2026-03-10 13:49:52 +01:00
sanitize_sensitive_info(stream_url)
);
2026-02-12 19:18:16 +01:00
return value.into_response();
}
}
2025-11-30 11:51:45 +01:00
None
};
2025-01-14 10:47:37 +01:00
2025-03-19 09:54:59 +01:00
let stream_options = get_stream_options(app_state);
2026-02-10 11:53:42 +01:00
let mut stream_details = match create_stream_response_details(
2025-07-24 14:27:57 +02:00
app_state,
&stream_options,
stream_url,
2025-10-30 13:00:25 +01:00
fingerprint,
2025-07-24 14:27:57 +02:00
req_headers,
input,
item_type,
share_stream,
connection_permission,
None,
2026-02-24 12:51:51 +01:00
true,
2026-01-20 18:51:46 +01:00
stream_channel.virtual_id,
2026-03-09 17:12:46 +01:00
user.priority,
2026-02-14 17:52:35 +01:00
)
.await
{
2026-02-10 11:53:42 +01:00
Ok(stream_details) => stream_details,
Err(err) => {
error!("Failed to stream: {err}");
return StatusCode::INTERNAL_SERVER_ERROR.into_response();
}
};
2025-11-06 15:02:22 +01:00
2026-03-11 12:25:54 +01:00
let deferred_grace_hold_stream =
!stream_details.has_stream() && stream_details.provider_grace_active && stream_details.grace_period.hold_stream;
if stream_details.has_stream() || deferred_grace_hold_stream {
2025-03-19 09:54:59 +01:00
// let content_length = get_stream_content_length(provider_response.as_ref());
2025-07-24 14:27:57 +02:00
let provider_response = stream_details
.stream_info
.as_ref()
2025-11-10 01:04:13 +01:00
.map(|(h, sc, response_url, cvt)| (h.clone(), *sc, response_url.clone(), *cvt));
2025-10-25 20:40:58 +02:00
let provider_name = stream_details.provider_name.clone();
let actual_request_url = stream_details.request_url.clone().unwrap_or_else(|| Arc::<str>::from(stream_url));
2026-03-10 13:49:52 +01:00
let log_actual_request_url = resolve_request_url_for_logging(input, actual_request_url.as_ref());
2026-02-23 21:16:45 +01:00
debug_if_enabled!(
"Provider request mapping: allocated_provider={} actual_request_url={}",
sanitize_sensitive_info(provider_name.as_deref().unwrap_or("?")),
2026-03-10 13:49:52 +01:00
sanitize_sensitive_info(log_actual_request_url.as_ref())
2026-02-23 21:16:45 +01:00
);
2025-07-24 14:27:57 +02:00
2026-01-20 18:51:46 +01:00
if let Some((headers, status, _response_url, Some(CustomVideoStreamType::Provisioning))) =
stream_details.stream_info.as_ref()
{
2026-02-14 17:52:35 +01:00
debug_if_enabled!("panel_api provisioning response to client: status={} headers={:?}", status, headers);
2026-01-20 18:51:46 +01:00
}
let mut is_stream_shared = share_stream;
2026-03-11 12:25:54 +01:00
if stream_details.provider_grace_active {
is_stream_shared = false;
}
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;
}
}
2026-03-11 12:25:54 +01:00
let provider_handle = if is_stream_shared { stream_details.provider_handle.take() } else { None };
stream_channel.shared = is_stream_shared;
2026-02-14 17:52:35 +01:00
let stream = create_active_client_stream(
stream_details,
app_state,
user,
connection_permission,
fingerprint,
stream_channel,
Some(session_token),
req_headers,
)
.await;
let stream_resp = if is_stream_shared {
debug_if_enabled!(
"Streaming shared stream request from {}",
2026-03-10 13:49:52 +01:00
sanitize_sensitive_info(log_actual_request_url.as_ref())
);
2025-03-19 09:54:59 +01:00
// Shared Stream response
2026-02-14 17:52:35 +01:00
let shared_headers = provider_response.as_ref().map_or_else(Vec::new, |(h, _, _, _)| h.clone());
2025-07-27 20:47:35 +02:00
2025-10-29 16:40:17 +01:00
if let Some((broadcast_stream, _shared_provider)) = SharedStreamManager::register_shared_stream(
2025-07-24 14:27:57 +02:00
app_state,
stream_url,
stream,
2025-11-01 00:26:30 +01:00
&fingerprint.addr,
2025-07-24 14:27:57 +02:00
shared_headers,
stream_options.buffer_size,
2025-11-04 21:46:44 +01:00
provider_handle,
)
2026-02-14 17:52:35 +01:00
.await
{
2025-07-24 14:27:57 +02:00
let (status_code, header_map) =
2025-11-10 01:04:13 +01:00
get_stream_response_with_headers(provider_response.map(|(h, s, _, _)| (h, s)));
2025-07-24 14:27:57 +02:00
let mut response = axum::response::Response::builder().status(status_code);
2025-03-11 14:56:10 +01:00
for (key, value) in &header_map {
response = response.header(key, value);
}
try_unwrap_body!(response.body(axum::body::Body::from_stream(broadcast_stream)))
2025-03-19 09:54:59 +01:00
} else {
2026-01-20 10:09:58 +01:00
StatusCode::BAD_REQUEST.into_response()
2025-03-19 09:54:59 +01:00
}
} else {
// Previously, we would always check if the provider redirected the request.
// If the provider redirected a movie from /movie/... to a temporary /live/... URL,
// We would save that redirected URL in your session.
// When we tried to seek or pause/resume, we would use that saved /live/ URL.
// However, providers often make these redirect links ephemeral or restricted—they
// might not support seeking, or they might trigger a 509 error if accessed again.
// For Movies/Series: We now ignore the redirect and always save the original,
// canonical URL (the one starting with /movie/) in your session.
// This ensures that every time you seek, we start "fresh" with the correct provider handshake,
// preventing the session from being "poisoned" by a temporary redirect.
// For everything else (Live): It continues to work as before, using the redirected URL if available,
// which is often desirable for live streams to stay on the same edge server.
2026-02-23 21:16:45 +01:00
let session_url: Cow<'_, str> = if matches!(
2026-02-14 17:52:35 +01:00
item_type,
PlaylistItemType::Catchup
| PlaylistItemType::Video
| PlaylistItemType::LocalVideo
| PlaylistItemType::Series
| PlaylistItemType::LocalSeries
) {
2026-02-23 21:16:45 +01:00
Cow::Owned(actual_request_url.to_string())
} else {
provider_response
.as_ref()
.and_then(|(_, _, u, _)| u.as_ref())
2026-02-23 21:16:45 +01:00
.map_or_else(|| Cow::Owned(actual_request_url.to_string()), |url| Cow::Owned(url.to_string()))
};
2026-03-10 13:49:52 +01:00
let log_session_url = resolve_request_url_for_logging(input, session_url.as_ref());
2025-06-04 21:40:14 +02:00
if log_enabled!(log::Level::Debug) {
2026-03-10 13:49:52 +01:00
if log_session_url.eq(log_actual_request_url.as_ref()) {
debug!("Streaming stream request from {}", sanitize_sensitive_info(log_actual_request_url.as_ref()));
2025-06-04 21:40:14 +02:00
} else {
2025-07-24 14:27:57 +02:00
debug!(
"Streaming stream request for {} from {}",
2026-03-10 13:49:52 +01:00
sanitize_sensitive_info(log_actual_request_url.as_ref()),
sanitize_sensitive_info(log_session_url.as_ref())
2025-07-24 14:27:57 +02:00
);
2025-06-04 21:40:14 +02:00
}
}
2025-07-24 14:27:57 +02:00
let (status_code, header_map) =
2025-11-10 01:04:13 +01:00
get_stream_response_with_headers(provider_response.map(|(h, s, _, _)| (h, s)));
2025-04-24 18:12:29 +02:00
let mut response = axum::response::Response::builder().status(status_code);
2025-03-19 09:54:59 +01:00
for (key, value) in &header_map {
response = response.header(key, value);
}
2025-03-26 15:06:05 +01:00
2025-04-22 17:08:31 +02:00
if let Some(provider) = provider_name {
2025-07-24 14:27:57 +02:00
if matches!(
item_type,
PlaylistItemType::LiveHls
| PlaylistItemType::LiveDash
| PlaylistItemType::Video
| PlaylistItemType::Series
2025-12-29 16:43:33 +01:00
| PlaylistItemType::LocalSeries
2025-07-24 14:27:57 +02:00
| PlaylistItemType::Catchup
) {
let _ = app_state
.active_users
.create_user_session(
user,
session_token,
virtual_id,
&provider,
&session_url,
2025-10-30 13:00:25 +01:00
&fingerprint.addr,
connection_permission,
)
.await;
2025-04-22 17:08:31 +02:00
}
}
2026-01-20 18:51:46 +01:00
let body_stream = prepare_body_stream(app_state, item_type, stream);
try_unwrap_body!(response.body(body_stream))
2025-03-19 09:54:59 +01:00
};
2025-01-24 19:01:53 +01:00
2025-03-19 09:54:59 +01:00
return stream_resp.into_response();
}
2025-11-10 01:04:13 +01:00
app_state.connection_manager.release_provider_handle(stream_details.provider_handle).await;
2026-01-20 10:09:58 +01:00
StatusCode::BAD_REQUEST.into_response()
}
2024-12-07 15:45:00 +01:00
2025-12-19 21:13:12 +01:00
fn get_stream_throttle(app_state: &Arc<AppState>) -> u64 {
2025-07-24 14:27:57 +02:00
app_state
.app_config
.config
.load()
2025-03-26 15:06:05 +01:00
.reverse_proxy
.as_ref()
.and_then(|reverse_proxy| reverse_proxy.stream.as_ref())
2025-07-24 14:27:57 +02:00
.map(|stream| stream.throttle_kbps)
.unwrap_or_default()
2025-03-26 15:06:05 +01:00
}
2025-11-20 13:00:01 +01:00
#[allow(clippy::too_many_arguments)]
2025-11-04 21:46:44 +01:00
async fn try_shared_stream_response_if_any(
2025-12-12 20:58:55 +01:00
app_state: &Arc<AppState>,
2025-07-24 14:27:57 +02:00
stream_url: &str,
2025-10-30 13:00:25 +01:00
fingerprint: &Fingerprint,
2025-07-24 14:27:57 +02:00
user: &ProxyUserCredentials,
connect_permission: UserConnectionPermission,
2025-10-25 20:40:58 +02:00
mut stream_channel: StreamChannel,
2025-11-20 13:00:01 +01:00
session_token: &str,
2025-10-25 20:40:58 +02:00
req_headers: &HeaderMap,
2025-07-24 14:27:57 +02:00
) -> Option<impl IntoResponse> {
2026-03-11 12:25:54 +01:00
if connect_permission == UserConnectionPermission::GracePeriod {
return None;
}
2025-10-27 08:27:35 +01:00
if let Some((stream, provider)) =
2025-11-01 00:26:30 +01:00
SharedStreamManager::subscribe_shared_stream(app_state, stream_url, &fingerprint.addr).await
2025-07-24 14:27:57 +02:00
{
2026-03-10 13:49:52 +01:00
debug_if_enabled!("Using shared stream {}", sanitize_sensitive_info(stream_url));
2026-02-14 17:52:35 +01:00
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 mut grace_period_options = app_state.get_grace_options();
if connect_permission != UserConnectionPermission::GracePeriod {
grace_period_options.period_millis = 0;
}
let mut stream_details = StreamDetails::from_stream(stream, grace_period_options);
2025-10-27 08:27:35 +01:00
stream_details.provider_name = provider;
2025-10-24 19:47:39 +02:00
stream_channel.shared = true;
2026-02-14 17:52:35 +01:00
let stream = create_active_client_stream(
stream_details,
app_state,
user,
connect_permission,
fingerprint,
stream_channel,
Some(session_token),
req_headers,
)
.await
.boxed();
2025-07-24 14:27:57 +02:00
let mut response = axum::response::Response::builder().status(status_code);
2025-03-11 14:56:10 +01:00
for (key, value) in &header_map {
response = response.header(key, value);
}
return response.body(axum::body::Body::from_stream(stream)).ok();
2024-12-07 15:45:00 +01:00
}
}
None
}
2024-12-09 16:43:19 +01:00
2025-12-16 16:23:15 +01:00
#[allow(clippy::too_many_arguments, clippy::too_many_lines)]
pub async fn local_stream_response(
fingerprint: &Fingerprint,
app_state: &Arc<AppState>,
pli: StreamChannel,
req_headers: &HeaderMap,
_input: &ConfigInput,
_target: &ConfigTarget,
_user: &ProxyUserCredentials,
connection_permission: UserConnectionPermission,
2026-01-20 10:09:58 +01:00
check_path: bool,
2025-12-16 16:23:15 +01:00
) -> impl IntoResponse + Send {
if log_enabled!(log::Level::Trace) {
trace!("Try to open stream {}", sanitize_sensitive_info(&pli.url));
}
if connection_permission == UserConnectionPermission::Exhausted {
return create_custom_video_stream_response(
app_state,
&fingerprint.addr,
CustomVideoStreamType::UserConnectionsExhausted,
2026-02-14 17:52:35 +01:00
)
.into_response();
2025-12-16 16:23:15 +01:00
}
let path = PathBuf::from(pli.url.strip_prefix("file://").unwrap_or(&pli.url));
2025-12-18 15:55:54 +01:00
// Canonicalize and validate the path
let path = match path.canonicalize() {
Ok(canonical) => canonical,
Err(err) => {
2026-01-20 10:09:58 +01:00
error!("Local file path is corrupt {}: {err}", path.display());
2025-12-18 17:25:13 +01:00
return StatusCode::NOT_FOUND.into_response();
}
2025-12-18 15:55:54 +01:00
};
2026-01-21 17:00:13 +01:00
if check_path {
2026-02-14 17:52:35 +01:00
let Some(library_paths) = app_state
.app_config
.config
.load()
.library
.as_ref()
.map(|lib| lib.scan_directories.iter().map(|dir| dir.path.clone()).collect::<Vec<_>>())
2026-01-21 17:00:13 +01:00
else {
return StatusCode::NOT_FOUND.into_response();
};
// Verify path is within allowed media directories
// (requires configuration of allowed base paths)
if !is_path_within_allowed_directories(&path, &library_paths) {
return StatusCode::FORBIDDEN.into_response();
}
2025-12-18 16:51:09 +01:00
}
2025-12-18 15:55:54 +01:00
2025-12-16 16:23:15 +01:00
let Ok(mut file) = tokio::fs::File::open(&path).await else { return StatusCode::NOT_FOUND.into_response() };
2026-01-20 10:09:58 +01:00
let Ok(metadata) = file.metadata().await else { return internal_server_error!() };
2025-12-16 16:23:15 +01:00
let file_size = metadata.len();
2026-02-14 17:52:35 +01:00
let range = req_headers.get("range").and_then(|v| v.to_str().ok()).and_then(parse_range);
2025-12-16 16:23:15 +01:00
2025-12-16 20:35:19 +01:00
let (start, end) = if let Some((req_start, req_end)) = range {
if file_size == 0 || req_start >= file_size {
return StatusCode::RANGE_NOT_SATISFIABLE.into_response();
}
let end = req_end.unwrap_or(file_size - 1).min(file_size - 1);
2025-12-16 20:46:09 +01:00
if end < req_start {
return StatusCode::RANGE_NOT_SATISFIABLE.into_response();
}
2025-12-16 20:35:19 +01:00
(req_start, end)
2025-12-16 16:23:15 +01:00
} else {
if file_size == 0 {
// Serve empty file
let body = axum::body::Body::empty();
let mut response = Response::new(body);
*response.status_mut() = StatusCode::OK;
let headers = response.headers_mut();
if let Some(ext) = get_file_extension(&pli.url) {
let ct = content_type_from_ext(&ext);
headers.insert(header::CONTENT_TYPE, HeaderValue::from_static(ct));
} else {
headers.insert(header::CONTENT_TYPE, HeaderValue::from_static("application/octet-stream"));
}
headers.insert("Accept-Ranges", HeaderValue::from_static("bytes"));
headers.insert(header::CONTENT_LENGTH, HeaderValue::from_static("0"));
return response.into_response();
}
(0, file_size - 1)
};
let content_length = end - start + 1;
if start > 0 {
if let Err(_err) = file.seek(SeekFrom::Start(start)).await {
2026-01-20 10:09:58 +01:00
return internal_server_error!();
2025-12-16 16:23:15 +01:00
}
}
let stream = ReaderStream::new(file.take(content_length));
2026-02-14 17:52:35 +01:00
let body_stream =
prepare_body_stream::<_>(app_state, pli.item_type, stream.map_err(|err| StreamError::Stream(err.to_string())));
2025-12-16 16:23:15 +01:00
2025-12-18 15:55:54 +01:00
let mut response = Response::new(body_stream);
2025-12-16 16:23:15 +01:00
2026-02-14 17:52:35 +01:00
*response.status_mut() = if range.is_some() { StatusCode::PARTIAL_CONTENT } else { StatusCode::OK };
2025-12-16 16:23:15 +01:00
let headers = response.headers_mut();
if let Some(ext) = get_file_extension(&pli.url) {
let ct = content_type_from_ext(&ext);
headers.insert(header::CONTENT_TYPE, HeaderValue::from_static(ct));
} else {
headers.insert(header::CONTENT_TYPE, HeaderValue::from_static("application/octet-stream"));
}
headers.insert("Accept-Ranges", HeaderValue::from_static("bytes"));
2026-01-04 09:55:50 +01:00
if let Ok(header_value) = HeaderValue::from_str(&content_length.to_string()) {
headers.insert(header::CONTENT_LENGTH, header_value);
}
2025-12-16 16:23:15 +01:00
if range.is_some() {
2026-01-04 09:55:50 +01:00
if let Ok(header_value) = HeaderValue::from_str(&format!("bytes {start}-{end}/{file_size}")) {
headers.insert(header::CONTENT_RANGE, header_value);
}
2025-12-16 16:23:15 +01:00
}
response
}
2025-12-18 17:25:13 +01:00
fn is_path_within_allowed_directories(sub_path: &Path, root_paths: &[String]) -> bool {
for root_path in root_paths {
2025-12-19 16:00:02 +01:00
if sub_path.starts_with(PathBuf::from(root_path)) {
2025-12-18 17:25:13 +01:00
return true;
}
}
false
2025-12-18 16:51:09 +01:00
}
2025-12-16 16:23:15 +01:00
2024-12-09 16:43:19 +01:00
pub fn is_stream_share_enabled(item_type: PlaylistItemType, target: &ConfigTarget) -> bool {
2025-07-24 14:27:57 +02:00
(item_type == PlaylistItemType::Live/* || item_type == PlaylistItemType::LiveHls */)
2026-02-14 17:52:35 +01:00
&& target.options.as_ref().is_some_and(|opt| opt.share_live_streams)
}
2025-03-11 14:56:10 +01:00
pub type HeaderFilter = Option<Box<dyn Fn(&str) -> bool + Send>>;
2026-02-14 17:52:35 +01:00
pub fn get_headers_from_request(req_headers: &HeaderMap, filter: &HeaderFilter) -> HashMap<String, Vec<u8>> {
2025-03-11 14:56:10 +01:00
req_headers
2025-01-08 21:17:57 +01:00
.iter()
.filter(|(k, _)| match &filter {
None => true,
2025-07-24 14:27:57 +02:00
Some(predicate) => predicate(k.as_str()),
})
2025-01-08 21:17:57 +01:00
.map(|(k, v)| (k.as_str().to_string(), v.as_bytes().to_vec()))
.collect()
}
2025-07-24 14:27:57 +02:00
fn get_add_cache_content(
res_url: &str,
2025-12-25 12:33:28 +01:00
mime_type: Option<String>,
2025-07-24 14:27:57 +02:00
cache: &Arc<ArcSwapOption<Mutex<LRUResourceCache>>>,
) -> Arc<dyn Fn(usize) + Send + Sync> {
let resource_url = String::from(res_url);
let cache = Arc::clone(cache);
2025-03-20 11:57:32 +01:00
let add_cache_content: Arc<dyn Fn(usize) + Send + Sync> = Arc::new(move |size| {
let res_url = resource_url.clone();
2025-12-25 12:33:28 +01:00
let mime_type = mime_type.clone();
2025-10-27 13:58:55 +01:00
// todo spawn, replace with unboundchannel
let cache = Arc::clone(&cache);
2025-03-11 14:56:10 +01:00
tokio::spawn(async move {
if let Some(cache) = cache.load().as_ref() {
2025-12-25 12:33:28 +01:00
let _ = cache.lock().await.add_content(&res_url, mime_type, size);
}
});
});
add_cache_content
}
2026-01-20 10:09:58 +01:00
fn get_mime_type(headers: &HeaderMap, resource_url: &str) -> Option<String> {
2025-12-25 12:33:28 +01:00
headers
.get(header::CONTENT_TYPE)
2026-02-14 17:52:35 +01:00
.and_then(|v| v.to_str().ok()) // Option<&str>
.map(ToString::to_string) // Option<String>
2025-12-25 12:33:28 +01:00
.or_else(|| {
2025-12-26 16:03:03 +01:00
// fallback to guess
2026-02-14 17:52:35 +01:00
mime_guess::from_path(resource_url).first_raw().map(ToString::to_string)
2025-12-25 12:33:28 +01:00
})
}
async fn build_resource_stream_response(
2025-12-19 21:13:12 +01:00
app_state: &Arc<AppState>,
resource_url: &str,
response: reqwest::Response,
) -> axum::response::Response {
2025-11-17 08:06:28 -06:00
let sanitized_resource_url = sanitize_sensitive_info(resource_url);
let status = response.status();
2026-02-14 17:52:35 +01:00
let mut response_builder = axum::response::Response::builder().status(status);
2025-12-25 12:33:28 +01:00
let mime_type = get_mime_type(response.headers(), resource_url);
2026-01-20 10:09:58 +01:00
let has_content_range = response.headers().contains_key(header::CONTENT_RANGE);
for (key, value) in response.headers() {
let name = key.as_str();
let is_hop_by_hop = matches!(
name.to_ascii_lowercase().as_str(),
"connection"
| "keep-alive"
| "proxy-authenticate"
| "proxy-authorization"
| "te"
| "trailer"
| "transfer-encoding"
| "upgrade"
);
if !is_hop_by_hop {
response_builder = response_builder.header(key, value);
}
}
if !response_builder.headers_ref().is_some_and(|h| h.contains_key(header::CACHE_CONTROL)) {
response_builder = response_builder.header(header::CACHE_CONTROL, "public, max-age=14400");
}
2026-02-14 17:52:35 +01:00
let byte_stream = response.bytes_stream().map_err(|err| StreamError::reqwest(&err));
// Cache only complete responses (200 OK without Content-Range)
2026-01-20 10:09:58 +01:00
let can_cache = status == StatusCode::OK && !has_content_range;
if can_cache {
2026-02-14 17:52:35 +01:00
debug!("Caching eligible resource stream {sanitized_resource_url}");
let cache_resource_path = if let Some(cache) = app_state.cache.load().as_ref() {
2025-12-25 12:33:28 +01:00
Some(cache.lock().await.store_path(resource_url, mime_type.as_deref()))
} else {
None
};
if let Some(resource_path) = cache_resource_path {
2025-11-17 08:06:28 -06:00
match create_new_file_for_write(&resource_path).await {
Ok(file) => {
2025-11-18 21:01:42 +01:00
debug!("Persisting resource stream {sanitized_resource_url} to {}", resource_path.display());
let writer = async_file_writer(file);
2025-12-25 12:33:28 +01:00
let add_cache_content = get_add_cache_content(resource_url, mime_type, &app_state.cache);
2025-12-10 17:19:57 +01:00
let tee = tee_stream(byte_stream, writer, &resource_path, add_cache_content);
return try_unwrap_body!(response_builder.body(axum::body::Body::from_stream(tee)));
2025-11-17 08:06:28 -06:00
}
Err(err) => {
2026-02-14 17:52:35 +01:00
warn!(
"Failed to create cache file {} for {sanitized_resource_url}: {err}",
resource_path.display()
);
2025-11-17 08:06:28 -06:00
}
}
2025-11-17 08:06:28 -06:00
} else {
2025-11-18 21:01:42 +01:00
debug!("Resource cache unavailable; streaming response for {sanitized_resource_url} without persistence");
}
}
try_unwrap_body!(response_builder.body(axum::body::Body::from_stream(byte_stream)))
}
async fn fetch_resource_with_retry(
2025-12-19 21:13:12 +01:00
app_state: &Arc<AppState>,
url: &Url,
resource_url: &str,
req_headers: &HashMap<String, Vec<u8>>,
input: Option<&ConfigInput>,
) -> Option<axum::response::Response> {
let config = app_state.app_config.config.load();
2026-01-20 18:51:46 +01:00
let default_user_agent = config.default_user_agent.clone();
drop(config);
2025-12-12 14:35:00 +01:00
let disabled_headers = app_state.get_disabled_headers();
2026-02-10 11:53:42 +01:00
let provider_config = input.and_then(|i| i.get_resolve_provider(url.as_str()));
2026-02-14 17:52:35 +01:00
let Ok(response) =
send_with_retry_and_provider(&app_state.app_config, url, provider_config.as_ref(), false, |resolved_url| {
request::get_client_request(
&app_state.http_client.load(),
input.map_or(InputFetchMethod::GET, |i| i.method),
input.map(|i| &i.headers),
2026-02-10 11:53:42 +01:00
resolved_url,
Some(req_headers),
disabled_headers.as_ref(),
default_user_agent.as_deref(),
)
2026-02-14 17:52:35 +01:00
})
.await
else {
return None;
};
let status = response.status();
if status.is_success() {
2026-02-14 17:52:35 +01:00
return Some(build_resource_stream_response(app_state, resource_url, response).await);
}
// Non-retriable Status → Upstream Response incl. Body
2026-02-14 17:52:35 +01:00
debug_if_enabled!("Failed to open resource got status {status} for {}", sanitize_sensitive_info(resource_url));
let mut response_builder = axum::response::Response::builder().status(status);
for (key, value) in response.headers() {
response_builder = response_builder.header(key, value);
}
2026-02-14 17:52:35 +01:00
let stream = response.bytes_stream().map_err(|err| StreamError::reqwest(&err));
2026-02-14 17:52:35 +01:00
Some(try_unwrap_body!(response_builder.body(axum::body::Body::from_stream(stream))))
}
2025-04-10 10:43:14 +02:00
/// # Panics
2025-07-24 14:27:57 +02:00
pub async fn resource_response(
2025-12-19 21:13:12 +01:00
app_state: &Arc<AppState>,
2025-07-24 14:27:57 +02:00
resource_url: &str,
req_headers: &HeaderMap,
input: Option<&ConfigInput>,
) -> impl IntoResponse + Send {
if resource_url.is_empty() {
2026-01-20 10:09:58 +01:00
return StatusCode::NO_CONTENT.into_response();
}
2026-02-14 17:52:35 +01:00
let filter: HeaderFilter = Some(Box::new(|key| key != "if-none-match" && key != "if-modified-since"));
2025-03-11 14:56:10 +01:00
let req_headers = get_headers_from_request(req_headers, &filter);
if let Some(cache) = app_state.cache.load().as_ref() {
2025-03-11 14:56:10 +01:00
let mut guard = cache.lock().await;
2025-12-25 12:33:28 +01:00
if let Some((resource_path, mime_type)) = guard.get_content(resource_url) {
2025-12-10 17:19:57 +01:00
trace_if_enabled!("Responding resource from cache {}", sanitize_sensitive_info(resource_url));
2026-02-14 17:52:35 +01:00
return serve_file(
&resource_path,
mime_type.unwrap_or_else(|| mime::APPLICATION_OCTET_STREAM.to_string()),
Some("public, max-age=14400"),
)
.await
.into_response();
2025-01-14 00:44:43 +01:00
}
}
2026-02-14 17:52:35 +01:00
trace_if_enabled!("Try to fetch resource {}", sanitize_sensitive_info(resource_url));
if let Ok(url) = Url::parse(resource_url) {
2026-02-14 17:52:35 +01:00
if let Some(resp) = fetch_resource_with_retry(app_state, &url, resource_url, &req_headers, input).await {
return resp;
}
// Upstream failure after retries
2026-01-20 10:09:58 +01:00
return StatusCode::BAD_GATEWAY.into_response();
}
error!("Url is malformed {}", sanitize_sensitive_info(resource_url));
2026-01-20 10:09:58 +01:00
StatusCode::BAD_REQUEST.into_response()
}
2025-01-31 23:08:23 +01:00
pub fn separate_number_and_remainder(input: &str) -> (String, Option<String>) {
2025-07-24 14:27:57 +02:00
input.rfind('.').map_or_else(
|| (input.to_string(), None),
|dot_index| {
let number_part = input[..dot_index].to_string();
let rest = input[dot_index..].to_string();
(number_part, if rest.len() < 2 { None } else { Some(rest) })
},
)
2025-01-31 23:08:23 +01:00
}
2025-04-10 10:43:14 +02:00
/// # Panics
2025-11-20 13:00:01 +01:00
pub fn empty_json_list_response() -> axum::response::Response {
try_unwrap_body!(axum::response::Response::builder()
2026-01-20 10:09:58 +01:00
.status(StatusCode::OK)
.header(header::CONTENT_TYPE, mime::APPLICATION_JSON.to_string())
2025-07-28 14:46:13 +02:00
.body("[]".to_owned()))
}
2025-07-24 14:27:57 +02:00
pub fn get_username_from_auth_header(token: &str, app_state: &Arc<AppState>) -> Option<String> {
2026-02-14 17:52:35 +01:00
if let Some(web_auth_config) = &app_state.app_config.config.load().web_ui.as_ref().and_then(|c| c.auth.as_ref()) {
2025-07-21 22:17:57 +02:00
let secret_key: &[u8] = web_auth_config.secret.as_ref();
2026-02-14 17:52:35 +01:00
if let Ok(token_data) =
decode::<Claims>(token, &DecodingKey::from_secret(secret_key), &Validation::new(Algorithm::HS256))
{
2025-03-11 14:56:10 +01:00
return Some(token_data.claims.username);
}
}
None
}
2025-03-11 14:56:10 +01:00
pub fn redirect(url: &str) -> impl IntoResponse {
try_unwrap_body!(axum::response::Response::builder()
2026-01-20 10:09:58 +01:00
.status(StatusCode::FOUND)
.header(header::LOCATION, url)
2026-02-01 13:37:49 +01:00
.body(Body::empty()))
2025-03-16 17:08:00 +01:00
}
2025-04-22 17:08:31 +02:00
2025-07-24 14:27:57 +02:00
pub async fn is_seek_request(cluster: XtreamCluster, req_headers: &HeaderMap) -> bool {
2025-04-24 18:12:29 +02:00
// seek only for non-live streams
if cluster == XtreamCluster::Live {
return false;
2025-04-24 18:57:18 +02:00
}
// seek requests contains range header
2026-02-14 17:52:35 +01:00
let range = req_headers.get("range").and_then(|h| h.to_str().ok()).map(ToString::to_string);
2025-04-24 18:12:29 +02:00
if let Some(range) = range {
if range.starts_with("bytes=") {
return true;
2025-04-28 19:52:59 +02:00
}
}
2025-06-04 21:40:14 +02:00
false
}
2025-09-29 15:39:27 +02:00
2025-09-30 10:54:30 +02:00
pub fn bin_response<T: Serialize>(data: &T) -> impl IntoResponse + Send {
match bin_serialize(data) {
2026-02-14 17:52:35 +01:00
Ok(body) => ([(header::CONTENT_TYPE, CONTENT_TYPE_CBOR)], body).into_response(),
2026-01-20 10:09:58 +01:00
Err(_) => internal_server_error!(),
2025-09-29 15:39:27 +02:00
}
}
pub fn json_response<T: Serialize>(data: &T) -> impl IntoResponse + Send {
2026-01-20 10:09:58 +01:00
(StatusCode::OK, axum::Json(data)).into_response()
2025-09-29 15:39:27 +02:00
}
2025-12-26 16:03:03 +01:00
pub fn json_or_bin_response<T: Serialize>(accept: Option<&str>, data: &T) -> impl IntoResponse + Send {
2026-01-22 14:07:25 +01:00
if accept.is_some_and(|a| a.contains(CONTENT_TYPE_CBOR)) {
2025-09-30 10:54:30 +02:00
return bin_response(data).into_response();
2025-09-29 15:39:27 +02:00
}
json_response(data).into_response()
2025-10-27 08:27:35 +01:00
}
2026-02-14 17:52:35 +01:00
pub fn stream_json_or_bin_response<P>(
accept: Option<&str>,
data: Box<dyn Iterator<Item = P> + Send>,
) -> axum::response::Response
2026-01-20 10:09:58 +01:00
where
P: serde::Serialize + Send + 'static,
{
2026-01-22 14:07:25 +01:00
if accept.is_some_and(|a| a.contains(CONTENT_TYPE_CBOR)) {
2026-01-20 10:09:58 +01:00
return stream_bin_array(data);
}
stream_json_array(data)
}
pub fn stream_json_or_bin_response_stream<P, S>(accept: Option<&str>, data: S) -> axum::response::Response
where
P: serde::Serialize + Send + 'static,
S: Stream<Item = P> + Send + Unpin + 'static,
{
if accept.is_some_and(|a| a.contains(CONTENT_TYPE_CBOR)) {
return stream_bin_array_stream(data);
}
stream_json_array_stream(data)
}
2026-01-20 10:09:58 +01:00
pub fn create_session_fingerprint(fingerprint: &Fingerprint, username: &str, virtual_id: u32) -> String {
concat_string!(&fingerprint.key, "|", username, "|", &virtual_id.to_string())
}
2026-02-14 17:52:35 +01:00
pub fn stream_json_array<P>(iter: Box<dyn Iterator<Item = P> + Send>) -> axum::response::Response
2026-01-20 10:09:58 +01:00
where
P: serde::Serialize + Send + 'static,
{
2026-02-14 17:52:35 +01:00
let stream = stream::unfold((iter, true), |(mut iter, first)| async move {
match iter.next() {
Some(item) => {
let mut json = String::new();
if !first {
json.push(',');
2026-01-20 10:09:58 +01:00
}
2026-02-14 17:52:35 +01:00
let element = serde_json::to_string(&item).ok()?;
json.push_str(&element);
Some((Ok::<Bytes, Infallible>(Bytes::from(json)), (iter, false)))
2026-01-20 10:09:58 +01:00
}
2026-02-14 17:52:35 +01:00
None => None,
}
});
2026-01-20 10:09:58 +01:00
let body = Body::from_stream(
stream::once(async { Ok::<_, Infallible>(Bytes::from_static(b"[")) })
.chain(stream)
2026-02-14 17:52:35 +01:00
.chain(stream::once(async { Ok::<_, Infallible>(Bytes::from_static(b"]")) })),
2026-01-20 10:09:58 +01:00
);
try_unwrap_body!(Response::builder().header(header::CONTENT_TYPE, CONTENT_TYPE_JSON).body(body))
2026-01-20 10:09:58 +01:00
}
2026-02-14 17:52:35 +01:00
pub fn stream_bin_array<P>(iter: Box<dyn Iterator<Item = P> + Send>) -> axum::response::Response
2026-01-20 10:09:58 +01:00
where
P: serde::Serialize + Send + 'static,
{
2026-02-14 17:52:35 +01:00
let stream = stream::unfold(iter, |mut iter| async move {
match iter.next() {
Some(item) => {
match bin_serialize(&item) {
Ok(buf) => Some((Ok::<Bytes, Infallible>(Bytes::from(buf)), iter)),
Err(err) => {
warn!("CBOR serialization error in stream: {err}");
Some((Ok::<Bytes, Infallible>(Bytes::new()), iter)) // skip errors, continue
2026-01-20 10:09:58 +01:00
}
}
}
2026-02-14 17:52:35 +01:00
None => None,
}
});
2026-01-20 10:09:58 +01:00
let body = Body::from_stream(
stream::once(async {
// CBOR: start indefinite-length array
Ok::<_, Infallible>(Bytes::from_static(&[0x9f]))
})
2026-02-14 17:52:35 +01:00
.chain(stream)
.chain(stream::once(async {
// CBOR: end indefinite-length array
Ok::<_, Infallible>(Bytes::from_static(&[0xff]))
})),
2026-01-20 10:09:58 +01:00
);
2026-02-14 17:52:35 +01:00
try_unwrap_body!(Response::builder().header(header::CONTENT_TYPE, CONTENT_TYPE_CBOR).body(body))
2026-01-20 10:09:58 +01:00
}
pub fn stream_json_array_stream<P, S>(stream: S) -> axum::response::Response
where
P: serde::Serialize + Send + 'static,
S: Stream<Item = P> + Send + Unpin + 'static,
{
let stream = stream::unfold((stream, true), |(mut stream, first)| async move {
match stream.next().await {
Some(item) => {
let mut json = String::new();
if !first {
json.push(',');
}
let element = serde_json::to_string(&item).ok()?;
json.push_str(&element);
Some((Ok::<Bytes, Infallible>(Bytes::from(json)), (stream, false)))
}
None => None,
}
});
let body = Body::from_stream(
stream::once(async { Ok::<_, Infallible>(Bytes::from_static(b"[")) })
.chain(stream)
.chain(stream::once(async { Ok::<_, Infallible>(Bytes::from_static(b"]")) })),
);
try_unwrap_body!(Response::builder().header(header::CONTENT_TYPE, CONTENT_TYPE_JSON).body(body))
}
pub fn stream_bin_array_stream<P, S>(stream: S) -> axum::response::Response
where
P: serde::Serialize + Send + 'static,
S: Stream<Item = P> + Send + Unpin + 'static,
{
let stream = stream::unfold(stream, |mut stream| async move {
match stream.next().await {
Some(item) => match bin_serialize(&item) {
Ok(buf) => Some((Ok::<Bytes, Infallible>(Bytes::from(buf)), stream)),
Err(err) => {
warn!("CBOR serialization error in stream: {err}");
Some((Ok::<Bytes, Infallible>(Bytes::new()), stream))
}
},
None => None,
}
});
let body = Body::from_stream(
stream::once(async { Ok::<_, Infallible>(Bytes::from_static(&[0x9f])) })
.chain(stream)
.chain(stream::once(async { Ok::<_, Infallible>(Bytes::from_static(&[0xff])) })),
);
try_unwrap_body!(Response::builder().header(header::CONTENT_TYPE, CONTENT_TYPE_CBOR).body(body))
}
2026-01-20 10:09:58 +01:00
pub fn create_api_proxy_user(app_state: &Arc<AppState>) -> ProxyUserCredentials {
let config = app_state.app_config.config.load();
let server = config
.web_ui
.as_ref()
.and_then(|web_ui| web_ui.player_server.as_ref())
.map_or("default", |server_name| server_name.as_str());
ProxyUserCredentials {
username: "api_user".to_string(),
password: "api_user".to_string(),
token: None,
proxy: ProxyType::Reverse(None),
server: Some(server.to_string()),
epg_timeshift: None,
epg_request_timeshift: None,
2026-01-20 10:09:58 +01:00
created_at: None,
exp_date: None,
max_connections: 0,
status: None,
ui_enabled: false,
comment: None,
2026-03-09 17:12:46 +01:00
priority: 0,
2026-03-05 00:04:53 +01:00
t_is_api_user: true,
2026-01-20 10:09:58 +01:00
}
}
2026-01-23 17:10:58 +01:00
pub fn empty_json_response_as_object() -> axum::http::Result<axum::response::Response> {
axum::response::Response::builder()
.status(axum::http::StatusCode::OK)
2026-02-14 17:52:35 +01:00
.header(axum::http::header::CONTENT_TYPE, mime::APPLICATION_JSON.to_string())
2026-01-23 17:10:58 +01:00
.body(axum::body::Body::from("{}".as_bytes()))
}
pub fn empty_json_response_as_array() -> axum::http::Result<axum::response::Response> {
axum::response::Response::builder()
.status(axum::http::StatusCode::OK)
2026-02-14 17:52:35 +01:00
.header(axum::http::header::CONTENT_TYPE, mime::APPLICATION_JSON.to_string())
2026-01-23 17:10:58 +01:00
.body(axum::body::Body::from("[]".as_bytes()))
}
#[cfg(test)]
mod tests {
use super::*;
use axum::http::HeaderMap;
use shared::model::XtreamCluster;
#[tokio::test]
async fn test_is_seek_request() {
let mut headers = HeaderMap::new();
// No range header
assert!(!is_seek_request(XtreamCluster::Video, &headers).await);
// Range: bytes=0- (Should be true now to allow session takeover on restart)
headers.insert("range", "bytes=0-".parse().unwrap());
assert!(is_seek_request(XtreamCluster::Video, &headers).await);
// Range: bytes=100- (Should be true)
headers.insert("range", "bytes=100-".parse().unwrap());
assert!(is_seek_request(XtreamCluster::Video, &headers).await);
// Range: bytes=100-200 (Should be true)
headers.insert("range", "bytes=100-200".parse().unwrap());
assert!(is_seek_request(XtreamCluster::Video, &headers).await);
// Live cluster should always return false
headers.insert("range", "bytes=100-".parse().unwrap());
assert!(!is_seek_request(XtreamCluster::Live, &headers).await);
}
}