From c2652a7ecebfcf9e8e49c0156f483887cd726cc0 Mon Sep 17 00:00:00 2001 From: euzu Date: Wed, 10 Dec 2025 17:19:57 +0100 Subject: [PATCH 1/2] Made cache storage more robust --- CHANGELOG.md | 1 + Cargo.lock | 6 +- backend/Cargo.toml | 2 +- backend/src/api/api_utils.rs | 13 +- backend/src/api/model/stream_error.rs | 4 +- .../api/model/streams/persist_pipe_stream.rs | 317 +++++++++++------- backend/src/tools/lru_cache.rs | 5 +- frontend/Cargo.toml | 4 +- shared/Cargo.toml | 2 +- 9 files changed, 212 insertions(+), 142 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 2a6ac5df6..c0c9d804a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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. diff --git a/Cargo.lock b/Cargo.lock index 5318875b2..2916ee3b3 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -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", diff --git a/backend/Cargo.toml b/backend/Cargo.toml index ffae5284f..bc9f6cd94 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -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 diff --git a/backend/src/api/api_utils.rs b/backend/src/api/api_utils.rs index 34a4c4a8e..376e836be 100644 --- a/backend/src/api/api_utils.rs +++ b/backend/src/api/api_utils.rs @@ -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(); diff --git a/backend/src/api/model/stream_error.rs b/backend/src/api/model/stream_error.rs index 498ccf5e5..a9a678012 100644 --- a/backend/src/api/model/stream_error.rs +++ b/backend/src/api/model/stream_error.rs @@ -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}") diff --git a/backend/src/api/model/streams/persist_pipe_stream.rs b/backend/src/api/model/streams/persist_pipe_stream.rs index fd0dfb4c7..8baac31fe 100644 --- a/backend/src/api/model/streams/persist_pipe_stream.rs +++ b/backend/src/api/model/streams/persist_pipe_stream.rs @@ -1,144 +1,213 @@ +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 { - inner: S, - completed: bool, - writer: W, - size: AtomicUsize, +pub fn tee_stream( + mut stream: S, + mut writer: W, + file_path: &Path, callback: Arc, - pending_writes: VecDeque, - current_offset: usize, -} - -impl PersistPipeStream -where - S: Stream + Unpin, - W: AsyncWrite + Unpin + 'static, +) -> ReceiverStream> +where S: tokio_stream::Stream> + Send + Unpin + 'static, + W: tokio::io::AsyncWrite + Send + Unpin + 'static, { - pub fn new(inner: S, writer: W, callback: Arc) -> 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::>(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 = 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; } - } + }); + + ReceiverStream::new(rx) } -impl Stream for PersistPipeStream -where - S: Stream> + Unpin, - W: AsyncWrite + Unpin + 'static, -{ - type Item = Result; - - fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { - 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)) - } - } - } -} +// +// +// /// `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 { +// inner: S, +// completed: bool, +// writer: W, +// size: AtomicUsize, +// callback: Arc, +// pending_writes: VecDeque, +// current_offset: usize, +// } +// +// impl PersistPipeStream +// where +// S: Stream + Unpin, +// W: AsyncWrite + Unpin + 'static, +// { +// pub fn new(inner: S, writer: W, callback: Arc) -> Self { +// Self { +// inner, +// completed: false, +// writer, +// size: AtomicUsize::new(0), +// callback, +// pending_writes: VecDeque::new(), +// current_offset: 0, +// } +// } +// +// fn enqueue_chunk(&mut self, bytes: Bytes) { +// self.pending_writes.push_back(bytes); +// } +// +// 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 { +// cx.waker().wake_by_ref(); +// // Avoid tight loop if writer makes no progress. +// return Poll::Pending; +// } +// self.current_offset += written; +// } +// 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; +// } +// } +// } +// +// 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(()) +// } +// } +// } +// +// 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); +// } +// } +// } +// +// impl Stream for PersistPipeStream +// where +// S: Stream> + Unpin, +// W: AsyncWrite + Unpin + 'static, +// { +// type Item = Result; +// +// fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { +// 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.pending_writes.is_empty() { +// if this.poll_pending_writes(cx).is_pending() { +// return Poll::Pending; +// } +// } +// +// if !this.pending_writes.is_empty() { +// 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)) +// } +// } +// } +// } diff --git a/backend/src/tools/lru_cache.rs b/backend/src/tools/lru_cache.rs index 7fe913ea7..43c43bd84 100644 --- a/backend/src/tools/lru_cache.rs +++ b/backend/src/tools/lru_cache.rs @@ -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!("Responding resource from cache with key: {key} for 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); diff --git a/frontend/Cargo.toml b/frontend/Cargo.toml index 456024b89..5f7b88406 100644 --- a/frontend/Cargo.toml +++ b/frontend/Cargo.toml @@ -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" diff --git a/shared/Cargo.toml b/shared/Cargo.toml index 3295729e8..85658c486 100644 --- a/shared/Cargo.toml +++ b/shared/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "shared" -version = "3.2.16" +version = "3.2.17" edition = "2021" [dependencies] From 07991f172ea4d0922c73c1cd2e2cf9c22bc34077 Mon Sep 17 00:00:00 2001 From: euzu Date: Wed, 10 Dec 2025 17:30:34 +0100 Subject: [PATCH 2/2] Made cache storage more robust --- .../api/model/streams/persist_pipe_stream.rs | 145 ------------------ backend/src/tools/lru_cache.rs | 2 +- 2 files changed, 1 insertion(+), 146 deletions(-) diff --git a/backend/src/api/model/streams/persist_pipe_stream.rs b/backend/src/api/model/streams/persist_pipe_stream.rs index 8baac31fe..10f48880e 100644 --- a/backend/src/api/model/streams/persist_pipe_stream.rs +++ b/backend/src/api/model/streams/persist_pipe_stream.rs @@ -66,148 +66,3 @@ where S: tokio_stream::Stream> + Send + Unpin ReceiverStream::new(rx) } - -// -// -// /// `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 { -// inner: S, -// completed: bool, -// writer: W, -// size: AtomicUsize, -// callback: Arc, -// pending_writes: VecDeque, -// current_offset: usize, -// } -// -// impl PersistPipeStream -// where -// S: Stream + Unpin, -// W: AsyncWrite + Unpin + 'static, -// { -// pub fn new(inner: S, writer: W, callback: Arc) -> Self { -// Self { -// inner, -// completed: false, -// writer, -// size: AtomicUsize::new(0), -// callback, -// pending_writes: VecDeque::new(), -// current_offset: 0, -// } -// } -// -// fn enqueue_chunk(&mut self, bytes: Bytes) { -// self.pending_writes.push_back(bytes); -// } -// -// 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 { -// cx.waker().wake_by_ref(); -// // Avoid tight loop if writer makes no progress. -// return Poll::Pending; -// } -// self.current_offset += written; -// } -// 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; -// } -// } -// } -// -// 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(()) -// } -// } -// } -// -// 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); -// } -// } -// } -// -// impl Stream for PersistPipeStream -// where -// S: Stream> + Unpin, -// W: AsyncWrite + Unpin + 'static, -// { -// type Item = Result; -// -// fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { -// 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.pending_writes.is_empty() { -// if this.poll_pending_writes(cx).is_pending() { -// return Poll::Pending; -// } -// } -// -// if !this.pending_writes.is_empty() { -// 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)) -// } -// } -// } -// } diff --git a/backend/src/tools/lru_cache.rs b/backend/src/tools/lru_cache.rs index 43c43bd84..d666ab1c4 100644 --- a/backend/src/tools/lru_cache.rs +++ b/backend/src/tools/lru_cache.rs @@ -127,7 +127,7 @@ impl LRUResourceCache { return Some(path.clone()); } { - trace_if_enabled!("Responding resource from cache with key: {key} for url: {}", sanitize_sensitive_info(url)); + 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);