diff --git a/backend/src/api/model/streams/persist_pipe_stream.rs b/backend/src/api/model/streams/persist_pipe_stream.rs index d568bc036..42642a913 100644 --- a/backend/src/api/model/streams/persist_pipe_stream.rs +++ b/backend/src/api/model/streams/persist_pipe_stream.rs @@ -6,8 +6,7 @@ use std::sync::Arc; use tokio::io::AsyncWriteExt; use tokio_stream::{StreamExt}; use tokio_stream::wrappers::ReceiverStream; - -const FLUSH_INTERVAL: usize = 50; +use crate::utils::IO_BUFFER_SIZE; pub fn tee_stream( mut stream: S, @@ -36,8 +35,8 @@ where S: tokio_stream::Stream> + Send + Unpin writer_active = false; write_err = Some(StreamError::StdIo(e.to_string())); } else { - write_counter += 1; - if write_counter > FLUSH_INTERVAL { + write_counter += bytes.len(); + if write_counter >= IO_BUFFER_SIZE { write_counter = 0; if let Err(err) = writer.flush().await { writer_active = false; diff --git a/backend/src/model/xmltv.rs b/backend/src/model/xmltv.rs index 5ec19fc6c..a47652c6d 100644 --- a/backend/src/model/xmltv.rs +++ b/backend/src/model/xmltv.rs @@ -126,7 +126,7 @@ impl Epg { } } write_counter += 1; - if write_counter > 50 { + if write_counter >= 50 { writer.get_mut().flush().await?; // flush underlying writer write_counter = 0; } diff --git a/backend/src/processing/processor/xtream_series.rs b/backend/src/processing/processor/xtream_series.rs index ee037cb8f..d4fdd7ec7 100644 --- a/backend/src/processing/processor/xtream_series.rs +++ b/backend/src/processing/processor/xtream_series.rs @@ -19,7 +19,7 @@ use log::{error, info, log_enabled, warn, Level}; use crate::model::{XtreamSeriesEpisode, XtreamSeriesInfoEpisode}; use crate::utils; use crate::processing::processor::xtream::normalize_json_content; -use crate::utils::bincode_serialize; +use crate::utils::{bincode_serialize, IO_BUFFER_SIZE}; create_resolve_options_function_for_xtream_target!(series); @@ -93,10 +93,10 @@ async fn playlist_resolve_series_info(cfg: &AppConfig, client: &reqwest::Client, |err| errors.push(notify_err!(format!("Failed to resolve series, could not write to wal file {err}")))); processed_info_ids.insert(provider_id, ts); content_updated = true; - write_counter += 1; + write_counter += normalized_str.len(); // periodic flush to bound BufWriter memory - if write_counter >= FLUSH_INTERVAL { + if write_counter >= IO_BUFFER_SIZE { write_counter = 0; if let Err(err) = content_writer.flush() { errors.push(notify_err!(format!("Failed periodic flush of wal content writer {err}"))); @@ -121,9 +121,9 @@ async fn playlist_resolve_series_info(cfg: &AppConfig, client: &reqwest::Client, // record_wal contains provider_id and timestamp if content_updated { handle_error!(content_writer.flush(), - |err| errors.push(notify_err!(format!("Failed to resolve vod, could not write to wal file {err}")))); + |err| errors.push(notify_err!(format!("Failed to resolve series, could not write to wal file {err}")))); handle_error!(record_writer.flush(), - |err| errors.push(notify_err!(format!("Failed to resolve vod tmdb, could not write to wal file {err}")))); + |err| errors.push(notify_err!(format!("Failed to resolve series tmdb, could not write to wal file {err}")))); handle_error!(content_writer.get_ref().sync_all(), |err| errors.push(notify_err!(format!("Failed to sync series info to wal file {err}")))); handle_error!(record_writer.get_ref().sync_all(), |err| errors.push(notify_err!(format!("Failed to sync series info record to wal file {err}")))); drop(content_writer); diff --git a/backend/src/processing/processor/xtream_vod.rs b/backend/src/processing/processor/xtream_vod.rs index a9189f72f..48500781e 100644 --- a/backend/src/processing/processor/xtream_vod.rs +++ b/backend/src/processing/processor/xtream_vod.rs @@ -15,6 +15,7 @@ use std::time::Instant; use log::{info, log_enabled, Level}; use crate::utils; use crate::processing::processor::xtream::normalize_json_content; +use crate::utils::IO_BUFFER_SIZE; create_resolve_options_function_for_xtream_target!(vod); @@ -63,8 +64,6 @@ fn should_update_vod_info(pli: &mut PlaylistItem, processed_provider_ids: &HashM should_update_info(pli, processed_provider_ids, crate::model::XC_TAG_VOD_INFO_ADDED) } -const FLUSH_INTERVAL: usize = 50; - pub async fn playlist_resolve_vod(app_config: &AppConfig, client: &reqwest::Client, target: &ConfigTarget, errors: &mut Vec, fpl: &mut FetchedPlaylist<'_>) { let (resolve_movies, resolve_delay) = get_resolve_vod_options(target, fpl); if !resolve_movies { return; } @@ -108,9 +107,9 @@ pub async fn playlist_resolve_vod(app_config: &AppConfig, client: &reqwest::Clie |err| errors.push(notify_err!(format!("Failed to resolve vod, could not write to wal file {err}")))); processed_info_ids.insert(provider_id, ts); content_updated = true; - write_counter += 1; + write_counter += normalized_str.len(); // periodic flush to bound BufWriter memory - if write_counter >= FLUSH_INTERVAL { + if write_counter >= IO_BUFFER_SIZE { write_counter = 0; if let Err(err) = content_writer.flush() { errors.push(notify_err!(format!("Failed periodic flush of wal content writer {err}"))); diff --git a/backend/src/repository/strm_repository.rs b/backend/src/repository/strm_repository.rs index 599ffe759..0cfdeb873 100644 --- a/backend/src/repository/strm_repository.rs +++ b/backend/src/repository/strm_repository.rs @@ -944,7 +944,7 @@ async fn write_strm_index_file( .await .map_err(|err| format!("Failed to write strm index entry: {err}"))?; writer - .write(new_line) + .write_all(new_line) .await .map_err(|err| format!("Failed to write strm index entry: {err}"))?; if write_counter >= IO_BUFFER_SIZE { diff --git a/backend/src/utils/file/file_utils.rs b/backend/src/utils/file/file_utils.rs index 698e8a41a..ba677119e 100644 --- a/backend/src/utils/file/file_utils.rs +++ b/backend/src/utils/file/file_utils.rs @@ -3,7 +3,7 @@ use std::collections::HashSet; use std::fs::{File, OpenOptions}; use std::path::{Path, PathBuf}; use std::{env, fs}; -use std::io::Read as IORead; +use std::io::Read; use shared::error::str_to_io_error; use crate::utils::debug_if_enabled; use shared::utils::{API_PROXY_FILE, CONFIG_FILE, CONFIG_PATH, MAPPING_FILE, SOURCE_FILE, USER_FILE}; diff --git a/backend/src/utils/network/request.rs b/backend/src/utils/network/request.rs index b184291c8..9f8baaa82 100644 --- a/backend/src/utils/network/request.rs +++ b/backend/src/utils/network/request.rs @@ -263,6 +263,8 @@ async fn get_remote_content_as_file(client: &reqwest::Client, input: &ConfigInpu } } Err(err) => { + let _ = writer.flush().await; + let _ = writer.shutdown().await; return Err(str_to_io_error(&format!("Failed to read chunk: {err}"))); } }