Using async file operations

This commit is contained in:
euzu
2025-11-14 22:45:08 +01:00
parent f7b213b5af
commit 6f2be9693c
3 changed files with 18 additions and 26 deletions
@@ -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<S, W> Stream for PersistPipeStream<S, W>
where
S: Stream<Item = Result<bytes::Bytes, StreamError>> + Unpin,
S: Stream<Item=Result<bytes::Bytes, StreamError>> + Unpin,
W: AsyncWrite + Unpin + 'static,
{
type Item = Result<Bytes, StreamError>;
@@ -107,19 +107,17 @@ where
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() {
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<S, W> Drop for PersistPipeStream<S, W> {
fn drop(&mut self) {
self.finalize();
}
}
+3 -3
View File
@@ -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
}
}
}
+2 -2
View File
@@ -217,12 +217,12 @@ async fn decode_local_file_bytes(content: Vec<u8>) -> Result<String, Error> {
}
}
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<String, Error> {
pub async fn get_local_file_content(file_path: &Path) -> Result<String, Error> {
match tokio::fs::read(file_path).await {
Ok(content) => decode_local_file_bytes(content).await,
Err(_) => Err(local_file_not_found(file_path)),