mirror of
https://github.com/euzu/tuliprox.git
synced 2026-10-04 06:52:26 +02:00
reconnect stream wip
This commit is contained in:
+25
-102
@@ -1,27 +1,24 @@
|
||||
use crate::api::model::app_state::AppState;
|
||||
use crate::api::model::request::UserApiRequest;
|
||||
use crate::api::model::shared_stream::SharedStream;
|
||||
use crate::model::api_proxy::{ProxyUserCredentials};
|
||||
use crate::debug_if_enabled;
|
||||
use crate::model::api_proxy::ProxyUserCredentials;
|
||||
use crate::model::config::{ConfigInput, ConfigTarget};
|
||||
use crate::model::playlist::PlaylistItemType;
|
||||
use crate::utils::request_utils;
|
||||
use crate::utils::request_utils::mask_sensitive_info;
|
||||
use actix_web::http::header::DATE;
|
||||
use actix_web::http::header::{HeaderValue, CACHE_CONTROL};
|
||||
use actix_web::http::header::{DATE};
|
||||
use actix_web::{HttpRequest, HttpResponse};
|
||||
use async_std::prelude::Stream;
|
||||
use async_std::stream::StreamExt;
|
||||
use bytes::Bytes;
|
||||
use chrono::Utc;
|
||||
use log::{debug, error};
|
||||
use log::{ error};
|
||||
use std::collections::HashMap;
|
||||
use std::path::Path;
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
use tokio::sync::broadcast;
|
||||
use tokio_stream::wrappers::BroadcastStream;
|
||||
use tokio_stream::wrappers::{BroadcastStream};
|
||||
use url::Url;
|
||||
use crate::debug_if_enabled;
|
||||
use crate::model::playlist::PlaylistItemType;
|
||||
use crate::api::model::buffered_stream;
|
||||
use crate::api::model::buffered_stream::get_stream_response_with_headers;
|
||||
|
||||
pub async fn serve_file(file_path: &Path, req: &HttpRequest, mime_type: mime::Mime) -> HttpResponse {
|
||||
if file_path.exists() {
|
||||
@@ -73,61 +70,8 @@ async fn create_notify_stream(
|
||||
}
|
||||
}
|
||||
|
||||
/// 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
|
||||
S: Stream<Item=Result<Bytes, E>> + Unpin + 'static,
|
||||
{
|
||||
// Create a broadcast channel for the shared stream
|
||||
let (tx, _) = broadcast::channel(100);
|
||||
let sender = Arc::new(tx);
|
||||
|
||||
// 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,
|
||||
},
|
||||
);
|
||||
|
||||
let shared_streams_map = Arc::clone(&app_state.shared_streams);
|
||||
let mut source_stream = Box::pin(bytes_stream);
|
||||
let streaming_url = stream_url.to_string();
|
||||
// Spawn a task to forward items from the source stream to the broadcast channel
|
||||
actix_rt::spawn(async move {
|
||||
while let Some(item) = source_stream.next().await {
|
||||
if let Ok(data) = item {
|
||||
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 {
|
||||
debug_if_enabled!("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;
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
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();
|
||||
debug_if_enabled!("Try to open stream {}", mask_sensitive_info(stream_url));
|
||||
if share_stream {
|
||||
if let Some(value) = shared_stream_response(app_state, stream_url).await {
|
||||
@@ -136,55 +80,34 @@ pub async fn stream_response(app_state: &AppState, stream_url: &str, req: &HttpR
|
||||
}
|
||||
|
||||
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) => {
|
||||
let status = response.status();
|
||||
if status.is_success() {
|
||||
let mut response_builder = HttpResponse::Ok();
|
||||
let mut header = HashMap::new();
|
||||
response.headers().iter().for_each(|(k, v)| {
|
||||
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()));
|
||||
});
|
||||
if share_stream {
|
||||
create_shared_stream(app_state, stream_url, header, response.bytes_stream()).await;
|
||||
if let Some(stream) = create_notify_stream(app_state, stream_url).await {
|
||||
debug!("Creating shared channel {stream_url}");
|
||||
return response_builder.body(actix_web::body::BodyStream::new(stream));
|
||||
}
|
||||
} else {
|
||||
return response_builder.body(actix_web::body::BodyStream::new(response.bytes_stream()));
|
||||
}
|
||||
}
|
||||
debug_if_enabled!("Failed to open stream got status {} for {}", status, mask_sensitive_info(stream_url));
|
||||
}
|
||||
Err(err) => {
|
||||
error!("Received failure from server {}: {}", mask_sensitive_info(stream_url), err);
|
||||
let mut buffered_stream_handler = buffered_stream::BufferedStreamHandler::new(&url, req, input);
|
||||
let stream = buffered_stream_handler.get_stream();
|
||||
return if share_stream {
|
||||
SharedStream::register(app_state, stream_url, stream).await;
|
||||
if let Some(broadcast_stream) = create_notify_stream(app_state, stream_url).await {
|
||||
let body_stream = actix_web::body::BodyStream::new(broadcast_stream);
|
||||
let mut response_builder = get_stream_response_with_headers();
|
||||
response_builder.body(body_stream)
|
||||
} else {
|
||||
HttpResponse::BadRequest().finish()
|
||||
}
|
||||
} else {
|
||||
let mut response_builder = get_stream_response_with_headers();
|
||||
let body_stream = actix_web::body::BodyStream::new(stream);
|
||||
response_builder.body(body_stream)
|
||||
}
|
||||
} else {
|
||||
error!("Url is malformed {}", mask_sensitive_info(stream_url));
|
||||
}
|
||||
error!("Url is malformed {}", mask_sensitive_info(stream_url));
|
||||
HttpResponse::BadRequest().finish()
|
||||
}
|
||||
|
||||
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 {
|
||||
debug_if_enabled!("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[..]));
|
||||
}
|
||||
if app_state.shared_streams.lock().await.get(stream_url).is_some() {
|
||||
let mut response_builder = get_stream_response_with_headers();
|
||||
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)));
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user