diff --git a/Cargo.lock b/Cargo.lock index 2916ee3b3..cceb3baf2 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1096,7 +1096,7 @@ dependencies = [ [[package]] name = "frontend" -version = "3.2.17" +version = "3.2.18" dependencies = [ "anyhow", "base64", @@ -3793,7 +3793,7 @@ dependencies = [ [[package]] name = "shared" -version = "3.2.17" +version = "3.2.18" dependencies = [ "base64", "bitflags 2.10.0", @@ -4356,7 +4356,7 @@ checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b" [[package]] name = "tuliprox" -version = "3.2.17" +version = "3.2.18" dependencies = [ "arc-swap", "async-compression", diff --git a/backend/Cargo.toml b/backend/Cargo.toml index bc9f6cd94..cbba13ba1 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "tuliprox" -version = "3.2.17" +version = "3.2.18" 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 376e836be..88615182d 100644 --- a/backend/src/api/api_utils.rs +++ b/backend/src/api/api_utils.rs @@ -12,7 +12,7 @@ use crate::api::model::{ProviderAllocation, ProviderConfig, ProviderStreamState, use crate::model::{ConfigInput, ResourceRetryConfig}; use crate::model::{ConfigTarget, ProxyUserCredentials}; use crate::tools::lru_cache::LRUResourceCache; -use crate::utils::create_new_file_for_write; +use crate::utils::{async_file_reader, async_file_writer, create_new_file_for_write}; use crate::utils::request; use crate::utils::{debug_if_enabled, trace_if_enabled}; use crate::BUILD_TIMESTAMP; @@ -35,7 +35,6 @@ use std::collections::HashMap; use std::path::Path; use std::sync::Arc; use std::time::Duration; -use tokio::io::BufWriter; use tokio::sync::Mutex; use url::Url; @@ -155,7 +154,7 @@ pub async fn serve_file(file_path: &Path, mime_type: mime::Mime) -> impl IntoRes match tokio::fs::File::open(file_path).await { Ok(file) => { - let reader = tokio::io::BufReader::new(file); + let reader = async_file_reader(file); let stream = tokio_util::io::ReaderStream::new(reader); let body = axum::body::Body::from_stream(stream); @@ -1107,7 +1106,7 @@ async fn build_stream_response( match create_new_file_for_write(&resource_path).await { Ok(file) => { debug!("Persisting resource stream {sanitized_resource_url} to {}", resource_path.display()); - let writer = BufWriter::new(file); + let writer = async_file_writer(file); let add_cache_content = get_add_cache_content(resource_url, &app_state.cache); 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))); diff --git a/backend/src/api/endpoints/download_api.rs b/backend/src/api/endpoints/download_api.rs index ff69feec6..618d3cc20 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}; +use crate::utils::{async_file_writer, request, WRITER_BUFFER_SIZE}; use tokio::sync::RwLock; use futures::stream::TryStreamExt; use log::info; @@ -36,8 +36,8 @@ async fn download_file(active: Arc>>, client: &reqwe if let Some(chunk) = item { match buf_writer.write_all(&chunk).await { Ok(()) => { - write_counter += 1; - if write_counter > 50 { + write_counter += chunk.len(); + if write_counter >= WRITER_BUFFER_SIZE { buf_writer.flush().await.map_err(|err| err.to_string())?; write_counter = 0; } diff --git a/backend/src/api/endpoints/xmltv_api.rs b/backend/src/api/endpoints/xmltv_api.rs index 2599281a2..8369651d2 100644 --- a/backend/src/api/endpoints/xmltv_api.rs +++ b/backend/src/api/endpoints/xmltv_api.rs @@ -8,7 +8,7 @@ use crate::repository::storage::get_target_storage_path; use crate::repository::storage_const; use crate::repository::xtream_repository::{xtream_get_epg_file_path, xtream_get_storage_path}; use crate::utils; -use crate::utils::{deobscure_text, obscure_text}; +use crate::utils::{async_file_reader, deobscure_text, obscure_text}; use axum::response::IntoResponse; use chrono::{DateTime, Duration, FixedOffset, NaiveDateTime, Offset, TimeZone, Utc}; use chrono_tz::Tz; @@ -206,7 +206,7 @@ async fn serve_epg_with_rewrites( let rewrite_base_url = base_url.to_owned(); match tokio::fs::File::open(epg_path).await { Ok(file) => { - let reader = tokio::io::BufReader::new(file); + let reader = async_file_reader(file); let (tx, rx) = tokio::io::duplex(8192); tokio::spawn(async move { let mut encoder = async_compression::tokio::write::GzipEncoder::new(tx); @@ -220,7 +220,7 @@ async fn serve_epg_with_rewrites( error!("EPG: Failed to write epg doc type {err}"); } - let mut xml_reader = quick_xml::reader::Reader::from_reader(tokio::io::BufReader::new(reader)); + let mut xml_reader = quick_xml::reader::Reader::from_reader(async_file_reader(reader)); let mut xml_writer = quick_xml::writer::Writer::new(encoder); // TODO howto avoid BytesText to escape the doctype "xmltv.dtd" which is written as "e;xmltv.dtd"e; diff --git a/backend/src/model/xmltv.rs b/backend/src/model/xmltv.rs index 7896246b3..5ec19fc6c 100644 --- a/backend/src/model/xmltv.rs +++ b/backend/src/model/xmltv.rs @@ -12,6 +12,7 @@ use tokio::io::{AsyncRead, AsyncWrite, AsyncWriteExt}; use url::Url; use shared::utils::sanitize_sensitive_info; use crate::api::model::AppState; +use crate::utils::async_file_reader; use crate::utils::request::{get_remote_content_as_stream}; pub const EPG_TAG_TV: &str = "tv"; @@ -247,7 +248,7 @@ pub fn get_attr_value(attr: &quick_xml::events::attributes::Attribute) -> Option #[allow(clippy::too_many_lines)] async fn parse_xmltv_for_web_ui(reader: R) -> Result { - let mut reader = quick_xml::reader::Reader::from_reader(tokio::io::BufReader::new(reader)); + let mut reader = quick_xml::reader::Reader::from_reader(async_file_reader(reader)); let mut buf = Vec::new(); let mut channels = Vec::new(); diff --git a/backend/src/processing/parser/xmltv.rs b/backend/src/processing/parser/xmltv.rs index b94dcd537..8cf35cd5e 100644 --- a/backend/src/processing/parser/xmltv.rs +++ b/backend/src/processing/parser/xmltv.rs @@ -13,6 +13,7 @@ use std::collections::HashMap; use std::mem; use std::sync::{Mutex}; use tokio::io::AsyncRead; +use crate::utils::async_file_reader; use crate::utils::compressed_file_reader_async::CompressedFileReaderAsync; /// Splits a string at the first delimiter if the prefix matches a known country code. @@ -412,7 +413,7 @@ where F: FnMut(XmlTag), { let mut stack: Vec = vec![]; - let mut xml_reader = quick_xml::reader::Reader::from_reader(tokio::io::BufReader::new(content)); + let mut xml_reader = quick_xml::reader::Reader::from_reader(async_file_reader(content)); let mut buf = Vec::::new(); loop { match xml_reader.read_event_into_async(&mut buf).await { diff --git a/backend/src/processing/processor/xtream_series.rs b/backend/src/processing/processor/xtream_series.rs index 158388dd8..79c3b1aef 100644 --- a/backend/src/processing/processor/xtream_series.rs +++ b/backend/src/processing/processor/xtream_series.rs @@ -96,7 +96,8 @@ async fn playlist_resolve_series_info(cfg: &AppConfig, client: &reqwest::Client, write_counter += 1; // periodic flush to bound BufWriter memory - if write_counter.is_multiple_of(FLUSH_INTERVAL) { + if write_counter >= FLUSH_INTERVAL { + write_counter = 0; if let Err(err) = content_writer.flush() { errors.push(notify_err!(format!("Failed periodic flush of wal content writer {err}"))); } @@ -202,7 +203,8 @@ async fn process_series_info( } write_counter +=1; // periodic flush to bound BufWriter memory - if write_counter.is_multiple_of(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/processing/processor/xtream_vod.rs b/backend/src/processing/processor/xtream_vod.rs index 8880fb509..a9189f72f 100644 --- a/backend/src/processing/processor/xtream_vod.rs +++ b/backend/src/processing/processor/xtream_vod.rs @@ -109,9 +109,9 @@ pub async fn playlist_resolve_vod(app_config: &AppConfig, client: &reqwest::Clie processed_info_ids.insert(provider_id, ts); content_updated = true; write_counter += 1; - // periodic flush to bound BufWriter memory - if write_counter.is_multiple_of(FLUSH_INTERVAL) { + if write_counter >= FLUSH_INTERVAL { + 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/epg_repository.rs b/backend/src/repository/epg_repository.rs index baeabe6ae..3903fdfea 100644 --- a/backend/src/repository/epg_repository.rs +++ b/backend/src/repository/epg_repository.rs @@ -2,40 +2,56 @@ use crate::model::Epg; use crate::model::{Config, ConfigTarget, TargetOutput}; use crate::repository::m3u_repository::m3u_get_epg_file_path; use crate::repository::xtream_repository::{xtream_get_epg_file_path, xtream_get_storage_path}; -use crate::utils::debug_if_enabled; +use crate::utils::{async_file_writer, debug_if_enabled}; use shared::error::{notify_err, TuliproxError}; use std::path::Path; use tokio::io::AsyncWriteExt; -// Due to an error in quick_xml we cant write doc type through event. The quotes are escaped and the xml file is invalid. +const XML_PREAMBLE: &str = r#" + +"#; + +// Due to a bug in quick_xml we cannot write the DOCTYPE via event; quotes are escaped and the XML becomes invalid. +// Keep the manual header/doctype write workaround below. // -// // XML Header +// // XML Header via events (DO NOT USE, kept for documentation): // writer.write_event_async(quick_xml::events::Event::Decl(quick_xml::events::BytesDecl::new("1.0", Some("utf-8"), None))) // .await.map_err(|e| notify_err!(format!("failed to write XML header: {}", e)))?; // -// // DOCTYPE +// // DOCTYPE via events (DO NOT USE): // writer.write_event_async(quick_xml::events::Event::DocType(quick_xml::events::BytesText::new(r#"tv SYSTEM "xmltv.dtd""#))) // .await.map_err(|e| notify_err!(format!("failed to write doctype: {}", e)))?; pub async fn epg_write_file(target: &ConfigTarget, epg: &Epg, path: &Path) -> Result<(), TuliproxError> { let file = tokio::fs::File::create(path).await - .map_err(|e| notify_err!(format!("failed to create epg file: {}", e)))?; - let mut buf_writer = tokio::io::BufWriter::new(file); + .map_err(|e| notify_err!(format!("failed to create epg file {}: {}", path.display(), e)))?; - // Work-Around BytesText DocType escape, see below - buf_writer.write_all(b"\n").await - .map_err(|e| notify_err!(format!("failed to write XML header: {}", e)))?; + // Use a larger buffer for sequential writes to reduce syscalls + let mut buf_writer = async_file_writer(file); - buf_writer.write_all(b"\n").await - .map_err(|e| notify_err!(format!("failed to write doctype: {}", e)))?; + // Header/DOCTYPE workaround: write both lines in a single call for fewer syscalls + buf_writer + .write_all(XML_PREAMBLE.as_bytes()) + .await + .map_err(|e| notify_err!(format!("failed to write XML preamble {}: {}", path.display(), e)))?; + // EPG content streamed via quick_xml writer (compact output) let mut writer = quick_xml::writer::Writer::new(buf_writer); + epg + .write_to_async(&mut writer) + .await + .map_err(|e| notify_err!(format!("failed to write epg {}: {}", path.display(), e)))?; - // EPG Content - epg.write_to_async(&mut writer).await.map_err(|e| notify_err!(format!("failed to write epg: {}", e)))?; + // Ensure buffers are flushed to the OS and capture any I/O error + let mut buf_writer = writer.into_inner(); + buf_writer + .flush() + .await + .map_err(|e| notify_err!(format!("failed to flush epg {}: {}", path.display(), e)))?; + buf_writer.shutdown().await.map_err(|e| notify_err!(format!("failed to write epg {}: {}", path.display(), e)))?; - debug_if_enabled!("Epg for target {} written to {}", target.name, path.to_str().unwrap_or("?")); + debug_if_enabled!("Epg for target {} written to {}", target.name, path.display()); Ok(()) } @@ -46,7 +62,7 @@ pub async fn epg_write(cfg: &Config, target: &ConfigTarget, target_path: &Path, match xtream_get_storage_path(cfg, &target.name) { Some(path) => { let epg_path = xtream_get_epg_file_path(&path); - debug_if_enabled!("writing xtream epg to {}", epg_path.to_str().unwrap_or("?")); + debug_if_enabled!("writing xtream epg to {}", epg_path.display()); epg_write_file(target, epg_data, &epg_path).await?; } None => return Err(notify_err!(format!("failed to serialize epg for target: {}, storage path not found", target.name))), @@ -54,7 +70,7 @@ pub async fn epg_write(cfg: &Config, target: &ConfigTarget, target_path: &Path, } TargetOutput::M3u(_) => { let path = m3u_get_epg_file_path(target_path); - debug_if_enabled!("writing m3u epg to {}", path.to_str().unwrap_or("?")); + debug_if_enabled!("writing m3u epg to {}", path.display()); epg_write_file(target, epg_data, &path).await?; } TargetOutput::Strm(_) | TargetOutput::HdHomeRun(_) => {} diff --git a/backend/src/repository/indexed_document.rs b/backend/src/repository/indexed_document.rs index e03a1fa22..b0ab04eee 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}; +use crate::utils::{bincode_deserialize, bincode_serialize, WRITER_BUFFER_SIZE}; const BLOCK_SIZE: usize = 4096; const LEN_SIZE: usize = 4; @@ -163,54 +163,67 @@ where where T: ?Sized + serde::Serialize, { - let encoded_bytes = bincode_serialize(doc).map_err(|_| Error::new(ErrorKind::InvalidData, "Failed to serialize document"))?; - let mut new_record_appended = false; // do i need to change the index and set the new offset + // Determine target position (overwrite in place or append) if let Some(&offset) = self.index_tree.query(&doc_id) { + // Existing entry: check size and choose strategy self.main_file.seek(SeekFrom::Start(u64::from(offset)))?; - let size = IndexedDocument::read_content_size(&mut self.main_file)?; - if size == encoded_bytes.len() { - // check if it is equal - let mut record_buffer = vec![0; size]; - // record_buffer.resize(size, 0); + let old_size = IndexedDocument::read_content_size(&mut self.main_file)?; - self.main_file.read_exact(&mut record_buffer)?; - if record_buffer == encoded_bytes { - return Ok(()); - } - } + // Prepare new payload (bincode utils, compatible with project) + let encoded = bincode_serialize(doc) + .map_err(|e| str_to_io_error(&format!("Failed to serialize document for {}: {e}", self.main_path.display())))?; + let new_size: SizeType = SizeType::try_from(encoded.len() as u64) + .map_err(|e| str_to_io_error(&format!("Encoded document size does not fit into u32 for {}: {e}", self.main_path.display())))?; - if encoded_bytes.len() > size { - // does not fit we need to append, file is fragmented + if usize::try_from(new_size).unwrap_or(usize::MAX) > old_size { + // Does not fit -> append at end of file (mark fragmentation) if !self.fragmented { self.fragmented = true; IndexedDocument::write_fragmentation(&mut self.main_file, true)?; } self.main_file.seek(SeekFrom::End(0))?; - new_record_appended = true; - } else { - self.main_file.seek(SeekFrom::Start(u64::from(offset)))?; - } - } else { - self.main_file.seek(SeekFrom::End(0))?; - new_record_appended = true; - } + // We write: length + payload + self.dirty = true; + self.main_file.write_all(&new_size.to_le_bytes())?; + self.main_file.write_all(&encoded)?; + self.index_tree.insert(doc_id, self.main_offset); + // Increase main_offset: LEN_SIZE + new_size + let written_bytes_u64 = (LEN_SIZE as u64) + u64::from(new_size); + let written_bytes: SizeType = SizeType::try_from(written_bytes_u64) + .map_err(|e| str_to_io_error(&format!("Written byte count overflow for {}: {e}", self.main_path.display())))?; + self.main_offset = self.main_offset.checked_add(written_bytes) + .ok_or_else(|| str_to_io_error(&format!("main_offset overflow while appending to {}", self.main_path.display())))?; + return Ok(()); + } + // Fits (<=): overwrite at the same position + self.main_file.seek(SeekFrom::Start(u64::from(offset)))?; + self.dirty = true; + self.main_file.write_all(&new_size.to_le_bytes())?; + self.main_file.write_all(&encoded)?; + return Ok(()); + } + // New entry -> append + self.main_file.seek(SeekFrom::End(0))?; + // Determine size and write it, then write payload + let encoded = bincode_serialize(doc) + .map_err(|e| str_to_io_error(&format!("Failed to serialize document for {}: {e}", self.main_path.display())))?; + let new_size: SizeType = SizeType::try_from(encoded.len() as u64) + .map_err(|e| str_to_io_error(&format!("Encoded document size does not fit into u32 for {}: {e}", self.main_path.display())))?; self.dirty = true; + self.main_file + .write_all(&new_size.to_le_bytes()) + .map_err(|e| str_to_io_error(&format!("Failed to write length prefix to {}: {e}", self.main_path.display())))?; + self.main_file + .write_all(&encoded) + .map_err(|e| str_to_io_error(&format!("Failed to write document payload to {}: {e}", self.main_path.display())))?; - let encoded_bytes_len = SizeType::try_from(encoded_bytes.len()).map_err(to_io_error)?; - self.main_file.write_all(&encoded_bytes_len.to_le_bytes())?; - match utils::check_write(&self.main_file.write_all(&encoded_bytes)) { - Ok(()) => { - if new_record_appended { - self.index_tree.insert(doc_id, self.main_offset); - let written_bytes = SizeType::try_from(encoded_bytes.len() + LEN_SIZE).map_err(to_io_error)?; - self.main_offset += written_bytes; - } - } - Err(err) => { - return Err(str_to_io_error(&format!("failed to write document: {} - {}", self.main_path.display(), err))); - } - } + self.index_tree.insert(doc_id, self.main_offset); + let written_bytes_u64 = (LEN_SIZE as u64) + u64::from(new_size); + let written_bytes: SizeType = SizeType::try_from(written_bytes_u64) + .map_err(|e| str_to_io_error(&format!("Written byte count overflow for {}: {e}", self.main_path.display())))?; + self.main_offset = self.main_offset.checked_add(written_bytes) + .ok_or_else(|| str_to_io_error(&format!("main_offset overflow while appending to {}", self.main_path.display())))?; Ok(()) } } @@ -237,6 +250,7 @@ where { main_file: BufReader, index_tree: IndexedDocumentIndex, + buffer: Vec, t_type: PhantomData, } @@ -252,8 +266,10 @@ where let index_tree = IndexedDocumentIndex::::load(index_path)?; Ok(Self { - main_file: utils::file_reader(main_file), + // Larger read buffer for sequential access + main_file: BufReader::with_capacity(BLOCK_SIZE * 64, main_file), index_tree, + buffer: Vec::with_capacity(BLOCK_SIZE), t_type: PhantomData, }) } else { @@ -264,9 +280,12 @@ 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)?; - let mut buffer: Vec = vec![0; buf_size]; - self.main_file.read_exact(&mut buffer)?; - if let Ok(item) = bincode_deserialize::(&buffer) { + 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]) { return Ok(item); } } @@ -308,7 +327,8 @@ where Ok(file) => { Ok(Self { main_path: main_path.to_path_buf(), - main_file: utils::file_reader(file), + // Larger read buffer for sequential iteration + main_file: BufReader::with_capacity(BLOCK_SIZE * 128, file), offsets, index: 0, failed: false, @@ -467,6 +487,8 @@ where self.index_tree.traverse(|keys, values| { keys.iter().zip(values.iter()).for_each(|(key, &offset)| key_offset.push((key.clone(), offset))); }); + // For better disk locality: sort by offset (sequential reads) + key_offset.sort_by_key(|(_, off)| *off); let mut gc_writer = utils::file_writer(&gc_file); let fragmented_byte = 0u8.to_le_bytes(); @@ -475,6 +497,8 @@ where let mut gc_offset = 1usize; // offset is 1 because of fragment bit let mut buffer: Vec = Vec::with_capacity(BLOCK_SIZE); let mut size_bytes = [0u8; LEN_SIZE]; + let mut write_counter = 0usize; + for (key, offset) in key_offset { // read old content self.main_file.seek(SeekFrom::Start(u64::from(offset)))?; @@ -489,6 +513,11 @@ 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 { + write_counter = 0; + gc_writer.flush()?; + } let pointer = OffsetPointer::try_from(gc_offset).map_err(to_io_error)?; self.index_tree.insert(key, pointer); diff --git a/backend/src/repository/m3u_playlist_iterator.rs b/backend/src/repository/m3u_playlist_iterator.rs index b0cdd1f3e..dd24600f0 100644 --- a/backend/src/repository/m3u_playlist_iterator.rs +++ b/backend/src/repository/m3u_playlist_iterator.rs @@ -10,6 +10,7 @@ use crate::repository::storage_const; use crate::repository::user_repository::user_get_bouquet_filter; use crate::utils::FileReadGuard; use std::collections::HashSet; +// concat_string! macro from shared utils is used for efficient String building #[allow(clippy::struct_excessive_bools)] pub struct M3uPlaylistIterator { @@ -67,26 +68,39 @@ impl M3uPlaylistIterator { } fn get_rewritten_url(&self, m3u_pli: &M3uPlaylistItem, typed: bool, prefix_path: &str) -> String { - if typed { - let stream_type = match m3u_pli.item_type { + // Build URL efficiently with a single allocation using concat_string! macro + let stream_type: &str = if typed { + match m3u_pli.item_type { PlaylistItemType::Live | PlaylistItemType::Catchup | PlaylistItemType::LiveUnknown | PlaylistItemType::LiveHls | PlaylistItemType::LiveDash => "live", PlaylistItemType::Video => "movie", - PlaylistItemType::Series - | PlaylistItemType::SeriesInfo => "series", - }; - format!("{}/{prefix_path}/{stream_type}/{}/{}/{}", - &self.base_url, - &self.username, - &self.password, - m3u_pli.virtual_id + PlaylistItemType::Series | PlaylistItemType::SeriesInfo => "series", + } + } else { + "" + }; + + let mut cap = self.base_url.len() + + prefix_path.len() + + self.username.len() + + self.password.len() + + 32; // separators and id + if typed { cap += stream_type.len() + 1; } + + if typed { + shared::concat_string!( + cap = cap; + &self.base_url, "/", prefix_path, "/", stream_type, "/", + &self.username, "/", &self.password, "/", m3u_pli.virtual_id ) } else { - format!("{}/{prefix_path}/{}/{}/{}", - &self.base_url, &self.username, &self.password, m3u_pli.virtual_id + shared::concat_string!( + cap = cap; + &self.base_url, "/", prefix_path, "/", + &self.username, "/", &self.password, "/", m3u_pli.virtual_id ) } } @@ -101,14 +115,15 @@ impl M3uPlaylistIterator { fn get_next(&mut self) -> Option<(M3uPlaylistItem, bool)> { let entry = if let Some(set) = &self.filter { if let Some((current_item, _)) = self.lookup_item.take() { - let next_valid = self.reader.find(|(pli, _)| set.contains(&pli.group.clone())); + // Avoid cloning strings while filtering + let next_valid = self.reader.find(|(pli, _)| set.contains(pli.group.as_str())); self.lookup_item = next_valid; let has_next = self.lookup_item.is_some(); Some((current_item, has_next)) } else { - let current_item = self.reader.find(|(item, _)| set.contains(&item.group.clone())); + let current_item = self.reader.find(|(item, _)| set.contains(item.group.as_str())); if let Some((item, _)) = current_item { - self.lookup_item = self.reader.find(|(item, _)| set.contains(&item.group.clone())); + self.lookup_item = self.reader.find(|(item, _)| set.contains(item.group.as_str())); let has_next = self.lookup_item.is_some(); Some((item, has_next)) } else { @@ -121,19 +136,29 @@ impl M3uPlaylistIterator { // TODO hls and unknown reverse proxy entry.map(|(mut m3u_pli, has_next)| { - let is_redirect = self.proxy_type.is_redirect(m3u_pli.item_type) || self.target_options.as_ref().and_then(|o| o.force_redirect.as_ref()).is_some_and(|f| f.has_cluster(m3u_pli.item_type)); - let should_rewrite_urls = if is_redirect { self.mask_redirect_url} else { true }; - let rewrite_urls = if should_rewrite_urls { - Some((self.get_stream_url(&m3u_pli, self.include_type_in_url), if self.rewrite_resource { Some(self.get_resource_url(&m3u_pli)) } else { None })) - } else { - None - }; - let url = m3u_pli.url.clone(); - let (stream_url, resource_url) = rewrite_urls - .map_or_else(|| (url, None), |(su, ru)| (su, ru.as_ref().map(String::to_string))); + let is_redirect = self.proxy_type.is_redirect(m3u_pli.item_type) + || self + .target_options + .as_ref() + .and_then(|o| o.force_redirect.as_ref()) + .is_some_and(|f| f.has_cluster(m3u_pli.item_type)); + let should_rewrite_urls = if is_redirect { self.mask_redirect_url } else { true }; + + if should_rewrite_urls { + let stream_url = self.get_stream_url(&m3u_pli, self.include_type_in_url); + let resource_url = if self.rewrite_resource { + Some(self.get_resource_url(&m3u_pli)) + } else { + None + }; + m3u_pli.t_stream_url = stream_url; + m3u_pli.t_resource_url = resource_url; + } else { + // Keep original URL (clone required because target field is distinct) + m3u_pli.t_stream_url = m3u_pli.url.clone(); + m3u_pli.t_resource_url = None; + } - m3u_pli.t_stream_url.clone_from(&stream_url); - m3u_pli.t_resource_url.clone_from(&resource_url); (m3u_pli, has_next) }) } diff --git a/backend/src/repository/m3u_repository.rs b/backend/src/repository/m3u_repository.rs index 7d97e0f3c..97c68384a 100644 --- a/backend/src/repository/m3u_repository.rs +++ b/backend/src/repository/m3u_repository.rs @@ -15,8 +15,9 @@ use std::path::{Path, PathBuf}; use std::sync::Arc; use log::error; use tokio::fs; -use tokio::io::{AsyncWriteExt, BufWriter as AsyncBufWriter}; +use tokio::io::{AsyncWriteExt}; use tokio::task; +use crate::utils::{async_file_writer, WRITER_BUFFER_SIZE}; macro_rules! cant_write_result { ($path:expr, $err:expr) => { @@ -53,13 +54,22 @@ async fn persist_m3u_playlist_as_text( let Some(m3u_filename) = utils::get_file_path(&cfg.working_dir, Some(PathBuf::from(filename))) else { return Ok(()); }; let file = await_playlist_write!(fs::File::create(&m3u_filename), "Can't write m3u plain playlist {} - {}", m3u_filename.display()); - let mut writer = AsyncBufWriter::new(file); + // Larger buffer for sequential writes to reduce syscalls + let mut writer = async_file_writer(file); await_playlist_write!(writer.write_all(b"#EXTM3U\n"), "Failed to write header to {} - {}", m3u_filename.display()); + let mut write_counter = 0usize; + for m3u in m3u_playlist.iter() { let line = m3u.to_m3u(target.options.as_ref(), false); - await_playlist_write!(writer.write_all(line.as_bytes()), "Failed to write entry to {} - {}", m3u_filename.display()); + let bytes = line.as_bytes(); + 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 { + await_playlist_write!(writer.flush(), "Failed to flush {} - {}", m3u_filename.display()); + write_counter = 0; + } } await_playlist_write!(writer.flush(), "Failed to flush {} - {}", m3u_filename.display()); diff --git a/backend/src/repository/playlist_repository.rs b/backend/src/repository/playlist_repository.rs index 519898ccb..327f4d89a 100644 --- a/backend/src/repository/playlist_repository.rs +++ b/backend/src/repository/playlist_repository.rs @@ -193,7 +193,6 @@ async fn load_m3u_target_storage(app_config: &AppConfig, target: &ConfigTarget) } - pub async fn load_playlists_into_memory_cache(app_state: &AppState) -> Result<(), TuliproxError> { for sources in &app_state.app_config.sources.load().sources { for target in &sources.targets { diff --git a/backend/src/repository/storage.rs b/backend/src/repository/storage.rs index b964b706c..132328f88 100644 --- a/backend/src/repository/storage.rs +++ b/backend/src/repository/storage.rs @@ -6,7 +6,8 @@ use crate::repository::storage_const; use crate::utils; pub(in crate::repository) fn get_target_id_mapping_file(target_path: &Path) -> PathBuf { - target_path.join(PathBuf::from(storage_const::FILE_ID_MAPPING)) + // Join directly with &str to avoid an intermediate PathBuf allocation + target_path.join(storage_const::FILE_ID_MAPPING) } pub fn ensure_target_storage_path(cfg: &Config, target_name: &str) -> Result { diff --git a/backend/src/repository/strm_repository.rs b/backend/src/repository/strm_repository.rs index 0c57af254..f3628b215 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::{normalize_string_path, truncate_filename, FileReadGuard}; +use crate::utils::{async_file_reader, async_file_writer, normalize_string_path, truncate_filename, FileReadGuard, WRITER_BUFFER_SIZE}; use chrono::Datelike; use filetime::{set_file_times, FileTime}; use log::{error, trace}; @@ -18,7 +18,7 @@ use std::collections::{HashMap, HashSet, VecDeque}; use std::path::{Path, PathBuf}; use std::sync::Arc; use tokio::fs::{create_dir_all, remove_dir, remove_file, File}; -use tokio::io::{AsyncBufReadExt, AsyncReadExt, AsyncWriteExt, BufReader, BufWriter}; +use tokio::io::{AsyncBufReadExt, AsyncReadExt, AsyncWriteExt}; use shared::model::{ClusterFlags, FieldGetAccessor, PlaylistGroup, PlaylistItem, PlaylistItemType, StrmExportStyle, UUIDType}; use crate::utils; // Import the new MediaQuality struct @@ -932,22 +932,34 @@ async fn write_strm_index_file( let file = File::create(index_file_path) .await .map_err(|err| format!("Failed to create strm index file: {} {err}", index_file_path.display()))?; - let mut writer = BufWriter::new(file); + // Use a larger buffered writer for sequential writes to reduce syscalls + let mut writer = async_file_writer(file); + let mut write_counter = 0usize; let new_line = "\n".as_bytes(); for entry in entries { + let bytes = entry.as_bytes(); + write_counter += bytes.len() + 1; writer - .write_all(entry.as_bytes()) + .write_all(bytes) .await .map_err(|err| format!("Failed to write strm index entry: {err}"))?; writer .write(new_line) .await .map_err(|err| format!("Failed to write strm index entry: {err}"))?; + if write_counter >= WRITER_BUFFER_SIZE { + write_counter = 0; + writer.flush().await.map_err(|err| format!("Failed to flush: {err}"))?; + } } writer .flush() .await .map_err(|err| format!("failed to write strm index entry: {err}"))?; + writer + .shutdown() + .await + .map_err(|err| format!("failed to write strm index entry: {err}"))?; Ok(()) } @@ -989,7 +1001,7 @@ async fn write_strm_file( async fn has_strm_file_same_hash(file_path: &PathBuf, content_hash: UUIDType) -> bool { if let Ok(file) = File::open(&file_path).await { - let mut reader = BufReader::new(file); + let mut reader = async_file_reader(file); let mut buffer = Vec::new(); match reader.read_to_end(&mut buffer).await { Ok(_) => { @@ -1018,7 +1030,7 @@ fn get_credentials_and_server_info( async fn read_strm_file_index(strm_file_index_path: &Path) -> std::io::Result> { let file = File::open(strm_file_index_path).await?; - let reader = BufReader::new(file); + let reader = async_file_reader(file); let mut result = HashSet::new(); let mut lines = reader.lines(); while let Ok(Some(line)) = lines.next_line().await { diff --git a/backend/src/repository/xtream_playlist_iterator.rs b/backend/src/repository/xtream_playlist_iterator.rs index 122a0c795..8b64d3714 100644 --- a/backend/src/repository/xtream_playlist_iterator.rs +++ b/backend/src/repository/xtream_playlist_iterator.rs @@ -1,21 +1,22 @@ -use std::collections::HashSet; -use log::error; -use serde_json::Value; -use shared::model::{TargetType, XtreamCluster, XtreamPlaylistItem}; -use shared::error::info_err; -use shared::error::{TuliproxError}; -use crate::model::{xtream_playlistitem_to_document, AppConfig, ProxyUserCredentials}; -use crate::model::{ConfigTarget}; use crate::model::XtreamMappingOptions; -use crate::repository::indexed_document::{IndexedDocumentIterator}; +use crate::model::{xtream_playlistitem_to_document, AppConfig, ProxyUserCredentials}; +use crate::model::ConfigTarget; +use crate::repository::indexed_document::IndexedDocumentIterator; use crate::repository::user_repository::user_get_bouquet_filter; use crate::repository::xtream_repository::{xtream_get_file_paths, xtream_get_storage_path}; use crate::utils::FileReadGuard; +use log::error; +use serde_json::Value; +use shared::error::info_err; +use shared::error::TuliproxError; +use shared::model::{TargetType, XtreamCluster, XtreamPlaylistItem}; +use std::collections::HashSet; pub struct XtreamPlaylistIterator { reader: IndexedDocumentIterator, options: XtreamMappingOptions, - filter: Option>, + // Use parsed numeric filter to avoid per-item String allocations (no to_string per check) + filter_ids: Option>, base_url: String, user: ProxyUserCredentials, lookup_item: Option<(XtreamPlaylistItem, bool)>, // this is for filtered iteration @@ -49,11 +50,15 @@ impl XtreamPlaylistIterator { let server_info = app_config.get_user_server_info(user); let filter = user_get_bouquet_filter(&config, &user.username, category_id, TargetType::Xtream, cluster).await; + // Parse bouquet filter (strings) once into u32 set to minimize per-item allocations + let filter_ids: Option> = filter.as_ref().map(|set| { + set.iter().filter_map(|s| s.parse::().ok()).collect() + }); Ok(Self { reader, options, - filter, + filter_ids, _file_lock: file_lock, base_url: server_info.get_base_url(), user: user.clone(), @@ -69,16 +74,16 @@ impl XtreamPlaylistIterator { error!("Could not deserialize xtream item: {}", self.reader.get_path().display()); return None; } - if let Some(set) = &self.filter { + if let Some(set) = &self.filter_ids { if let Some((current_item, _)) = self.lookup_item.take() { - let next_valid = self.reader.find(|(pli, _)| set.contains(&pli.category_id.to_string())); + let next_valid = self.reader.find(|(pli, _)| set.contains(&pli.category_id)); self.lookup_item = next_valid; let has_next = self.lookup_item.is_some(); Some((current_item, has_next)) } else { - let current_item = self.reader.find(|(item, _)| set.contains(&item.category_id.to_string())); + let current_item = self.reader.find(|(item, _)| set.contains(&item.category_id)); if let Some((item, _)) = current_item { - self.lookup_item = self.reader.find(|(item, _)| set.contains(&item.category_id.to_string())); + self.lookup_item = self.reader.find(|(item, _)| set.contains(&item.category_id)); let has_next = self.lookup_item.is_some(); Some((item, has_next)) } else { @@ -89,7 +94,6 @@ impl XtreamPlaylistIterator { self.reader.next() } } - } impl Iterator for XtreamPlaylistIterator { @@ -105,12 +109,12 @@ pub struct XtreamPlaylistJsonIterator { } impl XtreamPlaylistJsonIterator { -pub async fn new( - cluster: XtreamCluster, - config: &AppConfig, - target: &ConfigTarget, - category_id: Option, - user: &ProxyUserCredentials, + pub async fn new( + cluster: XtreamCluster, + config: &AppConfig, + target: &ConfigTarget, + category_id: Option, + user: &ProxyUserCredentials, ) -> Result { Ok(Self { inner: XtreamPlaylistIterator::new(cluster, config, target, category_id, user).await? @@ -122,7 +126,6 @@ pub fn to_doc(pli: &XtreamPlaylistItem, url: &str, options: &XtreamMappingOption xtream_playlistitem_to_document(pli, url, options, user) } - impl Iterator for XtreamPlaylistJsonIterator { type Item = (String, bool); fn next(&mut self) -> Option { diff --git a/backend/src/repository/xtream_repository.rs b/backend/src/repository/xtream_repository.rs index 6f90ed537..e3952ff8a 100644 --- a/backend/src/repository/xtream_repository.rs +++ b/backend/src/repository/xtream_repository.rs @@ -1076,15 +1076,12 @@ where { match channels { Some((_, chans)) => { - // Convert iterator items to Result + // Convert iterator items to Result with minimal allocations let mapped = chans.map(move |(item, has_next)| { match serde_json::to_string(&item) { - Ok(content) => { - Ok(Bytes::from(if has_next { - format!("{content},") - } else { - content - })) + Ok(mut content) => { + if has_next { content.push(','); } + Ok(Bytes::from(content)) } Err(_) => Ok(Bytes::from("")), } diff --git a/backend/src/utils/compression/compressed_file_reader.rs b/backend/src/utils/compression/compressed_file_reader.rs index 5326c7613..33bc34e22 100644 --- a/backend/src/utils/compression/compressed_file_reader.rs +++ b/backend/src/utils/compression/compressed_file_reader.rs @@ -26,7 +26,7 @@ impl CompressedFileReader { }; Ok(Self { - reader: BufReader::new(reader), + reader: file_reader(reader), }) } } diff --git a/backend/src/utils/compression/compressed_file_reader_async.rs b/backend/src/utils/compression/compressed_file_reader_async.rs index 9ea2b60bd..cc6d7e30f 100644 --- a/backend/src/utils/compression/compressed_file_reader_async.rs +++ b/backend/src/utils/compression/compressed_file_reader_async.rs @@ -6,7 +6,7 @@ use tokio::io::{ self, AsyncRead, BufReader, AsyncSeekExt, AsyncReadExt, 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 { @@ -17,7 +17,7 @@ impl CompressedFileReaderAsync { pub async fn new(path: &Path) -> std::io::Result { let file: File = tokio::fs::File::open(path).await?; - let mut buffered_file = BufReader::new(file); + let mut buffered_file = async_file_reader(file); let mut header = [0u8; 2]; buffered_file.read_exact(&mut header).await?; buffered_file.seek(io::SeekFrom::Start(0)).await?; @@ -31,7 +31,7 @@ impl CompressedFileReaderAsync { }; Ok(Self { - reader: BufReader::new(reader), + reader: async_file_reader(reader), }) } } @@ -46,19 +46,3 @@ impl AsyncRead for CompressedFileReaderAsync { } } -// -// impl AsyncBufRead for CompressedFileReaderAsync { -// fn poll_fill_buf( -// self: Pin<&mut Self>, -// cx: &mut Context<'_>, -// ) -> Poll> { -// unsafe { -// let this = self.get_unchecked_mut(); -// Pin::new_unchecked(&mut this.reader).poll_fill_buf(cx) -// } -// } -// -// fn consume(mut self: Pin<&mut Self>, amt: usize) { -// Pin::new(&mut self.reader).consume(amt) -// } -// } diff --git a/backend/src/utils/file/config_reader.rs b/backend/src/utils/file/config_reader.rs index c42d2429b..4b6e37feb 100644 --- a/backend/src/utils/file/config_reader.rs +++ b/backend/src/utils/file/config_reader.rs @@ -13,7 +13,7 @@ use shared::model::{ApiProxyConfigDto, AppConfigDto, ConfigDto, ConfigInputAlias use shared::utils::{CONSTANTS}; use std::env; use std::fs::File; -use std::io::{self, BufReader, Read}; +use std::io::{self, Read}; use std::path::PathBuf; use std::sync::Arc; use arc_swap::access::{Access}; @@ -40,7 +40,7 @@ pub fn config_file_reader(file: File, resolve_env: bool) -> impl Read if resolve_env { EitherReader::Left(EnvResolvingReader::new(file_reader(file))) } else { - EitherReader::Right(BufReader::new(file)) + EitherReader::Right(file_reader(file)) } } diff --git a/backend/src/utils/file/csv_input_reader.rs b/backend/src/utils/file/csv_input_reader.rs index b64e52f12..17fd6faf6 100644 --- a/backend/src/utils/file/csv_input_reader.rs +++ b/backend/src/utils/file/csv_input_reader.rs @@ -189,7 +189,7 @@ pub fn get_csv_file_path(file_uri: &str) -> Result { #[cfg(test)] mod tests { use crate::utils::file::csv_input_reader::csv_read_inputs_from_reader; - use crate::utils::resolve_env_var; + use crate::utils::{file_reader, resolve_env_var}; use std::io::{BufReader, Cursor}; use shared::model::InputType; @@ -209,7 +209,7 @@ input_2;de566567;de2345f43g5;http://provider_2.tv:8080;1;2028-12-23 13:12:34 #[test] fn test_read_inputs_xtream_as_m3u() { - let reader = BufReader::new(Cursor::new(XTREAM_BATCH)); + let reader = file_reader(Cursor::new(XTREAM_BATCH)); let result = csv_read_inputs_from_reader(InputType::M3uBatch, reader); assert!(result.is_ok()); let aliases = result.unwrap(); @@ -221,7 +221,7 @@ input_2;de566567;de2345f43g5;http://provider_2.tv:8080;1;2028-12-23 13:12:34 #[test] fn test_read_inputs_m3u_as_m3u() { - let reader = BufReader::new(Cursor::new(M3U_BATCH)); + let reader = file_reader(Cursor::new(M3U_BATCH)); let result = csv_read_inputs_from_reader(InputType::M3uBatch, reader); assert!(result.is_ok()); let aliases = result.unwrap(); @@ -233,7 +233,7 @@ input_2;de566567;de2345f43g5;http://provider_2.tv:8080;1;2028-12-23 13:12:34 #[test] fn test_read_inputs_xtream_as_xtream() { - let reader = BufReader::new(Cursor::new(XTREAM_BATCH)); + let reader = file_reader(Cursor::new(XTREAM_BATCH)); let result = csv_read_inputs_from_reader(InputType::XtreamBatch, reader); assert!(result.is_ok()); let aliases = result.unwrap(); @@ -245,7 +245,7 @@ input_2;de566567;de2345f43g5;http://provider_2.tv:8080;1;2028-12-23 13:12:34 #[test] fn test_read_inputs_m3u_as_xtream() { - let reader = BufReader::new(Cursor::new(M3U_BATCH)); + let reader = file_reader(Cursor::new(M3U_BATCH)); let result = csv_read_inputs_from_reader(InputType::XtreamBatch, reader); assert!(result.is_ok()); let aliases = result.unwrap(); diff --git a/backend/src/utils/file/file_utils.rs b/backend/src/utils/file/file_utils.rs index f53477cd4..1e719ded8 100644 --- a/backend/src/utils/file/file_utils.rs +++ b/backend/src/utils/file/file_utils.rs @@ -11,7 +11,7 @@ use log::{debug, error}; use path_clean::PathClean; use tokio::fs as tokio_fs; -const WRITER_BUFFER_SIZE: usize = 131_072; // 128kb +pub const WRITER_BUFFER_SIZE: usize = 256*1024; // 256kb pub fn file_writer(w: W) -> std::io::BufWriter where @@ -34,6 +34,13 @@ where tokio::io::BufWriter::with_capacity(WRITER_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) +} + pub fn get_exe_path() -> PathBuf { let default_path = std::path::PathBuf::from("./"); diff --git a/backend/src/utils/geoip.rs b/backend/src/utils/geoip.rs index ff3430cda..e5cadd7d1 100644 --- a/backend/src/utils/geoip.rs +++ b/backend/src/utils/geoip.rs @@ -100,13 +100,14 @@ mod test { use std::fs::File; use std::io::BufReader; use std::path::PathBuf; + use crate::utils::file_reader; #[test] pub fn test_csv() { let db_file = PathBuf::from("/projects/m3u-test/asn-country-ipv4.db"); let source = PathBuf::from("/projects/m3u-test/asn-country-ipv4.csv"); let file = File::open(source).expect("Could not open csv file"); - let reader = BufReader::new(file); + let reader = file_reader(file); let mut geo_ip = GeoIp::new(); let _ = geo_ip.import_ipv4_from_csv(reader, &db_file).expect("Could not import csv"); diff --git a/backend/src/utils/json_utils.rs b/backend/src/utils/json_utils.rs index 7ef9e90cc..43888607a 100644 --- a/backend/src/utils/json_utils.rs +++ b/backend/src/utils/json_utils.rs @@ -1,4 +1,4 @@ -use crate::utils::file_reader; +use crate::utils::{file_reader, file_writer}; use serde::Serialize; use serde_json::Value; use std::collections::{HashMap, HashSet}; @@ -49,7 +49,7 @@ where T: Serialize, { let file = std::fs::File::create(path)?; - let mut buf_writer = std::io::BufWriter::new(file); + let mut buf_writer = file_writer(file); serde_json::to_writer(&mut buf_writer, value)?; buf_writer.flush()?; buf_writer.into_inner()?.sync_all() diff --git a/backend/src/utils/network/request.rs b/backend/src/utils/network/request.rs index 4b15cde7b..3986c0eac 100644 --- a/backend/src/utils/network/request.rs +++ b/backend/src/utils/network/request.rs @@ -8,7 +8,7 @@ use log::{debug, error, log_enabled, trace, Level}; use reqwest::header::CONTENT_ENCODING; use reqwest::header::{HeaderMap, HeaderName, HeaderValue}; use tokio::fs::File; -use tokio::io::{AsyncBufReadExt, AsyncRead, AsyncReadExt, AsyncWriteExt, BufReader, BufWriter}; +use tokio::io::{AsyncBufReadExt, AsyncRead, AsyncReadExt, AsyncWriteExt}; use tokio_util::io::StreamReader; use url::Url; @@ -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::{debug_if_enabled}; +use crate::utils::{async_file_reader, async_file_writer, debug_if_enabled, WRITER_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}; @@ -209,7 +209,7 @@ pub async fn get_local_file_content(file_path: &Path) -> Result { if response.status().is_success() { // Open a file in write mode - let mut file = BufWriter::with_capacity(8192, File::create(file_path).await?); + let mut writer = async_file_writer(File::create(file_path).await?); + let mut write_counter = 0; // Stream the response body in chunks let mut stream = response.bytes_stream(); while let Some(chunk) = stream.next().await { match chunk { Ok(bytes) => { - file.write_all(&bytes).await?; + write_counter += bytes.len(); + writer.write_all(&bytes).await?; + if write_counter >= WRITER_BUFFER_SIZE { + writer.flush().await?; + write_counter = 0; + } } Err(err) => { return Err(str_to_io_error(&format!("Failed to read chunk: {err}"))); @@ -262,7 +268,8 @@ async fn get_remote_content_as_file(client: &reqwest::Client, input: &ConfigInpu } } - file.flush().await?; + writer.flush().await?; + writer.shutdown().await?; let elapsed = start_time.elapsed().as_secs(); debug!("File downloaded successfully to {}, took:{}", file_path.display(), format_elapsed_time(elapsed)); Ok(file_path.to_path_buf()) @@ -297,7 +304,7 @@ pub async fn get_remote_content_as_stream( let stream_reader = StreamReader::new( response.bytes_stream().map_err(std::io::Error::other), ); - let mut buf_reader = BufReader::new(stream_reader); + let mut buf_reader = async_file_reader(stream_reader); let peek = buf_reader.fill_buf().await?; if peek.len() >= 2 { diff --git a/frontend/Cargo.toml b/frontend/Cargo.toml index 5f7b88406..4c1a20cd0 100644 --- a/frontend/Cargo.toml +++ b/frontend/Cargo.toml @@ -1,10 +1,10 @@ [package] name = "frontend" -version = "3.2.17" +version = "3.2.18" edition = "2021" [dependencies] -shared = { version = "3.2.17", path = "../shared" } +shared = { version = "3.2.18", path = "../shared" } chrono = "0" yew = "0.21" yew-router = "0.18" diff --git a/shared/Cargo.toml b/shared/Cargo.toml index 85658c486..5dacf7d4a 100644 --- a/shared/Cargo.toml +++ b/shared/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "shared" -version = "3.2.17" +version = "3.2.18" edition = "2021" [dependencies] diff --git a/shared/src/utils/string_utils.rs b/shared/src/utils/string_utils.rs index 08d18ffa4..f71d15e6c 100644 --- a/shared/src/utils/string_utils.rs +++ b/shared/src/utils/string_utils.rs @@ -131,11 +131,34 @@ pub fn longest<'a>(a: &'a str, b: &'a str) -> &'a str { if a.len() >= b.len() { a } else { b } } +// ------------------------------------------------------------ +// Generic string concatenation macro with optional capacity hint +// Usage: +// let s = concat_string!("/", user, "/", pass, "/", id); +// let s = concat_string!(cap = 128; prefix, "/", value); +// The macro writes all arguments using Display into a preallocated String +// to minimize temporary allocations and copies. +// ------------------------------------------------------------ +#[macro_export] +macro_rules! concat_string { + (cap = $cap:expr; $($arg:expr),* $(,)?) => {{ + let mut __s = ::std::string::String::with_capacity($cap); + $( let _ = ::std::fmt::Write::write_fmt(&mut __s, format_args!("{}", $arg)); )* + __s + }}; + ( $($arg:expr),* $(,)?) => {{ + let mut __s = ::std::string::String::new(); + $( let _ = ::std::fmt::Write::write_fmt(&mut __s, format_args!("{}", $arg)); )* + __s + }}; +} + #[cfg(test)] mod test { use std::collections::HashSet; use crate::utils::Capitalize; use super::generate_random_string; + use crate as shared; // allow path-based macro call in tests #[test] fn test_generate_random_string() { @@ -151,4 +174,20 @@ mod test { assert_eq!("hELLO".capitalize(), "Hello"); } + #[test] + fn test_concat_string_basic() { + let a = "hello"; + let b = String::from("world"); + let n = 42; + let s = shared::concat_string!(a, " ", b, " ", n); + assert_eq!(s, "hello world 42"); + } + + #[test] + fn test_concat_string_with_cap() { + let part = "abc"; + let s = shared::concat_string!(cap = 16; part, "/", 123); + assert_eq!(s, "abc/123"); + } + }