Files
tuliprox/src/api/api_utils.rs
T

211 lines
8.9 KiB
Rust
Raw Normal View History

2024-12-05 19:54:34 +01:00
use crate::api::model::app_state::AppState;
use crate::api::model::request::UserApiRequest;
2024-12-07 15:45:00 +01:00
use crate::api::model::shared_stream::SharedStream;
2024-05-03 01:47:41 +02:00
use crate::model::api_proxy::{ApiProxyServerInfo, ProxyUserCredentials};
2024-12-05 19:54:34 +01:00
use crate::model::config::{Config, ConfigInput, ConfigTarget};
use crate::utils::request_utils;
use crate::utils::request_utils::mask_sensitive_info;
2024-12-05 19:54:34 +01:00
use actix_web::http::header::{HeaderValue, CACHE_CONTROL};
2024-12-07 15:45:00 +01:00
use actix_web::http::header::{DATE};
2024-12-05 19:54:34 +01:00
use actix_web::{HttpRequest, HttpResponse};
2024-12-06 17:05:53 +01:00
use async_std::prelude::Stream;
use async_std::stream::StreamExt;
use bytes::Bytes;
2024-12-07 15:45:00 +01:00
use chrono::Utc;
2024-12-05 19:54:34 +01:00
use log::{debug, error, log_enabled, Level};
use std::collections::HashMap;
use std::path::Path;
2024-12-06 17:05:53 +01:00
use std::sync::Arc;
2024-12-07 15:45:00 +01:00
use std::time::Duration;
2024-12-06 17:51:33 +01:00
use tokio::sync::broadcast;
2024-12-06 17:05:53 +01:00
use tokio_stream::wrappers::BroadcastStream;
2024-12-05 19:54:34 +01:00
use url::Url;
2024-12-09 16:43:19 +01:00
use crate::model::playlist::PlaylistItemType;
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)
}
2024-11-04 18:46:56 +01:00
pub fn get_user_server_info(cfg: &Config, user: &ProxyUserCredentials) -> ApiProxyServerInfo {
2024-05-10 12:01:48 +02:00
let server_info_list = cfg.t_api_proxy.read().unwrap().as_ref().unwrap().server.clone();
2024-11-04 18:46:56 +01:00
let server_info_name = user.server.as_ref().map_or("default", |server_name| server_name.as_str());
server_info_list.iter().find(|c| c.name.eq(server_info_name)).map_or_else(|| server_info_list.first().unwrap().clone(), std::clone::Clone::clone)
}
2024-12-06 17:51:33 +01:00
/// Creates a notify stream for the given URL if a shared stream exists.
async fn create_notify_stream(
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
}
}
2024-12-06 17:51:33 +01:00
/// Creates a shared stream and stores it in the shared state.
async fn create_shared_stream<S, E>(
app_state: &AppState,
stream_url: &str,
header: HashMap<String, Vec<u8>>,
bytes_stream: S,
) where
2024-12-07 15:45:00 +01:00
S: Stream<Item=Result<Bytes, E>> + Unpin + 'static,
2024-12-06 17:05:53 +01:00
{
2024-12-06 17:51:33 +01:00
// Create a broadcast channel for the shared stream
2024-12-07 15:45:00 +01:00
let (tx, _) = broadcast::channel(100);
2024-12-06 17:05:53 +01:00
let sender = Arc::new(tx);
2024-12-06 17:51:33 +01:00
// Insert the shared stream into the shared state
app_state
.shared_streams
.lock()
.await
.insert(
stream_url.to_string(),
SharedStream {
data_stream: sender.clone(),
header,
},
);
2024-12-07 15:45:00 +01:00
let shared_streams_map = Arc::clone(&app_state.shared_streams);
2024-12-06 17:05:53 +01:00
let mut source_stream = Box::pin(bytes_stream);
2024-12-07 15:45:00 +01:00
let streaming_url = stream_url.to_string();
// Spawn a task to forward items from the source stream to the broadcast channel
2024-12-06 17:05:53 +01:00
actix_rt::spawn(async move {
2024-12-06 17:51:33 +01:00
while let Some(item) = source_stream.next().await {
if let Ok(data) = item {
2024-12-07 15:45:00 +01:00
if sender.receiver_count() > 0 {
// if let Err(err) = sender.send(data) {
// debug!("{err}")
// }
if sender.send(data).is_err() {
// ignore
}
actix_web::rt::time::sleep(Duration::from_millis(20)).await;
} else {
if log_enabled!(Level::Debug) {
debug!("No active subscribers. Closing stream {}", mask_sensitive_info(&streaming_url));
}
// Cleanup for removing unused shared streams
let mut shared_streams = shared_streams_map.lock().await;
shared_streams.remove(&streaming_url);
return;
2024-12-06 17:51:33 +01:00
}
2024-12-06 17:05:53 +01:00
}
}
});
}
2024-12-05 19:57:13 +01:00
pub async fn stream_response(app_state: &AppState, stream_url: &str, req: &HttpRequest, input: Option<&ConfigInput>, share_stream: bool) -> HttpResponse {
let req_headers: HashMap<&str, &[u8]> = req.headers().iter().map(|(k, v)| (k.as_str(), v.as_bytes())).collect();
2024-12-03 14:49:23 +01:00
if log_enabled!(Level::Debug) {
debug!("Try to open stream {}", mask_sensitive_info(stream_url));
}
2024-12-05 19:57:13 +01:00
if share_stream {
2024-12-07 15:45:00 +01:00
if let Some(value) = shared_stream_response(app_state, stream_url).await {
return value;
2024-12-05 19:54:34 +01:00
}
}
if let Ok(url) = Url::parse(stream_url) {
let client = request_utils::get_client_request(input, &url, Some(&req_headers));
match client.send().await {
Ok(response) => {
2024-12-06 17:05:53 +01:00
let status = response.status();
if status.is_success() {
let mut response_builder = HttpResponse::Ok();
2024-12-06 17:51:33 +01:00
let mut header = HashMap::new();
response.headers().iter().for_each(|(k, v)| {
2024-12-07 15:45:00 +01:00
let key = k.as_str().to_string();
// ignore date it is dynamic
if !"date".eq(key.to_lowercase().as_str()) {
header.insert(k.as_str().to_string(), v.as_bytes().to_vec());
}
response_builder.insert_header((k.as_str(), v.as_ref()));
});
2024-12-06 17:05:53 +01:00
if share_stream {
2024-12-06 17:51:33 +01:00
create_shared_stream(app_state, stream_url, header, response.bytes_stream()).await;
2024-12-06 17:05:53 +01:00
if let Some(stream) = create_notify_stream(app_state, stream_url).await {
2024-12-06 17:51:33 +01:00
debug!("Creating shared channel {stream_url}");
2024-12-06 17:05:53 +01:00
return response_builder.body(actix_web::body::BodyStream::new(stream));
}
2024-12-05 19:57:13 +01:00
} else {
2024-12-06 17:05:53 +01:00
return response_builder.body(actix_web::body::BodyStream::new(response.bytes_stream()));
}
}
2024-12-03 14:49:23 +01:00
if log_enabled!(Level::Debug) {
2024-12-06 17:05:53 +01:00
debug!("Failed to open stream got status {} for {}", status, mask_sensitive_info(stream_url));
2024-12-03 14:49:23 +01:00
}
}
Err(err) => {
error!("Received failure from server {}: {}", mask_sensitive_info(stream_url), err);
}
}
} else {
error!("Url is malformed {}", mask_sensitive_info(stream_url));
}
HttpResponse::BadRequest().finish()
}
2024-12-07 15:45:00 +01:00
async fn shared_stream_response(app_state: &AppState, stream_url: &str) -> Option<HttpResponse> {
if let Some(stream) = create_notify_stream(app_state, stream_url).await {
if log_enabled!(Level::Debug) {
debug!("Using shared channel {}", mask_sensitive_info(stream_url));
}
// return HttpResponse::Ok().body(actix_web::body::BodyStream::new(stream));
if let Some(shared_stream) = app_state.shared_streams.lock().await.get(stream_url) {
let mut response_builder = HttpResponse::Ok();
for (key, value) in &shared_stream.header {
response_builder.insert_header((key.as_str(), &value[..]));
}
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((CACHE_CONTROL, "no-cache".as_bytes()));
// response_builder.insert_header((ACCEPT_RANGES, "bytes".as_bytes()));
return Some(response_builder.body(actix_web::body::BodyStream::new(stream)));
}
}
None
}
2024-12-09 16:43:19 +01:00
pub fn is_stream_share_enabled(item_type: PlaylistItemType, target: &ConfigTarget) -> bool {
item_type == PlaylistItemType::Live && target.options.as_ref().map_or(false, |opt| opt.share_live_streams)
}