Files
tuliprox/src/api/api_utils.rs
T

233 lines
10 KiB
Rust
Raw Normal View History

2024-12-05 19:54:34 +01:00
use crate::api::model::app_state::AppState;
2025-01-14 00:44:43 +01:00
use crate::api::model::buffered_stream;
2025-01-14 10:47:37 +01:00
use crate::api::model::buffered_stream::{get_provider_stream, get_stream_response_with_headers};
2024-12-05 19:54:34 +01:00
use crate::api::model::request::UserApiRequest;
2024-12-07 15:45:00 +01:00
use crate::api::model::shared_stream::SharedStream;
2025-01-07 14:00:13 +01:00
use crate::debug_if_enabled;
use crate::model::api_proxy::ProxyUserCredentials;
2024-12-30 12:40:22 +01:00
use crate::model::config::{ConfigInput, ConfigTarget};
2025-01-07 14:00:13 +01:00
use crate::model::playlist::PlaylistItemType;
use crate::utils::request_utils;
use crate::utils::request_utils::mask_sensitive_info;
2025-01-14 00:44:43 +01:00
use actix_files::NamedFile;
use actix_web::body::{BodyStream};
2025-01-07 14:00:13 +01:00
use actix_web::http::header::DATE;
2024-12-05 19:54:34 +01:00
use actix_web::http::header::{HeaderValue, CACHE_CONTROL};
use actix_web::{HttpRequest, HttpResponse};
2024-12-06 17:05:53 +01:00
use bytes::Bytes;
2024-12-07 15:45:00 +01:00
use chrono::Utc;
2025-01-10 13:07:35 +01:00
use log::{error, log_enabled, trace};
2024-12-05 19:54:34 +01:00
use std::collections::HashMap;
2025-01-14 00:44:43 +01:00
use std::path::{Path};
use std::sync::Arc;
use async_std::sync::Mutex;
use tokio_stream::wrappers::BroadcastStream;
2024-12-05 19:54:34 +01:00
use url::Url;
2025-01-14 00:44:43 +01:00
use crate::api::model::persist_pipe_stream::PersistPipeStream;
use crate::utils::file_utils::create_new_file_for_write;
use crate::utils::lru_cache::LRUResourceCache;
2024-11-04 18:46:56 +01:00
pub async fn serve_file(file_path: &Path, req: &HttpRequest, mime_type: mime::Mime) -> HttpResponse {
if file_path.exists() {
2023-12-01 21:38:23 +01:00
if let Ok(file) = actix_files::NamedFile::open_async(file_path).await {
let mut result = file.set_content_type(mime_type)
2023-12-01 21:38:23 +01:00
.disable_content_disposition().into_response(req);
let headers = result.headers_mut();
2024-11-04 18:46:56 +01:00
headers.insert(CACHE_CONTROL, HeaderValue::from_bytes(b"no-cache").unwrap());
2023-12-01 21:38:23 +01:00
return result;
}
}
2023-12-01 21:38:23 +01:00
HttpResponse::NoContent().finish()
}
2024-11-04 18:46:56 +01:00
pub fn get_user_target_by_credentials<'a>(username: &str, password: &str, api_req: &'a UserApiRequest,
2024-12-05 19:54:34 +01:00
app_state: &'a AppState) -> Option<(ProxyUserCredentials, &'a ConfigTarget)> {
if !username.is_empty() && !password.is_empty() {
2023-12-01 21:38:23 +01:00
app_state.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 {
app_state.config.get_target_for_user_by_token(token)
}
}
2023-11-22 21:06:58 +01:00
}
2024-12-03 14:54:42 +01:00
pub fn get_user_target<'a>(api_req: &'a UserApiRequest, app_state: &'a AppState) -> Option<(ProxyUserCredentials, &'a 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-01-14 10:47:37 +01:00
/// Creates a broadcast notify stream for the given URL if a shared stream exists.
async fn create_broadcast_stream(
2024-12-06 17:51:33 +01:00
app_state: &AppState,
stream_url: &str,
2024-12-07 15:45:00 +01:00
) -> Option<BroadcastStream<Bytes>> {
2024-12-06 17:05:53 +01:00
let notify_stream_url = stream_url.to_string();
2024-12-06 17:51:33 +01:00
// Acquire lock and check for existing stream
let shared_streams = app_state.shared_streams.lock().await;
2024-12-06 17:05:53 +01:00
if let Some(shared_stream) = shared_streams.get(&notify_stream_url) {
let rx = shared_stream.data_stream.subscribe();
2024-12-07 15:45:00 +01:00
Some(BroadcastStream::new(rx))
} else {
None
2024-12-06 17:05:53 +01:00
}
}
2025-01-14 10:47:37 +01:00
pub async fn stream_response(app_state: &AppState, stream_url: &str,
req: &HttpRequest, input: Option<&ConfigInput>,
item_type: PlaylistItemType, target: &ConfigTarget) -> HttpResponse {
2025-01-14 00:44:43 +01:00
if log_enabled!(log::Level::Trace) { trace!("Try to open stream {}", mask_sensitive_info(stream_url)); }
2025-01-08 21:17:57 +01:00
2025-01-07 14:34:33 +01:00
let share_stream = is_stream_share_enabled(item_type, target);
2024-12-05 19:57:13 +01:00
if share_stream {
2025-01-09 17:59:34 +01:00
if let Some(value) = shared_stream_response(app_state, stream_url, None).await {
2024-12-07 15:45:00 +01:00
return value;
2024-12-05 19:54:34 +01:00
}
}
2025-01-14 10:47:37 +01:00
let (stream_retry, buffer_enabled, buffer_size) = app_state
.config
.reverse_proxy
.as_ref()
.and_then(|reverse_proxy| reverse_proxy.stream.as_ref())
.map_or((false, false, 0), |stream| {
let (buffer_enabled, buffer_size) = stream
.buffer
.as_ref()
.map_or((false, 0), |buffer| (buffer.enabled, buffer.size));
(stream.retry, buffer_enabled, buffer_size)
});
if let Ok(url) = Url::parse(stream_url) {
2025-01-14 10:47:37 +01:00
let direct_pipe_provider_stream = !stream_retry && !buffer_enabled;
let (stream_opt, provider_response) = if direct_pipe_provider_stream {
get_provider_stream(&app_state.http_client, &url, req, input).await
2025-01-07 14:00:13 +01:00
} else {
2025-01-14 10:47:37 +01:00
let buffer_stream_options = (item_type, stream_retry, buffer_enabled, buffer_size);
buffered_stream::get_buffered_stream(&app_state.http_client, &url, req, input, buffer_stream_options).await
2025-01-14 00:44:43 +01:00
};
2025-01-14 10:47:37 +01:00
if let Some(stream) = stream_opt {
let use_buffer = !buffer_enabled || direct_pipe_provider_stream;
return if share_stream {
SharedStream::register(app_state, stream_url, stream, use_buffer).await;
if let Some(broadcast_stream) = create_broadcast_stream(app_state, stream_url).await {
let body_stream = BodyStream::new(broadcast_stream);
let mut response_builder = get_stream_response_with_headers(provider_response, stream_url);
response_builder.body(body_stream)
} else {
HttpResponse::BadRequest().finish()
}
} else {
let mut response_builder = get_stream_response_with_headers(provider_response, stream_url);
response_builder.streaming(stream)
};
}
}
2025-01-07 14:00:13 +01:00
error!("Url is malformed {}", mask_sensitive_info(stream_url));
HttpResponse::BadRequest().finish()
}
2024-12-07 15:45:00 +01:00
2025-01-10 13:07:35 +01:00
async fn shared_stream_response(app_state: &AppState, stream_url: &str, headers: Option<(Vec<(String, String)>, reqwest::StatusCode)>) -> Option<HttpResponse> {
2025-01-14 10:47:37 +01:00
if let Some(stream) = create_broadcast_stream(app_state, stream_url).await {
2024-12-09 19:26:39 +01:00
debug_if_enabled!("Using shared channel {}", mask_sensitive_info(stream_url));
2025-01-07 14:00:13 +01:00
if app_state.shared_streams.lock().await.get(stream_url).is_some() {
let mut response_builder = get_stream_response_with_headers(headers, stream_url);
2024-12-07 15:45:00 +01:00
let current_date = Utc::now().format("%a, %d %b %Y %H:%M:%S GMT").to_string();
response_builder.insert_header((DATE, current_date.as_bytes()));
// response_builder.insert_header((ACCEPT_RANGES, "bytes".as_bytes()));
2025-01-14 00:44:43 +01:00
return Some(response_builder.body(BodyStream::new(stream)));
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 {
2024-12-13 00:13:14 +01:00
item_type == PlaylistItemType::Live && target.options.as_ref().is_some_and(|opt| opt.share_live_streams)
}
pub type HeaderFilter = Option<Box<dyn Fn(&str) -> bool>>;
pub fn get_headers_from_request(req: &HttpRequest, filter: &HeaderFilter) -> HashMap<String, Vec<u8>> {
2025-01-08 21:17:57 +01:00
req.headers()
.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<Option<Mutex<LRUResourceCache>>>) -> Box<dyn Fn(usize)> {
let resource_url = String::from(res_url);
let cache = Arc::clone(cache);
let add_cache_content: Box<dyn Fn(usize)> = Box::new(move|size| {
let res_url = resource_url.clone();
let cache = Arc::clone(&cache);
actix_rt::spawn(async move {
if let Some(cache) = cache.as_ref() {
let mut guard = cache.lock().await;
let _ = guard.add_content(&res_url, size).await;
}
});
});
add_cache_content
}
2025-01-09 15:58:38 +01:00
pub async fn resource_response(app_state: &AppState, resource_url: &str, req: &HttpRequest, input: Option<&ConfigInput>) -> HttpResponse {
if resource_url.is_empty() {
return HttpResponse::NoContent().finish();
}
let filter: HeaderFilter = Some(Box::new(|key| key != "if-none-match" && key != "if-modified-since"));
let req_headers = get_headers_from_request(req, &filter);
2025-01-14 00:44:43 +01:00
if let Some(cache) = app_state.cache.as_ref() {
let mut guard = cache.lock().await;
if let Some(resource_path) = guard.get_content(resource_url).await {
if let Ok(named_file) = NamedFile::open_async(resource_path).await {
debug_if_enabled!("Cached resource {}", mask_sensitive_info(resource_url));
return named_file.into_response(req);
}
}
}
debug_if_enabled!("Try to fetch resource {}", mask_sensitive_info(resource_url));
if let Ok(url) = Url::parse(resource_url) {
2025-01-09 15:58:38 +01:00
let client = request_utils::get_client_request(&app_state.http_client, input.map(|i| &i.headers), &url, Some(&req_headers));
match client.send().await {
Ok(response) => {
let status = response.status();
if status.is_success() {
let mut response_builder = HttpResponse::Ok();
response.headers().iter().for_each(|(k, v)| {
response_builder.insert_header((k.as_str(), v.as_ref()));
});
2025-01-14 00:44:43 +01:00
let byte_stream = response.bytes_stream();
if let Some(cache) = app_state.cache.as_ref() {
let resource_path = {
let guard = cache.lock().await;
guard.store_path(resource_url)
};
if let Ok(file) = create_new_file_for_write(&resource_path) {
let writer = Arc::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 response_builder.body(BodyStream::new(stream));
}
}
return response_builder.body(BodyStream::new(byte_stream));
}
debug_if_enabled!("Failed to open resource got status {} for {}", status, mask_sensitive_info(resource_url));
}
Err(err) => {
error!("Received failure from server {}: {}", mask_sensitive_info(resource_url), err);
}
}
} else {
error!("Url is malformed {}", mask_sensitive_info(resource_url));
}
HttpResponse::BadRequest().finish()
}