diff --git a/backend/src/api/endpoints/download_api.rs b/backend/src/api/endpoints/download_api.rs index 618d3cc20..918d30812 100644 --- a/backend/src/api/endpoints/download_api.rs +++ b/backend/src/api/endpoints/download_api.rs @@ -1,7 +1,7 @@ use crate::api::model::AppState; use crate::api::model::{DownloadQueue, FileDownload, FileDownloadRequest}; use crate::model::{AppConfig, VideoDownloadConfig}; -use crate::utils::{async_file_writer, request, WRITER_BUFFER_SIZE}; +use crate::utils::{async_file_writer, request, IO_BUFFER_SIZE}; use tokio::sync::RwLock; use futures::stream::TryStreamExt; use log::info; @@ -37,7 +37,7 @@ async fn download_file(active: Arc>>, client: &reqwe match buf_writer.write_all(&chunk).await { Ok(()) => { write_counter += chunk.len(); - if write_counter >= WRITER_BUFFER_SIZE { + if write_counter >= IO_BUFFER_SIZE { buf_writer.flush().await.map_err(|err| err.to_string())?; write_counter = 0; } diff --git a/backend/src/api/model/streams/persist_pipe_stream.rs b/backend/src/api/model/streams/persist_pipe_stream.rs index 984ec1001..d568bc036 100644 --- a/backend/src/api/model/streams/persist_pipe_stream.rs +++ b/backend/src/api/model/streams/persist_pipe_stream.rs @@ -40,6 +40,7 @@ where S: tokio_stream::Stream> + Send + Unpin if write_counter > FLUSH_INTERVAL { write_counter = 0; if let Err(err) = writer.flush().await { + writer_active = false; write_err = Some(StreamError::StdIo(format!("Failed periodic flush of tee_stream writer {err}"))); } } diff --git a/backend/src/api/scheduler.rs b/backend/src/api/scheduler.rs index 7b511d092..bdf97f9bf 100644 --- a/backend/src/api/scheduler.rs +++ b/backend/src/api/scheduler.rs @@ -55,7 +55,6 @@ async fn start_scheduler(client: reqwest::Client, expression: &str, app_state: A loop { let mut upcoming = schedule.upcoming(offset).take(1); if let Some(datetime) = upcoming.next() { - let client = client.clone(); tokio::select! { () = tokio::time::sleep_until(tokio::time::Instant::from(datetime_to_instant(datetime))) => { let app_config = Arc::clone(&app_state.app_config); diff --git a/backend/src/processing/processor/xtream_series.rs b/backend/src/processing/processor/xtream_series.rs index 79c3b1aef..ee037cb8f 100644 --- a/backend/src/processing/processor/xtream_series.rs +++ b/backend/src/processing/processor/xtream_series.rs @@ -203,7 +203,7 @@ async fn process_series_info( } write_counter +=1; // periodic flush to bound BufWriter memory - if write_counter <= FLUSH_INTERVAL { + if write_counter >= FLUSH_INTERVAL { write_counter = 0; if let Err(err) = wal_writer.flush() { errors.push(notify_err!(format!("Failed periodic flush of wal content writer {err}"))); diff --git a/backend/src/repository/indexed_document.rs b/backend/src/repository/indexed_document.rs index b0ab04eee..9323e12d7 100644 --- a/backend/src/repository/indexed_document.rs +++ b/backend/src/repository/indexed_document.rs @@ -10,7 +10,7 @@ use serde::{Deserialize, Serialize}; use tempfile::NamedTempFile; use shared::error::{str_to_io_error, to_io_error}; use crate::utils; -use crate::utils::{bincode_deserialize, bincode_serialize, WRITER_BUFFER_SIZE}; +use crate::utils::{bincode_deserialize, bincode_serialize, IO_BUFFER_SIZE}; const BLOCK_SIZE: usize = 4096; const LEN_SIZE: usize = 4; @@ -280,9 +280,6 @@ where if let Some(offset) = self.index_tree.query(doc_id) { self.main_file.seek(SeekFrom::Start(u64::from(*offset)))?; let buf_size = IndexedDocument::read_content_size(&mut self.main_file)?; - if self.buffer.capacity() < buf_size { - self.buffer.reserve(buf_size - self.buffer.capacity()); - } self.buffer.resize(buf_size, 0u8); self.main_file.read_exact(&mut self.buffer[..buf_size])?; if let Ok(item) = bincode_deserialize::(&self.buffer[..buf_size]) { @@ -514,7 +511,7 @@ where gc_writer.write_all(&size_bytes)?; gc_writer.write_all(&buffer[0..buf_size])?; write_counter += buf_size; - if write_counter >= WRITER_BUFFER_SIZE { + if write_counter >= IO_BUFFER_SIZE { write_counter = 0; gc_writer.flush()?; } diff --git a/backend/src/repository/m3u_repository.rs b/backend/src/repository/m3u_repository.rs index 97c68384a..a4c5f7a5d 100644 --- a/backend/src/repository/m3u_repository.rs +++ b/backend/src/repository/m3u_repository.rs @@ -17,7 +17,7 @@ use log::error; use tokio::fs; use tokio::io::{AsyncWriteExt}; use tokio::task; -use crate::utils::{async_file_writer, WRITER_BUFFER_SIZE}; +use crate::utils::{async_file_writer, IO_BUFFER_SIZE}; macro_rules! cant_write_result { ($path:expr, $err:expr) => { @@ -66,7 +66,7 @@ async fn persist_m3u_playlist_as_text( await_playlist_write!(writer.write_all(bytes), "Failed to write entry to {} - {}", m3u_filename.display()); await_playlist_write!(writer.write_all(b"\n"), "Failed to write newline to {} - {}", m3u_filename.display()); write_counter += bytes.len() + 1; - if write_counter >= WRITER_BUFFER_SIZE { + if write_counter >= IO_BUFFER_SIZE { await_playlist_write!(writer.flush(), "Failed to flush {} - {}", m3u_filename.display()); write_counter = 0; } diff --git a/backend/src/repository/strm_repository.rs b/backend/src/repository/strm_repository.rs index f3628b215..599ffe759 100644 --- a/backend/src/repository/strm_repository.rs +++ b/backend/src/repository/strm_repository.rs @@ -8,7 +8,7 @@ use crate::repository::storage::{ensure_target_storage_path, get_input_storage_p use crate::repository::storage_const; use crate::repository::xtream_repository::{xtream_get_record_file_path, InputVodInfoRecord}; use shared::utils::{extract_extension_from_url, hash_bytes, hash_string_as_hex, truncate_string, ExportStyleConfig, CONSTANTS}; -use crate::utils::{async_file_reader, async_file_writer, normalize_string_path, truncate_filename, FileReadGuard, WRITER_BUFFER_SIZE}; +use crate::utils::{async_file_reader, async_file_writer, normalize_string_path, truncate_filename, FileReadGuard, IO_BUFFER_SIZE}; use chrono::Datelike; use filetime::{set_file_times, FileTime}; use log::{error, trace}; @@ -947,7 +947,7 @@ async fn write_strm_index_file( .write(new_line) .await .map_err(|err| format!("Failed to write strm index entry: {err}"))?; - if write_counter >= WRITER_BUFFER_SIZE { + if write_counter >= IO_BUFFER_SIZE { write_counter = 0; writer.flush().await.map_err(|err| format!("Failed to flush: {err}"))?; } diff --git a/backend/src/repository/xtream_playlist_iterator.rs b/backend/src/repository/xtream_playlist_iterator.rs index 8b64d3714..ddf703e7e 100644 --- a/backend/src/repository/xtream_playlist_iterator.rs +++ b/backend/src/repository/xtream_playlist_iterator.rs @@ -54,6 +54,14 @@ impl XtreamPlaylistIterator { let filter_ids: Option> = filter.as_ref().map(|set| { set.iter().filter_map(|s| s.parse::().ok()).collect() }); + let filter_ids: Option> = filter.as_ref().map(|set| { + set.iter().filter_map(|s| { + s.parse::().map_err(|e| { + error!("Failed to parse bouquet filter id '{}': {}", s, e); + e + }).ok() + }).collect() + }); Ok(Self { reader, diff --git a/backend/src/utils/compression/compressed_file_reader_async.rs b/backend/src/utils/compression/compressed_file_reader_async.rs index cc6d7e30f..a62c92b45 100644 --- a/backend/src/utils/compression/compressed_file_reader_async.rs +++ b/backend/src/utils/compression/compressed_file_reader_async.rs @@ -1,16 +1,16 @@ +use crate::utils::async_file_reader; +use crate::utils::compression::compression_utils::{is_deflate, is_gzip}; +use async_compression::tokio::bufread::{GzipDecoder, ZlibDecoder}; use std::path::Path; use std::pin::Pin; use std::task::{Context, Poll}; use tokio::fs::File; use tokio::io::{ - self, AsyncRead, BufReader, AsyncSeekExt, AsyncReadExt, ReadBuf, + self, AsyncRead, AsyncReadExt, AsyncSeekExt, ReadBuf, }; -use async_compression::tokio::bufread::{GzipDecoder, ZlibDecoder}; -use crate::utils::async_file_reader; -use crate::utils::compression::compression_utils::{is_deflate, is_gzip}; pub struct CompressedFileReaderAsync { - reader: BufReader>, + reader: Box, } impl CompressedFileReaderAsync { @@ -22,17 +22,13 @@ impl CompressedFileReaderAsync { buffered_file.read_exact(&mut header).await?; buffered_file.seek(io::SeekFrom::Start(0)).await?; - let reader: Box = if is_gzip(&header) { - Box::new(GzipDecoder::new(buffered_file)) + if is_gzip(&header) { + Ok(Self { reader: Box::new(GzipDecoder::new(buffered_file)) }) } else if is_deflate(&header) { - Box::new(ZlibDecoder::new(buffered_file)) + Ok(Self { reader: Box::new(ZlibDecoder::new(buffered_file)) }) } else { - Box::new(buffered_file) - }; - - Ok(Self { - reader: async_file_reader(reader), - }) + Ok(Self { reader: Box::new(buffered_file) }) + } } } diff --git a/backend/src/utils/file/file_utils.rs b/backend/src/utils/file/file_utils.rs index 1e719ded8..698e8a41a 100644 --- a/backend/src/utils/file/file_utils.rs +++ b/backend/src/utils/file/file_utils.rs @@ -11,34 +11,34 @@ use log::{debug, error}; use path_clean::PathClean; use tokio::fs as tokio_fs; -pub const WRITER_BUFFER_SIZE: usize = 256*1024; // 256kb +pub const IO_BUFFER_SIZE: usize = 256*1024; // 256kb pub fn file_writer(w: W) -> std::io::BufWriter where W: std::io::Write, { - std::io::BufWriter::with_capacity(WRITER_BUFFER_SIZE, w) + std::io::BufWriter::with_capacity(IO_BUFFER_SIZE, w) } pub fn file_reader(r: R) -> std::io::BufReader where R: std::io::Read, { - std::io::BufReader::with_capacity(WRITER_BUFFER_SIZE, r) + std::io::BufReader::with_capacity(IO_BUFFER_SIZE, r) } pub fn async_file_writer(w: W) -> tokio::io::BufWriter where W: tokio::io::AsyncWrite, { - tokio::io::BufWriter::with_capacity(WRITER_BUFFER_SIZE, w) + tokio::io::BufWriter::with_capacity(IO_BUFFER_SIZE, w) } pub fn async_file_reader(r: R) -> tokio::io::BufReader where R: tokio::io::AsyncRead, { - tokio::io::BufReader::with_capacity(WRITER_BUFFER_SIZE, r) + tokio::io::BufReader::with_capacity(IO_BUFFER_SIZE, r) } diff --git a/backend/src/utils/network/request.rs b/backend/src/utils/network/request.rs index 3986c0eac..b184291c8 100644 --- a/backend/src/utils/network/request.rs +++ b/backend/src/utils/network/request.rs @@ -20,7 +20,7 @@ use crate::model::{ConfigInput}; use crate::repository::storage::{get_input_storage_path}; use crate::repository::storage_const; use crate::utils::compression::compression_utils::{is_deflate, is_gzip}; -use crate::utils::{async_file_reader, async_file_writer, debug_if_enabled, WRITER_BUFFER_SIZE}; +use crate::utils::{async_file_reader, async_file_writer, debug_if_enabled, IO_BUFFER_SIZE}; use shared::utils::{filter_request_header, sanitize_sensitive_info, short_hash, ENCODING_DEFLATE, ENCODING_GZIP}; use crate::utils::{get_file_path, persist_file}; @@ -257,7 +257,7 @@ async fn get_remote_content_as_file(client: &reqwest::Client, input: &ConfigInpu Ok(bytes) => { write_counter += bytes.len(); writer.write_all(&bytes).await?; - if write_counter >= WRITER_BUFFER_SIZE { + if write_counter >= IO_BUFFER_SIZE { writer.flush().await?; write_counter = 0; }