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

836 lines
39 KiB
Rust
Raw Normal View History

2025-06-03 18:25:15 +02:00
use crate::api::endpoints::xtream_api::{get_xtream_player_api_stream_url, ApiStreamContext};
2025-04-27 13:26:21 +02:00
use crate::api::model::active_provider_manager::{ProviderAllocation, ProviderConnectionGuard};
2024-12-05 19:54:34 +01:00
use crate::api::model::app_state::AppState;
use crate::api::model::model_utils::get_stream_response_with_headers;
2025-03-20 11:57:32 +01:00
use crate::api::model::request::UserApiRequest;
2025-04-22 17:08:31 +02:00
use crate::api::model::stream::{BoxedProviderStream, ProviderStreamInfo, ProviderStreamResponse};
2025-03-20 11:57:32 +01:00
use crate::api::model::stream_error::StreamError;
use crate::api::model::streams::active_client_stream::ActiveClientStream;
2025-02-05 20:01:49 +01:00
use crate::api::model::streams::persist_pipe_stream::PersistPipeStream;
2025-05-09 17:13:37 +02:00
use crate::api::model::streams::provider_stream::{create_channel_unavailable_stream, create_custom_video_stream_response, create_provider_connections_exhausted_stream, CustomVideoStreamType};
use crate::api::model::streams::provider_stream_factory::{create_provider_stream, ProviderStreamFactoryOptions};
2025-03-20 11:57:32 +01:00
use crate::api::model::streams::shared_stream_manager::SharedStreamManager;
use crate::api::model::streams::throttled_stream::ThrottledStream;
2025-06-05 17:51:03 +02:00
use crate::auth::Claims;
use crate::model::ConfigInput;
2025-06-19 17:20:31 +02:00
use crate::model::{ConfigTarget, ProxyUserCredentials};
2025-04-22 17:08:31 +02:00
use crate::tools::atomic_once_flag::AtomicOnceFlag;
2025-02-05 20:01:49 +01:00
use crate::tools::lru_cache::LRUResourceCache;
2025-05-14 18:23:22 +02:00
use crate::utils::create_new_file_for_write;
2025-05-04 10:59:08 +02:00
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;
2025-03-20 11:57:32 +01:00
use axum::response::IntoResponse;
2025-04-22 17:08:31 +02:00
use chrono::{DateTime, Utc};
2025-03-11 14:56:10 +01:00
use futures::{StreamExt, TryStreamExt};
2025-03-20 11:57:32 +01:00
use jsonwebtoken::{decode, Algorithm, DecodingKey, Validation};
use log::{debug, error, log_enabled, trace};
use shared::model::{InputFetchMethod, PlaylistEntry, PlaylistItemType, TargetType, UserConnectionPermission, XtreamCluster};
2025-07-08 18:23:03 +02:00
use shared::utils::{default_grace_period_millis, human_readable_byte_size, trim_slash};
use shared::utils::{extract_extension_from_url, replace_url_extension, sanitize_sensitive_info, DASH_EXT, HLS_EXT};
2025-04-22 17:08:31 +02:00
use std::borrow::Cow;
2024-12-05 19:54:34 +01:00
use std::collections::HashMap;
2025-02-11 19:20:51 +01:00
use std::io::BufWriter;
2025-01-27 15:58:57 +01:00
use std::path::Path;
2025-03-20 11:57:32 +01:00
use std::sync::Arc;
2025-03-11 14:56:10 +01:00
use tokio::sync::Mutex;
2024-12-05 19:54:34 +01:00
use url::Url;
use crate::api::model::active_user_manager::UserSession;
use crate::api::model::provider_config::ProviderConfig;
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 => {
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
}
};
}
#[macro_export]
macro_rules! try_result_bad_request {
($option:expr, $msg_is_error:expr, $msg:expr) => {
match $option {
Ok(value) => value,
Err(_) => {
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 {
Ok(value) => value,
2025-03-11 14:56:10 +01:00
Err(_) => return axum::http::StatusCode::BAD_REQUEST.into_response(),
2025-01-31 23:08:23 +01:00
}
};
}
#[macro_export]
macro_rules! try_unwrap_body {
($body:expr) => {
$body.map_or_else(
|_| axum::http::StatusCode::INTERNAL_SERVER_ERROR.into_response(),
|resp| resp.into_response(),
)
};
}
2025-02-05 20:01:49 +01:00
pub use try_option_bad_request;
pub use try_result_bad_request;
pub use try_unwrap_body;
2025-04-22 00:21:45 +02:00
pub fn get_server_time() -> String {
chrono::offset::Local::now().with_timezone(&chrono::Local).format("%Y-%m-%d %H:%M:%S %Z").to_string()
}
pub fn get_build_time() -> Option<String> {
BUILD_TIMESTAMP.to_string().parse::<DateTime<Utc>>().ok().map(|datetime| datetime.format("%Y-%m-%d %H:%M:%S %Z").to_string())
}
pub fn get_memory_usage() -> String {
2025-05-04 10:59:08 +02:00
crate::utils::get_memory_usage().map_or(String::from("?"), human_readable_byte_size)
2025-04-22 00:21:45 +02:00
}
2025-02-05 20:01:49 +01:00
2025-04-10 00:06:48 +02:00
#[allow(clippy::missing_panics_doc)]
2025-07-15 12:28:25 +02:00
pub async fn serve_file(file_path: &Path, mime_type: mime::Mime) -> impl IntoResponse + Send {
if file_path.exists() {
2025-03-11 14:56:10 +01:00
return match tokio::fs::File::open(file_path).await {
Ok(file) => {
let reader = tokio::io::BufReader::new(file);
let stream = tokio_util::io::ReaderStream::new(reader);
let body = axum::body::Body::from_stream(stream);
try_unwrap_body!(axum::response::Response::builder()
2025-07-15 12:28:25 +02:00
.status(axum::http::StatusCode::OK)
2025-03-11 14:56:10 +01:00
.header(axum::http::header::CONTENT_TYPE, mime_type.to_string())
.header(axum::http::header::CACHE_CONTROL, axum::http::header::HeaderValue::from_static("no-cache"))
.body(body))
2025-03-11 14:56:10 +01:00
}
Err(_) => axum::http::StatusCode::INTERNAL_SERVER_ERROR.into_response(),
};
}
2025-03-11 14:56:10 +01:00
axum::http::StatusCode::NOT_FOUND.into_response()
}
2025-06-24 01:06:25 +02:00
pub fn get_user_target_by_username(username: &str, app_state: &AppState) -> 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
}
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-06-24 01:06:25 +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 stream_force_retry_secs: u32,
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,
/// - `stream_force_retry_secs`: the number of seconds to wait before a forced retry,
/// - `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`
/// - forced retry interval: `0`
/// - 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-03-16 17:08:00 +01:00
fn get_stream_options(app_state: &AppState) -> StreamOptions {
2025-03-21 17:53:21 +01:00
let (stream_retry, stream_force_retry_secs, buffer_enabled, buffer_size) = app_state
2025-06-24 12:31:38 +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-03-21 17:53:21 +01:00
.map_or((false, 0, false, 0), |stream| {
2025-01-14 10:47:37 +01:00
let (buffer_enabled, buffer_size) = stream
.buffer
.as_ref()
.map_or((false, 0), |buffer| (buffer.enabled, buffer.size));
2025-03-21 17:53:21 +01:00
(stream.retry, stream.forced_retry_interval_secs, 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;
2025-03-21 17:53:21 +01:00
StreamOptions { stream_retry, stream_force_retry_secs, buffer_enabled, buffer_size, pipe_provider_stream }
2025-01-27 15:58:57 +01:00
}
2025-03-11 14:56:10 +01:00
// fn get_stream_content_length(provider_response: Option<&(Vec<(String, String)>, StatusCode)>) -> u64 {
// let content_length = provider_response
// .as_ref()
// .and_then(|(headers, _)| headers.iter().find(|(h, _)| h.eq(axum::http::header::CONTENT_LENGTH.as_str())))
// .and_then(|(_, val)| val.parse::<u64>().ok())
// .unwrap_or(0);
// content_length
// }
2025-01-27 15:58:57 +01:00
2025-04-24 18:12:29 +02:00
pub fn get_stream_alternative_url(stream_url: &str, input: &ConfigInput, alias_input: &Arc<ProviderConfig>) -> String {
2025-03-19 09:54:59 +01:00
let Some(input_user_info) = input.get_user_info() else { return stream_url.to_owned() };
let Some(alt_input_user_info) = alias_input.get_user_info() else { return stream_url.to_owned() };
let modified = stream_url.replace(&input_user_info.base_url, &alt_input_user_info.base_url);
let modified = modified.replace(&input_user_info.username, &alt_input_user_info.username);
2025-04-10 00:06:48 +02:00
modified.replace(&input_user_info.password, &alt_input_user_info.password)
2025-03-19 09:54:59 +01:00
}
async fn get_redirect_alternative_url<'a>(app_state: &AppState, redirect_url: &'a str, input: &ConfigInput) -> Cow<'a, str> {
if let Some((base_url, username, password)) = input.get_matched_config_by_url(redirect_url) {
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) {
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);
return Cow::Owned(new_url);
}
// one has credentials the other not, something not right
return Cow::Borrowed(redirect_url);
}
return Cow::Owned(new_url);
}
}
Cow::Borrowed(redirect_url)
}
2025-03-20 11:57:32 +01:00
type StreamUrl = String;
type ProviderName = String;
2025-05-09 17:13:37 +02:00
enum ProviderStreamState {
2025-04-10 10:43:14 +02:00
Custom(ProviderStreamResponse),
Available(Option<ProviderName>, StreamUrl),
GracePeriod(Option<ProviderName>, StreamUrl),
2025-03-20 11:57:32 +01:00
}
pub struct StreamDetails {
pub stream: Option<BoxedProviderStream>,
stream_info: ProviderStreamInfo,
pub input_name: Option<String>,
2025-03-27 12:19:07 +01:00
pub grace_period_millis: u64,
2025-03-20 11:57:32 +01:00
pub reconnect_flag: Option<Arc<AtomicOnceFlag>>,
2025-07-21 22:17:57 +02:00
pub provider_connection_guard: Option<Arc<ProviderConnectionGuard>>,
2025-03-20 11:57:32 +01:00
}
impl StreamDetails {
pub fn from_stream(stream: BoxedProviderStream) -> Self {
Self {
stream: Some(stream),
stream_info: None,
input_name: None,
2025-03-27 12:19:07 +01:00
grace_period_millis: default_grace_period_millis(),
2025-03-20 11:57:32 +01:00
reconnect_flag: None,
2025-04-16 20:28:02 +02:00
provider_connection_guard: None,
2025-03-20 11:57:32 +01:00
}
}
#[inline]
2025-03-20 11:57:32 +01:00
pub fn has_stream(&self) -> bool {
self.stream.is_some()
}
#[inline]
pub fn has_grace_period(&self) -> bool {
self.grace_period_millis > 0
}
2025-03-20 11:57:32 +01:00
}
2025-05-09 17:13:37 +02:00
struct StreamingStrategy {
2025-07-21 22:17:57 +02:00
provider_connection_guard: Option<Arc<ProviderConnectionGuard>>,
2025-05-09 17:13:37 +02:00
provider_stream_state: ProviderStreamState,
input_headers: Option<HashMap<String, String>>,
}
/// 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-19 15:32:30 +02:00
async fn resolve_streaming_strategy(app_state: &AppState, stream_url: &str, addr: &str, input: &ConfigInput, force_provider: Option<&str>)
2025-05-09 17:13:37 +02:00
-> StreamingStrategy {
// allocate a provider connection
2025-04-24 18:12:29 +02:00
let provider_connection_guard = match force_provider {
2025-07-19 15:32:30 +02:00
Some(provider) => app_state.active_provider.force_exact_acquire_connection(provider, addr).await,
None => app_state.active_provider.acquire_connection(&input.name, addr).await
2025-04-24 18:12:29 +02:00
};
2025-07-21 22:17:57 +02:00
let stream_response_params = match &**provider_connection_guard {
2025-04-24 18:12:29 +02:00
ProviderAllocation::Exhausted => {
debug!("Input {} is exhausted. No connections allowed.", input.name);
2025-06-24 12:31:38 +02:00
let stream = create_provider_connections_exhausted_stream(&app_state.app_config, &[]);
2025-05-09 17:13:37 +02:00
ProviderStreamState::Custom(stream)
2025-04-24 18:12:29 +02:00
}
2025-07-19 15:32:30 +02:00
ProviderAllocation::Available(_, ref provider)
| ProviderAllocation::GracePeriod(_, ref provider) => {
2025-04-24 18:12:29 +02:00
// 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
2025-04-24 18:12:29 +02:00
let (provider, url) = if force_provider.is_some() || provider.id == input.id {
(input.name.to_string(), stream_url.to_string())
} else {
(provider.name.to_string(), get_stream_alternative_url(stream_url, input, provider))
};
2025-03-20 11:57:32 +01:00
2025-07-21 22:17:57 +02:00
if matches!(&**provider_connection_guard, ProviderAllocation::Available(_, _)) {
2025-05-09 17:13:37 +02:00
ProviderStreamState::Available(Some(provider), url)
2025-04-24 18:12:29 +02:00
} else {
2025-05-09 17:13:37 +02:00
ProviderStreamState::GracePeriod(Some(provider), url)
2025-03-19 09:54:59 +01:00
}
2025-04-24 18:12:29 +02:00
}
};
2025-05-09 17:13:37 +02:00
StreamingStrategy {
provider_connection_guard: Some(provider_connection_guard),
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-04-17 12:11:51 +02:00
2025-05-26 13:57:29 +02:00
fn get_grace_period_millis(connection_permission: UserConnectionPermission, stream_response_params: &ProviderStreamState, config_grace_period_millis: u64) -> u64 {
2025-04-17 12:11:51 +02:00
if config_grace_period_millis > 0 &&
2025-05-09 17:13:37 +02:00
(matches!(stream_response_params, ProviderStreamState::GracePeriod(_, _)) // provider grace period
2025-05-26 13:57:29 +02:00
|| connection_permission == UserConnectionPermission::GracePeriod // user grace period
2025-04-17 12:11:51 +02:00
) { config_grace_period_millis } else { 0 }
}
2025-04-16 18:14:49 +02:00
#[allow(clippy::too_many_arguments)]
2025-04-24 18:12:29 +02:00
async fn create_stream_response_details(app_state: &AppState,
stream_options: &StreamOptions,
stream_url: &str,
2025-07-19 15:32:30 +02:00
addr: &str,
2025-04-24 18:12:29 +02:00
req_headers: &HeaderMap,
input: &ConfigInput,
item_type: PlaylistItemType,
share_stream: bool,
2025-04-22 17:08:31 +02:00
connection_permission: UserConnectionPermission,
force_provider: Option<&str>) -> StreamDetails {
2025-05-09 17:13:37 +02:00
let mut streaming_strategy =
2025-07-19 15:32:30 +02:00
resolve_streaming_strategy(app_state, stream_url, addr, input, force_provider).await;
2025-06-24 12:31:38 +02:00
let config_grace_period_millis = app_state.app_config.config.load().reverse_proxy.as_ref()
2025-04-22 17:08:31 +02:00
.and_then(|r| r.stream.as_ref()).map_or_else(default_grace_period_millis, |s| s.grace_period_millis);
2025-05-26 13:57:29 +02:00
let grace_period_millis = get_grace_period_millis(connection_permission, &streaming_strategy.provider_stream_state, config_grace_period_millis);
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;
StreamDetails {
stream,
stream_info,
input_name: None,
2025-03-27 12:19:07 +01:00
grace_period_millis,
2025-03-20 11:57:32 +01:00
reconnect_flag: None,
2025-07-21 22:17:57 +02:00
provider_connection_guard: streaming_strategy.provider_connection_guard.clone(),
2025-03-20 11:57:32 +01:00
}
}
2025-05-09 17:13:37 +02:00
ProviderStreamState::Available(provider_name, request_url) |
ProviderStreamState::GracePeriod(provider_name, request_url) => {
2025-03-20 11:57:32 +01:00
let parsed_url = Url::parse(&request_url);
let ((stream, stream_info), reconnect_flag) = if let Ok(url) = parsed_url {
2025-05-09 17:13:37 +02:00
let provider_stream_factory_options = ProviderStreamFactoryOptions::new(item_type, share_stream, stream_options, &url, req_headers, streaming_strategy.input_headers.as_ref());
let reconnect_flag = provider_stream_factory_options.get_reconnect_flag_clone();
let provider_stream = match create_provider_stream(Arc::clone(&app_state.app_config), Arc::clone(&app_state.http_client.load()), provider_stream_factory_options).await {
2025-05-09 17:13:37 +02:00
None => (None, None),
Some((stream, info)) => {
(Some(stream), info)
}
};
(provider_stream, Some(reconnect_flag))
2025-03-20 11:57:32 +01:00
} else {
((None, None), None)
};
2025-03-19 09:54:59 +01:00
// if we have no stream, we should release the provider
2025-03-20 11:57:32 +01:00
if stream.is_none() {
2025-05-09 17:13:37 +02:00
if let Some(guard) = streaming_strategy.provider_connection_guard.take() {
drop(guard);
}
2025-04-22 17:08:31 +02:00
error!("Cant open stream {}", sanitize_sensitive_info(&request_url));
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-06-03 18:25:15 +02:00
if let Some((headers, status_code, response_url)) = stream_info.as_ref() {
2025-03-20 11:57:32 +01:00
debug!(
"Responding stream request {} with status {}, headers {:?}",
2025-06-03 18:25:15 +02: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
2025-03-20 11:57:32 +01:00
StreamDetails {
stream,
stream_info,
input_name: provider_name,
2025-03-27 12:19:07 +01:00
grace_period_millis,
2025-03-20 11:57:32 +01:00
reconnect_flag,
2025-05-09 17:13:37 +02:00
provider_connection_guard: streaming_strategy.provider_connection_guard.take(),
2025-03-20 11:57:32 +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 {
let extension = self.stream_ext.map_or_else(
|| extract_extension_from_url(url).map_or_else(String::new, std::string::ToString::to_string),
std::string::ToString::to_string);
// 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() {
format!("{provider_id}{extension}")
} else {
2025-07-08 18:23:03 +02:00
format!("{}/{provider_id}{extension}", trim_slash(self.action_path))
2025-04-19 15:01:23 +02:00
}
}
}
2025-04-19 13:50:41 +02:00
2025-04-19 15:01:23 +02:00
pub async fn redirect_response<'a, P>(app_state: &AppState, params: &'a RedirectParams<'a, P>) -> Option<impl IntoResponse + Send>
where
P: PlaylistEntry,
{
2025-04-19 13:50:41 +02:00
let item_type = params.item.get_item_type();
let provider_url = &params.item.get_provider_url();
let redirect_request = params.user.proxy.is_redirect(item_type) || params.target.is_force_redirect(item_type);
2025-04-19 13:50:41 +02:00
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);
if params.target_type == TargetType::M3u {
if redirect_request || is_dash_request {
let redirect_url = if is_hls_request { &replace_url_extension(provider_url, HLS_EXT) } else { provider_url };
let redirect_url = if is_dash_request { &replace_url_extension(redirect_url, DASH_EXT) } else { redirect_url };
2025-04-24 18:12:29 +02:00
let redirect_url = get_redirect_alternative_url(app_state, redirect_url, params.input).await;
debug_if_enabled!("Redirecting stream request to {}", sanitize_sensitive_info(&redirect_url));
return Some(redirect(&redirect_url).into_response());
2025-04-19 13:50:41 +02:00
}
} else if params.target_type == TargetType::Xtream {
let Some(provider_id) = params.provider_id else {
2025-07-15 12:28:25 +02:00
return Some(axum::http::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();
let username = params.input.username.as_ref().map_or("", |v| v);
let password = params.input.password.as_ref().map_or("", |v| v);
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}");
2025-04-20 22:53:01 +02:00
debug_if_enabled!("Redirecting stream request to {}", sanitize_sensitive_info(&stream_url));
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-05-26 13:57:29 +02:00
let stream_url = match get_xtream_player_api_stream_url(params.input, params.req_context, &params.get_query_path(provider_id, provider_url), provider_url) {
2025-04-19 15:01:23 +02:00
None => {
error!("Cant find stream url for target {target_name}, context {}, stream_id {virtual_id}", params.req_context);
2025-07-15 12:28:25 +02:00
return Some(axum::http::StatusCode::BAD_REQUEST.into_response());
2025-04-19 15:01:23 +02:00
}
Some(url) => {
2025-04-24 18:12:29 +02:00
match app_state.active_provider.get_next_provider(&params.input.name).await {
Some(provider_cfg) => get_stream_alternative_url(&url, params.input, &provider_cfg),
2025-04-19 15:01:23 +02:00
None => url,
}
}
};
2025-04-19 13:50:41 +02:00
2025-04-19 15:01:23 +02:00
// hls or dash redirect
if is_dash_request {
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 {}", sanitize_sensitive_info(redirect_url));
return Some(redirect(redirect_url).into_response());
2025-04-19 13:50:41 +02:00
}
debug_if_enabled!("Redirecting stream request to {}", sanitize_sensitive_info(&stream_url));
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-06-03 18:25:15 +02:00
throttle_kbps > 0 && matches!(item_type, PlaylistItemType::Video | PlaylistItemType::Series | PlaylistItemType::SeriesInfo | PlaylistItemType::Catchup)
2025-04-22 17:08:31 +02:00
}
2025-07-15 12:28:25 +02:00
fn prepare_body_stream(app_state: &AppState, item_type: PlaylistItemType, stream: ActiveClientStream) -> axum::body::Body {
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) {
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
pub async fn force_provider_stream_response(addr: &str,
app_state: &AppState,
user_session: &UserSession,
2025-04-24 18:12:29 +02:00
item_type: PlaylistItemType,
req_headers: &HeaderMap,
input: &ConfigInput,
user: &ProxyUserCredentials) -> impl IntoResponse + Send {
let stream_options = get_stream_options(app_state);
let share_stream = false;
let connection_permission = UserConnectionPermission::Allowed;
2025-04-22 17:08:31 +02:00
let mut stream_details =
2025-07-19 15:32:30 +02:00
create_stream_response_details(app_state, &stream_options, &user_session.stream_url, addr, req_headers, input, item_type, share_stream, connection_permission, Some(&user_session.provider)).await;
2025-04-22 17:08:31 +02:00
if stream_details.has_stream() {
let provider_response = stream_details.stream_info.as_ref().map(|(h, sc, url)| (h.clone(), *sc, url.clone()));
2025-07-19 15:32:30 +02:00
app_state.active_users.update_session_addr(&user.username, &user_session.token, addr);
let stream = ActiveClientStream::new(stream_details, app_state, user, connection_permission, addr);
2025-04-22 17:08:31 +02: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
}
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 try_unwrap_body!(response.body(body_stream));
2025-04-22 17:08:31 +02:00
}
drop(stream_details.provider_connection_guard.take());
2025-05-09 17:13:37 +02:00
if let (Some(stream), _stream_info) =
2025-07-15 12:28:25 +02:00
create_channel_unavailable_stream(&app_state.app_config, &[], axum::http::StatusCode::BAD_GATEWAY)
2025-05-09 17:13:37 +02:00
{
debug!("Streaming custom stream");
try_unwrap_body!(axum::response::Response::builder().status(axum::http::StatusCode::OK).body(axum::body::Body::from_stream(stream)))
2025-05-09 17:13:37 +02:00
} else {
2025-07-15 12:28:25 +02:00
axum::http::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-04-16 18:14:49 +02:00
#[allow(clippy::too_many_arguments)]
pub async fn stream_response(addr: &str,
app_state: &AppState, session_token: &str, virtual_id: u32,
item_type: PlaylistItemType, stream_url: &str, req_headers: &HeaderMap,
input: &ConfigInput, target: &ConfigTarget, user: &ProxyUserCredentials,
connection_permission: UserConnectionPermission) -> impl IntoResponse + Send {
2025-01-27 15:58:57 +01:00
if log_enabled!(log::Level::Trace) { trace!("Try to open stream {}", sanitize_sensitive_info(stream_url)); }
2025-04-16 18:04:57 +02:00
if connection_permission == UserConnectionPermission::Exhausted {
2025-06-24 12:31:38 +02:00
return create_custom_video_stream_response(&app_state.app_config, CustomVideoStreamType::UserConnectionsExhausted).into_response();
2025-04-16 18:04:57 +02:00
}
2025-01-27 15:58:57 +01:00
let share_stream = is_stream_share_enabled(item_type, target);
if share_stream {
2025-07-19 15:32:30 +02:00
if let Some(value) = shared_stream_response(app_state, stream_url, addr, user, connection_permission) {
2025-03-11 14:56:10 +01:00
return value.into_response();
2025-01-27 15:58:57 +01:00
}
}
2025-01-14 10:47:37 +01:00
2025-03-19 09:54:59 +01:00
let stream_options = get_stream_options(app_state);
2025-04-16 20:28:02 +02:00
let mut stream_details =
2025-07-19 15:32:30 +02:00
create_stream_response_details(app_state, &stream_options, stream_url, addr, req_headers, input, item_type, share_stream, connection_permission, None).await;
2025-03-20 11:57:32 +01:00
if stream_details.has_stream() {
2025-03-19 09:54:59 +01:00
// let content_length = get_stream_content_length(provider_response.as_ref());
2025-06-03 18:25:15 +02:00
let provider_response = stream_details.stream_info.as_ref().map(|(h, sc, response_url)| (h.clone(), *sc, response_url.clone()));
2025-07-21 22:17:57 +02:00
let provider_name = stream_details.provider_connection_guard.as_ref().and_then(|guard| guard.get_provider_name());
2025-04-22 17:08:31 +02:00
2025-07-19 15:32:30 +02:00
let provider_guard = if share_stream { stream_details.provider_connection_guard.take() } else { None };
let stream = ActiveClientStream::new(stream_details, app_state, user, connection_permission, addr);
2025-03-19 09:54:59 +01:00
let stream_resp = if share_stream {
2025-04-24 18:12:29 +02:00
debug_if_enabled!("Streaming shared stream request from {}", sanitize_sensitive_info(stream_url));
2025-03-19 09:54:59 +01:00
// Shared Stream response
2025-06-03 18:25:15 +02:00
let shared_headers = provider_response.as_ref().map_or_else(Vec::new, |(h, _, _)| h.clone());
SharedStreamManager::subscribe(app_state, stream_url, stream, shared_headers, stream_options.buffer_size, provider_guard);
if let Some(broadcast_stream) = SharedStreamManager::subscribe_shared_stream(app_state, stream_url, Some(addr)) {
let (status_code, header_map) = get_stream_response_with_headers(provider_response.map(|(h, s, _)| (h, s)));
2025-03-11 14:56:10 +01:00
let mut response = axum::response::Response::builder()
.status(status_code);
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 {
axum::http::StatusCode::BAD_REQUEST.into_response()
}
} else {
let session_url = provider_response.as_ref().and_then(|(_, _, u)| u.as_ref()).map_or_else(|| Cow::Borrowed(stream_url), |url| Cow::Owned(url.to_string()));
2025-06-04 21:40:14 +02:00
if log_enabled!(log::Level::Debug) {
if session_url.eq(&stream_url) {
debug!("Streaming stream request from {}", sanitize_sensitive_info(stream_url));
} else {
debug!("Streaming stream request for {} from {}", sanitize_sensitive_info(stream_url), sanitize_sensitive_info(&session_url));
}
}
let (status_code, header_map) = 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-06-03 18:25:15 +02:00
if matches!(item_type, PlaylistItemType::LiveHls | PlaylistItemType::LiveDash | PlaylistItemType::Video | PlaylistItemType::Series | PlaylistItemType::Catchup) {
2025-07-19 15:32:30 +02:00
let _ = app_state.active_users.create_user_session(user, session_token, virtual_id, &provider, &session_url, addr, connection_permission);
2025-04-22 17:08:31 +02: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-04-16 20:28:02 +02:00
drop(stream_details.provider_connection_guard.take());
2025-03-11 14:56:10 +01:00
axum::http::StatusCode::BAD_REQUEST.into_response()
}
2024-12-07 15:45:00 +01:00
2025-03-28 17:31:24 +01:00
fn get_stream_throttle(app_state: &AppState) -> u64 {
2025-06-24 12:31:38 +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-03-28 17:31:24 +01:00
.map(|stream| stream.throttle_kbps).unwrap_or_default()
2025-03-26 15:06:05 +01:00
}
2025-07-19 15:32:30 +02:00
fn shared_stream_response(app_state: &AppState, stream_url: &str, addr: &str, user: &ProxyUserCredentials, connect_permission: UserConnectionPermission) -> Option<impl IntoResponse> {
if let Some(stream) = SharedStreamManager::subscribe_shared_stream(app_state, stream_url, Some(addr)) {
2025-04-24 18:12:29 +02:00
debug_if_enabled!("Using shared stream {}", sanitize_sensitive_info(stream_url));
2025-07-19 15:32:30 +02:00
if let Some(headers) = app_state.shared_stream_manager.get_shared_state_headers(stream_url) {
2025-07-15 12:28:25 +02:00
let (status_code, header_map) = get_stream_response_with_headers(Some((headers.clone(), axum::http::StatusCode::OK)));
2025-03-20 11:57:32 +01:00
let stream_details = StreamDetails::from_stream(stream);
2025-07-19 15:32:30 +02:00
let stream = ActiveClientStream::new(stream_details, app_state, user, connect_permission, addr).boxed();
2025-03-11 14:56:10 +01:00
let mut response = axum::response::Response::builder()
.status(status_code);
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
pub fn is_stream_share_enabled(item_type: PlaylistItemType, target: &ConfigTarget) -> bool {
2025-03-19 09:54:59 +01:00
(item_type == PlaylistItemType::Live /* || item_type == PlaylistItemType::LiveHls */) && 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>>;
pub fn get_headers_from_request(req_headers: &HeaderMap, filter: &HeaderFilter) -> HashMap<String, Vec<u8>> {
req_headers
2025-01-08 21:17:57 +01:00
.iter()
.filter(|(k, _)| match &filter {
None => true,
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()
}
fn get_add_cache_content(res_url: &str, 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();
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-03-11 14:56:10 +01:00
let _ = cache.lock().await.add_content(&res_url, size);
}
});
});
add_cache_content
}
2025-04-10 10:43:14 +02:00
/// # Panics
2025-07-15 12:28:25 +02:00
pub async fn resource_response(app_state: &AppState, resource_url: &str, req_headers: &HeaderMap, input: Option<&ConfigInput>) -> impl IntoResponse + Send {
if resource_url.is_empty() {
2025-03-11 14:56:10 +01:00
return axum::http::StatusCode::NO_CONTENT.into_response();
}
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-02-06 10:49:34 +01:00
if let Some(resource_path) = guard.get_content(resource_url) {
2025-03-11 14:56:10 +01:00
trace_if_enabled!("Responding resource from cache {}", sanitize_sensitive_info(resource_url));
return serve_file(&resource_path, mime::APPLICATION_OCTET_STREAM).await.into_response();
2025-01-14 00:44:43 +01:00
}
}
trace_if_enabled!("Try to fetch resource {}", sanitize_sensitive_info(resource_url));
if let Ok(url) = Url::parse(resource_url) {
let client = request::get_client_request(&app_state.http_client.load(), input.map_or(InputFetchMethod::GET, |i| i.method), input.map(|i| &i.headers), &url, Some(&req_headers));
match client.send().await {
Ok(response) => {
let status = response.status();
if status.is_success() {
2025-03-11 14:56:10 +01:00
let mut response_builder = axum::response::Response::builder()
2025-07-15 12:28:25 +02:00
.status(axum::http::StatusCode::OK);
2025-03-11 14:56:10 +01:00
for (key, value) in response.headers() {
response_builder = response_builder.header(key, value);
}
2025-01-14 00:44:43 +01:00
2025-01-27 15:58:57 +01:00
let byte_stream = response.bytes_stream().map_err(|err| StreamError::reqwest(&err));
let cache_resource_path = {
if let Some(cache) = app_state.cache.load().as_ref() {
Some(cache.lock().await.store_path(resource_url))
} else {
None
}
};
if let Some(resource_path) = cache_resource_path {
if let Ok(file) = create_new_file_for_write(&resource_path) {
2025-02-11 19:20:51 +01:00
let writer = BufWriter::new(file);
let add_cache_content = get_add_cache_content(resource_url, &app_state.cache);
2025-01-14 00:44:43 +01:00
let stream = PersistPipeStream::new(byte_stream, writer, add_cache_content);
return try_unwrap_body!(response_builder.body(axum::body::Body::from_stream(stream)));
2025-01-14 00:44:43 +01:00
}
}
return try_unwrap_body!(response_builder.body(axum::body::Body::from_stream(byte_stream)));
}
debug_if_enabled!("Failed to open resource got status {} for {}", status, sanitize_sensitive_info(resource_url));
}
Err(err) => {
error!("Received failure from server {}: {}", sanitize_sensitive_info(resource_url), err);
}
}
} else {
error!("Url is malformed {}", sanitize_sensitive_info(resource_url));
}
2025-03-11 14:56:10 +01:00
axum::http::StatusCode::BAD_REQUEST.into_response()
}
2025-01-31 23:08:23 +01:00
pub fn separate_number_and_remainder(input: &str) -> (String, Option<String>) {
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-04-10 10:43:14 +02:00
/// # Panics
2025-07-15 12:28:25 +02:00
pub fn empty_json_list_response() -> impl IntoResponse + Send {
try_unwrap_body!(axum::response::Response::builder()
2025-07-15 12:28:25 +02:00
.status(axum::http::StatusCode::OK)
2025-03-11 14:56:10 +01:00
.header("Content-Type", mime::APPLICATION_JSON.to_string())
.body("[]".to_string()))
}
2025-03-11 14:56:10 +01:00
pub fn get_username_from_auth_header(
token: &str,
app_state: &Arc<AppState>,
) -> Option<String> {
2025-06-24 12:31:38 +02: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();
2025-03-11 14:56:10 +01:00
if let Ok(token_data) = decode::<Claims>(
token,
2025-07-21 22:17:57 +02:00
&DecodingKey::from_secret(secret_key),
2025-03-11 14:56:10 +01:00
&Validation::new(Algorithm::HS256),
) {
return Some(token_data.claims.username);
}
}
None
}
2025-03-11 14:56:10 +01:00
2025-04-10 10:43:14 +02:00
/// # Panics
2025-03-11 14:56:10 +01:00
pub fn redirect(url: &str) -> impl IntoResponse {
try_unwrap_body!(axum::response::Response::builder()
2025-07-15 12:28:25 +02:00
.status(axum::http::StatusCode::FOUND)
2025-03-11 14:56:10 +01:00
.header("Location", url)
.body(axum::body::Body::empty()))
2025-03-16 17:08:00 +01:00
}
2025-04-22 17:08:31 +02:00
2025-04-30 11:04:01 +02:00
pub async fn is_seek_request(
2025-04-24 18:12:29 +02:00
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
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=0-") {
return false;
2025-04-22 17:08:31 +02:00
}
if range.starts_with("bytes=") {
return true;
2025-04-28 19:52:59 +02:00
}
}
2025-06-04 21:40:14 +02:00
false
}