mirror of
https://github.com/euzu/tuliprox.git
synced 2026-10-10 09:52:18 +02:00
Merge branch 'develop' into feature/source_editor
This commit is contained in:
@@ -34,6 +34,7 @@
|
||||
- Fixed race conditions during simultaneous access by the same user.
|
||||
- Added extended debug logging for client requests and ID chain (request/action/virtual) to trace stream resolution.
|
||||
- Fixed xtream series/catchup lookups using the series-info virtual_id so episode requests now keep their own virtual_id/session.
|
||||
- Made cache storage more robust. Incomplete downloads will be deleted from cache.
|
||||
|
||||
# 3.2.0 (2025-11-14)
|
||||
- Added `name` attribute to Staged Input.
|
||||
|
||||
Generated
+3
-3
@@ -1096,7 +1096,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "frontend"
|
||||
version = "3.2.16"
|
||||
version = "3.2.17"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"base64",
|
||||
@@ -3793,7 +3793,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "shared"
|
||||
version = "3.2.16"
|
||||
version = "3.2.17"
|
||||
dependencies = [
|
||||
"base64",
|
||||
"bitflags 2.10.0",
|
||||
@@ -4356,7 +4356,7 @@ checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b"
|
||||
|
||||
[[package]]
|
||||
name = "tuliprox"
|
||||
version = "3.2.16"
|
||||
version = "3.2.17"
|
||||
dependencies = [
|
||||
"arc-swap",
|
||||
"async-compression",
|
||||
|
||||
+1
-1
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "tuliprox"
|
||||
version = "3.2.16"
|
||||
version = "3.2.17"
|
||||
edition = "2021"
|
||||
|
||||
# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html
|
||||
|
||||
@@ -1,10 +1,10 @@
|
||||
use crate::api::endpoints::xtream_api::{get_xtream_player_api_stream_url, ApiStreamContext};
|
||||
use crate::api::model::{UserSession};
|
||||
use crate::api::model::{tee_stream, UserSession};
|
||||
use crate::api::model::{
|
||||
create_channel_unavailable_stream, create_custom_video_stream_response,
|
||||
create_provider_connections_exhausted_stream, create_provider_stream,
|
||||
get_stream_response_with_headers, ActiveClientStream, AppState,
|
||||
CustomVideoStreamType, PersistPipeStream,
|
||||
CustomVideoStreamType,
|
||||
ProviderStreamFactoryOptions, SharedStreamManager,
|
||||
StreamError, ThrottledStream, UserApiRequest,
|
||||
};
|
||||
@@ -1109,8 +1109,8 @@ async fn build_stream_response(
|
||||
debug!("Persisting resource stream {sanitized_resource_url} to {}", resource_path.display());
|
||||
let writer = BufWriter::new(file);
|
||||
let add_cache_content = get_add_cache_content(resource_url, &app_state.cache);
|
||||
let stream = PersistPipeStream::new(byte_stream, writer, add_cache_content);
|
||||
return try_unwrap_body!(response_builder.body(axum::body::Body::from_stream(stream)));
|
||||
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)));
|
||||
}
|
||||
Err(err) => {
|
||||
warn!("Failed to create cache file {} for {sanitized_resource_url}: {err}", resource_path.display());
|
||||
@@ -1247,10 +1247,7 @@ pub async fn resource_response(
|
||||
if let Some(cache) = app_state.cache.load().as_ref() {
|
||||
let mut guard = cache.lock().await;
|
||||
if let Some(resource_path) = guard.get_content(resource_url) {
|
||||
trace_if_enabled!(
|
||||
"Responding resource from cache {}",
|
||||
sanitize_sensitive_info(resource_url)
|
||||
);
|
||||
trace_if_enabled!("Responding resource from cache {}", sanitize_sensitive_info(resource_url));
|
||||
return serve_file(&resource_path, mime::APPLICATION_OCTET_STREAM)
|
||||
.await
|
||||
.into_response();
|
||||
|
||||
@@ -3,7 +3,7 @@ use tokio_stream::wrappers::errors::BroadcastStreamRecvError;
|
||||
#[derive(Debug, Clone)]
|
||||
pub enum StreamError {
|
||||
Reqwest(String),
|
||||
// StdIo(std::io::Error),
|
||||
StdIo(String),
|
||||
// ReceiverClosed,
|
||||
ReceiverError(BroadcastStreamRecvError),
|
||||
LockError(String)
|
||||
@@ -24,7 +24,7 @@ impl std::fmt::Display for StreamError {
|
||||
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
||||
match self {
|
||||
StreamError::Reqwest(e) => write!(f, "Reqwest error: {e}"),
|
||||
// StreamError::StdIo(e) => write!(f, "IO error: {e}"),
|
||||
StreamError::StdIo(e) => write!(f, "IO error: {e}"),
|
||||
// StreamError::ReceiverClosed => write!(f, "Receiver closed"),
|
||||
StreamError::ReceiverError(e) => write!(f, "Receiver error {e}"),
|
||||
StreamError::LockError(e) => write!(f, "{e}")
|
||||
|
||||
@@ -1,144 +1,68 @@
|
||||
use std::path::Path;
|
||||
use crate::api::model::StreamError;
|
||||
use bytes::Bytes;
|
||||
use log::{debug, error, warn};
|
||||
use std::collections::VecDeque;
|
||||
use std::pin::Pin;
|
||||
use std::sync::atomic::{AtomicUsize, Ordering};
|
||||
use log::{debug, error};
|
||||
use std::sync::Arc;
|
||||
use std::task::{Context, Poll};
|
||||
use tokio::io::AsyncWrite;
|
||||
use tokio_stream::Stream;
|
||||
use tokio::io::AsyncWriteExt;
|
||||
use tokio_stream::{StreamExt};
|
||||
use tokio_stream::wrappers::ReceiverStream;
|
||||
|
||||
/// `PersistPipeStream`
|
||||
///
|
||||
/// Pipes bytes from an upstream stream to an async writer while tracking total size.
|
||||
/// Once the stream completes and the writer is flushed, the provided callback is invoked
|
||||
/// with the total number of bytes written.
|
||||
pub struct PersistPipeStream<S, W> {
|
||||
inner: S,
|
||||
completed: bool,
|
||||
writer: W,
|
||||
size: AtomicUsize,
|
||||
pub fn tee_stream<S, W>(
|
||||
mut stream: S,
|
||||
mut writer: W,
|
||||
file_path: &Path,
|
||||
callback: Arc<dyn Fn(usize) + Send + Sync>,
|
||||
pending_writes: VecDeque<Bytes>,
|
||||
current_offset: usize,
|
||||
}
|
||||
|
||||
impl<S, W> PersistPipeStream<S, W>
|
||||
where
|
||||
S: Stream + Unpin,
|
||||
W: AsyncWrite + Unpin + 'static,
|
||||
) -> ReceiverStream<Result<Bytes, StreamError>>
|
||||
where S: tokio_stream::Stream<Item = Result<Bytes, StreamError>> + Send + Unpin + 'static,
|
||||
W: tokio::io::AsyncWrite + Send + Unpin + 'static,
|
||||
{
|
||||
pub fn new(inner: S, writer: W, callback: Arc<dyn Fn(usize) + Send + Sync>) -> Self {
|
||||
Self {
|
||||
inner,
|
||||
completed: false,
|
||||
writer,
|
||||
size: AtomicUsize::new(0),
|
||||
callback,
|
||||
pending_writes: VecDeque::new(),
|
||||
current_offset: 0,
|
||||
}
|
||||
}
|
||||
let (tx, rx) = tokio::sync::mpsc::channel::<Result<Bytes, StreamError>>(32);
|
||||
let resource_path = file_path.to_owned();
|
||||
|
||||
fn enqueue_chunk(&mut self, bytes: Bytes) {
|
||||
self.pending_writes.push_back(bytes);
|
||||
}
|
||||
tokio::spawn(async move {
|
||||
let mut total_size = 0usize;
|
||||
let mut writer_active = true;
|
||||
let mut write_err: Option<StreamError> = None;
|
||||
|
||||
fn poll_pending_writes(&mut self, cx: &mut Context<'_>) -> Poll<()> {
|
||||
while let Some(chunk) = self.pending_writes.front() {
|
||||
let chunk_len = chunk.len();
|
||||
if self.current_offset >= chunk_len {
|
||||
if let Some(finished) = self.pending_writes.pop_front() {
|
||||
self.size.fetch_add(finished.len(), Ordering::AcqRel);
|
||||
}
|
||||
self.current_offset = 0;
|
||||
continue;
|
||||
}
|
||||
|
||||
let remaining = &chunk[self.current_offset..];
|
||||
match Pin::new(&mut self.writer).poll_write(cx, remaining) {
|
||||
Poll::Pending => return Poll::Pending,
|
||||
Poll::Ready(Ok(written)) => {
|
||||
if written == 0 {
|
||||
// Avoid tight loop if writer makes no progress.
|
||||
return Poll::Pending;
|
||||
while let Some(chunk) = stream.next().await {
|
||||
match chunk {
|
||||
Ok(bytes) => {
|
||||
if writer_active {
|
||||
total_size += bytes.len();
|
||||
if let Err(e) = writer.write_all(&bytes).await {
|
||||
writer_active = false;
|
||||
write_err = Some(StreamError::StdIo(e.to_string()));
|
||||
}
|
||||
}
|
||||
self.current_offset += written;
|
||||
|
||||
let _ = tx.send(Ok(bytes)).await;
|
||||
}
|
||||
Poll::Ready(Err(err)) => {
|
||||
warn!(
|
||||
"Dropping {} buffered bytes after persistence write error: {err}",
|
||||
chunk_len.saturating_sub(self.current_offset)
|
||||
);
|
||||
error!("Error writing to resource file: {err}");
|
||||
self.pending_writes.pop_front();
|
||||
self.current_offset = 0;
|
||||
Err(e) => {
|
||||
let _ = tx.send(Err(e)).await;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
self.current_offset = 0;
|
||||
Poll::Ready(())
|
||||
}
|
||||
|
||||
fn poll_flush(&mut self, cx: &mut Context<'_>) -> Poll<()> {
|
||||
match Pin::new(&mut self.writer).poll_flush(cx) {
|
||||
Poll::Pending => Poll::Pending,
|
||||
Poll::Ready(Ok(())) => Poll::Ready(()),
|
||||
Poll::Ready(Err(err)) => {
|
||||
error!("Error flushing resource file: {err}");
|
||||
Poll::Ready(())
|
||||
// final flush & shutdown
|
||||
if writer_active {
|
||||
if let Err(e) = writer.flush().await {
|
||||
writer_active = false;
|
||||
write_err = Some(StreamError::StdIo(e.to_string()));
|
||||
}
|
||||
}
|
||||
}
|
||||
let _ = writer.shutdown().await;
|
||||
|
||||
fn finalize(&mut self) {
|
||||
if !self.completed {
|
||||
self.completed = true;
|
||||
let size = self.size.load(Ordering::Acquire);
|
||||
debug!("Persisted {size} bytes to cache resource");
|
||||
(self.callback)(size);
|
||||
if writer_active {
|
||||
debug!("Persisted {total_size} bytes to cache resource");
|
||||
(callback)(total_size);
|
||||
} else {
|
||||
if let Some(err) = write_err {
|
||||
error!("Persisted stream error: {err}.");
|
||||
}
|
||||
drop(writer);
|
||||
let _ = tokio::fs::remove_file(&resource_path).await;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl<S, W> Stream for PersistPipeStream<S, W>
|
||||
where
|
||||
S: Stream<Item=Result<bytes::Bytes, StreamError>> + Unpin,
|
||||
W: AsyncWrite + Unpin + 'static,
|
||||
{
|
||||
type Item = Result<Bytes, StreamError>;
|
||||
|
||||
fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
|
||||
let this = self.get_mut();
|
||||
|
||||
if !this.pending_writes.is_empty() && this.poll_pending_writes(cx).is_pending() {
|
||||
return Poll::Pending;
|
||||
}
|
||||
|
||||
match Pin::new(&mut this.inner).poll_next(cx) {
|
||||
Poll::Pending => Poll::Pending,
|
||||
Poll::Ready(None) => {
|
||||
if this.poll_pending_writes(cx).is_pending() {
|
||||
return Poll::Pending;
|
||||
}
|
||||
if this.poll_flush(cx).is_pending() {
|
||||
return Poll::Pending;
|
||||
}
|
||||
this.finalize();
|
||||
Poll::Ready(None)
|
||||
}
|
||||
Poll::Ready(Some(item)) => {
|
||||
if let Ok(bytes) = &item {
|
||||
this.enqueue_chunk(bytes.clone());
|
||||
}
|
||||
// Try to drain pending bytes after queuing new data.
|
||||
if this.poll_pending_writes(cx).is_pending() {
|
||||
// fall through: we still return the chunk to the caller even if persistence is pending
|
||||
}
|
||||
Poll::Ready(Some(item))
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
ReceiverStream::new(rx)
|
||||
}
|
||||
|
||||
@@ -1,9 +1,10 @@
|
||||
use crate::utils::{traverse_dir};
|
||||
use shared::utils::{hash_string_as_hex, human_readable_byte_size};
|
||||
use shared::utils::{hash_string_as_hex, human_readable_byte_size, sanitize_sensitive_info};
|
||||
use log::{debug, error, info, trace};
|
||||
use std::collections::{HashMap, VecDeque};
|
||||
use std::fs;
|
||||
use std::path::PathBuf;
|
||||
use crate::utils::{trace_if_enabled};
|
||||
|
||||
/// `LRUResourceCache`
|
||||
///
|
||||
@@ -119,12 +120,14 @@ impl LRUResourceCache {
|
||||
{
|
||||
if let Some((path, size)) = self.cache.get(&key) {
|
||||
if path.exists() {
|
||||
trace_if_enabled!("Responding resource from cache with key: {key} for url: {}", sanitize_sensitive_info(url));
|
||||
// Move to the end of the queue
|
||||
self.usage_order.retain(|k| k != &key); // remove from queue
|
||||
self.usage_order.push_back(key); // add to the to end
|
||||
return Some(path.clone());
|
||||
}
|
||||
{
|
||||
trace_if_enabled!("Cache inconsistency: file missing for key: {key}, url: {}", sanitize_sensitive_info(url));
|
||||
// this should not happen, someone deleted the file manually and the cache is not in sync
|
||||
self.current_size -= size;
|
||||
self.cache.remove(&key);
|
||||
|
||||
+2
-2
@@ -1,10 +1,10 @@
|
||||
[package]
|
||||
name = "frontend"
|
||||
version = "3.2.16"
|
||||
version = "3.2.17"
|
||||
edition = "2021"
|
||||
|
||||
[dependencies]
|
||||
shared = { version = "3.2.16", path = "../shared" }
|
||||
shared = { version = "3.2.17", path = "../shared" }
|
||||
chrono = "0"
|
||||
yew = "0.21"
|
||||
yew-router = "0.18"
|
||||
|
||||
+1
-1
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "shared"
|
||||
version = "3.2.16"
|
||||
version = "3.2.17"
|
||||
edition = "2021"
|
||||
|
||||
[dependencies]
|
||||
|
||||
Reference in New Issue
Block a user