From 6f2be9693c0a8488ee3e82daeb04ecdced5a408b Mon Sep 17 00:00:00 2001 From: euzu Date: Fri, 14 Nov 2025 22:45:08 +0100 Subject: [PATCH] Using async file operations --- .../api/model/streams/persist_pipe_stream.rs | 34 +++++++------------ backend/src/repository/m3u_repository.rs | 6 ++-- backend/src/utils/network/request.rs | 4 +-- 3 files changed, 18 insertions(+), 26 deletions(-) diff --git a/backend/src/api/model/streams/persist_pipe_stream.rs b/backend/src/api/model/streams/persist_pipe_stream.rs index dd36d8ecf..550384350 100644 --- a/backend/src/api/model/streams/persist_pipe_stream.rs +++ b/backend/src/api/model/streams/persist_pipe_stream.rs @@ -1,13 +1,13 @@ -use std::collections::VecDeque; -use std::pin::Pin; -use std::sync::Arc; -use std::sync::atomic::{AtomicUsize, Ordering}; -use std::task::{Context, Poll}; +use crate::api::model::StreamError; use bytes::Bytes; use log::error; -use tokio::io::{AsyncWrite, AsyncWriteExt}; +use std::collections::VecDeque; +use std::pin::Pin; +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::sync::Arc; +use std::task::{Context, Poll}; +use tokio::io::AsyncWrite; use tokio_stream::Stream; -use crate::api::model::StreamError; /// `PersistPipeStream` /// @@ -99,7 +99,7 @@ where impl Stream for PersistPipeStream where - S: Stream> + Unpin, + S: Stream> + Unpin, W: AsyncWrite + Unpin + 'static, { type Item = Result; @@ -107,19 +107,17 @@ where fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { let this = self.get_mut(); - if !this.pending_writes.is_empty() { - if let Poll::Pending = this.poll_pending_writes(cx) { - return Poll::Pending; - } + 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 let Poll::Pending = this.poll_pending_writes(cx) { + if this.poll_pending_writes(cx).is_pending() { return Poll::Pending; } - if let Poll::Pending = this.poll_flush(cx) { + if this.poll_flush(cx).is_pending() { return Poll::Pending; } this.finalize(); @@ -130,7 +128,7 @@ where this.enqueue_chunk(bytes.clone()); } // Try to drain pending bytes after queuing new data. - if let Poll::Pending = this.poll_pending_writes(cx) { + 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)) @@ -138,9 +136,3 @@ where } } } - -impl Drop for PersistPipeStream { - fn drop(&mut self) { - self.finalize(); - } -} \ No newline at end of file diff --git a/backend/src/repository/m3u_repository.rs b/backend/src/repository/m3u_repository.rs index 72430e6c3..552160385 100644 --- a/backend/src/repository/m3u_repository.rs +++ b/backend/src/repository/m3u_repository.rs @@ -108,8 +108,8 @@ pub async fn m3u_write_playlist( Err(err) => Err(cant_write_result!(&m3u_path_clone, err)), } }) - .await - .map_err(|err| create_tuliprox_error!(TuliproxErrorKind::Notify, "failed to write m3u playlist: {} - {err}", m3u_path.display()))??; + .await + .map_err(|err| create_tuliprox_error!(TuliproxErrorKind::Notify, "failed to write m3u playlist: {} - {err}", m3u_path.display()))??; Ok(()) } @@ -156,4 +156,4 @@ pub async fn iter_raw_m3u_playlist(config: &AppConfig, target: &ConfigTarget) -> Ok(reader) => Some((file_lock, reader)), Err(_) => None } -} +} \ No newline at end of file diff --git a/backend/src/utils/network/request.rs b/backend/src/utils/network/request.rs index e9652ff87..46d613bf5 100644 --- a/backend/src/utils/network/request.rs +++ b/backend/src/utils/network/request.rs @@ -217,12 +217,12 @@ async fn decode_local_file_bytes(content: Vec) -> Result { } } -fn local_file_not_found(path: &PathBuf) -> Error { +fn local_file_not_found(path: &Path) -> Error { let file_str = path.to_str().unwrap_or("?"); Error::new(ErrorKind::InvalidData, format!("Cant find file {file_str}")) } -pub async fn get_local_file_content(file_path: &PathBuf) -> Result { +pub async fn get_local_file_content(file_path: &Path) -> Result { match tokio::fs::read(file_path).await { Ok(content) => decode_local_file_bytes(content).await, Err(_) => Err(local_file_not_found(file_path)),