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..10f48880e 100644 --- a/backend/src/api/model/streams/persist_pipe_stream.rs +++ b/backend/src/api/model/streams/persist_pipe_stream.rs @@ -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 { - 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; } - } -} - -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)) - } - } - } + }); + + ReceiverStream::new(rx) } diff --git a/backend/src/tools/lru_cache.rs b/backend/src/tools/lru_cache.rs index 7fe913ea7..d666ab1c4 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!("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); 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]