From ab8e6c0d53db086a362a4005af540bc43eaaa322 Mon Sep 17 00:00:00 2001 From: euzu Date: Thu, 11 Dec 2025 14:32:04 +0100 Subject: [PATCH] write operations optimized for less memory usage --- backend/src/api/endpoints/download_api.rs | 16 +- .../api/model/streams/persist_pipe_stream.rs | 11 + backend/src/model/xmltv.rs | 14 +- .../src/processing/processor/xtream_series.rs | 86 ++++--- backend/src/repository/bplustree.rs | 241 ++++++++++-------- backend/src/repository/epg_repository.rs | 2 - backend/src/utils/file/file_utils.rs | 25 +- 7 files changed, 241 insertions(+), 154 deletions(-) diff --git a/backend/src/api/endpoints/download_api.rs b/backend/src/api/endpoints/download_api.rs index 3ec401d1f..ff69feec6 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::request; +use crate::utils::{async_file_writer, request}; use tokio::sync::RwLock; use futures::stream::TryStreamExt; use log::info; @@ -25,15 +25,23 @@ async fn download_file(active: Arc>>, client: &reqwe if let Some(file_path_str) = file_download.file_path.to_str() { info!("Downloading {file_path_str}"); match File::create(&file_download.file_path).await { - Ok(mut file) => { + Ok(file) => { + let mut buf_writer = async_file_writer(file); let mut downloaded: u64 = 0; let mut stream = response.bytes_stream().map_err(to_io_error); + let mut write_counter = 0; loop { match stream.try_next().await { Ok(item) => { if let Some(chunk) = item { - match file.write_all(&chunk).await { + match buf_writer.write_all(&chunk).await { Ok(()) => { + write_counter += 1; + if write_counter > 50 { + buf_writer.flush().await.map_err(|err| err.to_string())?; + write_counter = 0; + } + downloaded += chunk.len() as u64; if let Some(lock) = active.write().await.as_mut() { lock.size = downloaded; @@ -47,6 +55,8 @@ async fn download_file(active: Arc>>, client: &reqwe if let Some(lock) = active.write().await.as_mut() { lock.size = downloaded; } + buf_writer.flush().await.map_err(|err| err.to_string())?; + buf_writer.shutdown().await.map_err(|err| err.to_string())?; return Ok(()); } } diff --git a/backend/src/api/model/streams/persist_pipe_stream.rs b/backend/src/api/model/streams/persist_pipe_stream.rs index 10f48880e..984ec1001 100644 --- a/backend/src/api/model/streams/persist_pipe_stream.rs +++ b/backend/src/api/model/streams/persist_pipe_stream.rs @@ -7,6 +7,8 @@ use tokio::io::AsyncWriteExt; use tokio_stream::{StreamExt}; use tokio_stream::wrappers::ReceiverStream; +const FLUSH_INTERVAL: usize = 50; + pub fn tee_stream( mut stream: S, mut writer: W, @@ -23,6 +25,7 @@ where S: tokio_stream::Stream> + Send + Unpin let mut total_size = 0usize; let mut writer_active = true; let mut write_err: Option = None; + let mut write_counter = 0usize; while let Some(chunk) = stream.next().await { match chunk { @@ -32,6 +35,14 @@ where S: tokio_stream::Stream> + Send + Unpin if let Err(e) = writer.write_all(&bytes).await { writer_active = false; write_err = Some(StreamError::StdIo(e.to_string())); + } else { + write_counter += 1; + if write_counter > FLUSH_INTERVAL { + write_counter = 0; + if let Err(err) = writer.flush().await { + write_err = Some(StreamError::StdIo(format!("Failed periodic flush of tee_stream writer {err}"))); + } + } } } diff --git a/backend/src/model/xmltv.rs b/backend/src/model/xmltv.rs index 8954f41f2..7896246b3 100644 --- a/backend/src/model/xmltv.rs +++ b/backend/src/model/xmltv.rs @@ -8,7 +8,7 @@ use std::collections::HashMap; use std::path::{Path, PathBuf}; use std::sync::Arc; use futures::TryFutureExt; -use tokio::io::{AsyncRead, AsyncWrite}; +use tokio::io::{AsyncRead, AsyncWrite, AsyncWriteExt}; use url::Url; use shared::utils::sanitize_sensitive_info; use crate::api::model::AppState; @@ -93,6 +93,8 @@ impl Epg { .map(|c| (c, false)) .collect(); + let mut write_counter = 0usize; + while let Some((tag, ended)) = stack.pop() { if ended { // End-Event @@ -122,10 +124,20 @@ impl Epg { } } } + write_counter += 1; + if write_counter > 50 { + writer.get_mut().flush().await?; // flush underlying writer + write_counter = 0; + } } // write tv-end writer.write_event_async(Event::End(BytesEnd::new("tv"))).await?; + + let inner = writer.get_mut(); + inner.flush().await?; + inner.shutdown().await?; + Ok(()) } } diff --git a/backend/src/processing/processor/xtream_series.rs b/backend/src/processing/processor/xtream_series.rs index 6a788132b..158388dd8 100644 --- a/backend/src/processing/processor/xtream_series.rs +++ b/backend/src/processing/processor/xtream_series.rs @@ -49,9 +49,11 @@ fn should_update_series_info(pli: &mut PlaylistItem, processed_provider_ids: &Ha should_update_info(pli, processed_provider_ids, crate::model::XC_TAG_SERIES_INFO_LAST_MODIFIED) } +const FLUSH_INTERVAL: usize = 50; + async fn playlist_resolve_series_info(cfg: &AppConfig, client: &reqwest::Client, errors: &mut Vec, fpl: &mut FetchedPlaylist<'_>, resolve_delay: u16) -> bool { - let mut processed_info_ids = read_processed_series_info_ids(cfg, errors, fpl).await; + let mut processed_info_ids: HashMap = read_processed_series_info_ids(cfg, errors, fpl).await; let mut fetched_in_run: HashSet = HashSet::new(); // we cant write to the indexed-document directly because of the write lock and time-consuming operation. // All readers would be waiting for the lock and the app would be unresponsive. @@ -63,45 +65,57 @@ async fn playlist_resolve_series_info(cfg: &AppConfig, client: &reqwest::Client, let mut record_writer = utils::file_writer(&wal_record_file); let mut content_updated = false; - // TODO merge both filters to one let series_info_count = fpl.playlistgroups.iter() .filter(|&plg| plg.xtream_cluster == XtreamCluster::Series) .flat_map(|plg| &plg.channels) .filter(|&pli| pli.header.item_type == PlaylistItemType::SeriesInfo).count(); - let series_info_iter = fpl.playlistgroups.iter_mut() - .filter(|plg| plg.xtream_cluster == XtreamCluster::Series) - .flat_map(|plg| &mut plg.channels) - .filter(|pli| pli.header.item_type == PlaylistItemType::SeriesInfo); - info!("Found {series_info_count} series info to resolve"); - let start_time = Instant::now(); + let mut last_log_time = Instant::now(); let mut processed_series_info_count = 0; - let mut last_processed_series_info_count = 0; - for pli in series_info_iter { - let (should_update, provider_id, ts) = should_update_series_info(pli, &processed_info_ids); - if should_update && provider_id != 0 && fetched_in_run.insert(provider_id) { - if let Some(content) = playlist_resolve_download_playlist_item(client, pli, fpl.input, errors, resolve_delay, XtreamCluster::Series).await { - let normalized_content = normalize_json_content(content); - handle_error_and_return!(write_series_info_to_wal_file(provider_id, ts, &normalized_content, &mut content_writer, &mut record_writer), - |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; - } + let mut write_counter = 0usize; + + for plg in &mut fpl.playlistgroups { + if plg.xtream_cluster != XtreamCluster::Series { + continue; } - if log_enabled!(Level::Info) { - processed_series_info_count += 1; - let elapsed = start_time.elapsed().as_secs(); - if elapsed > 0 && ((processed_series_info_count - last_processed_series_info_count) > 50) && elapsed.is_multiple_of(30) { - info!("resolved {processed_series_info_count}/{series_info_count} series info"); - last_processed_series_info_count = processed_series_info_count; + for pli in &mut plg.channels { + if pli.header.item_type != PlaylistItemType::SeriesInfo { + continue; + } + let (should_update, provider_id, ts) = should_update_series_info(pli, &processed_info_ids); + if should_update && provider_id != 0 && fetched_in_run.insert(provider_id) { + if let Some(content) = playlist_resolve_download_playlist_item(client, pli, fpl.input, errors, resolve_delay, XtreamCluster::Series).await { + let normalized_content = normalize_json_content(content); + let normalized_str = normalized_content.as_str(); + handle_error_and_return!(write_series_info_to_wal_file(provider_id, ts, normalized_str, &mut content_writer, &mut record_writer), + |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; + + // periodic flush to bound BufWriter memory + if write_counter.is_multiple_of(FLUSH_INTERVAL) { + if let Err(err) = content_writer.flush() { + errors.push(notify_err!(format!("Failed periodic flush of wal content writer {err}"))); + } + if let Err(err) = record_writer.flush() { + errors.push(notify_err!(format!("Failed periodic flush of wal record writer {err}"))); + } + } + } + } + if log_enabled!(Level::Info) { + processed_series_info_count += 1; + if last_log_time.elapsed().as_secs() >= 30 { + info!("resolved {processed_series_info_count}/{series_info_count} series info"); + last_log_time = Instant::now(); + } } } } - if last_processed_series_info_count != processed_series_info_count { - info!("resolved {processed_series_info_count}/{series_info_count} series info"); - } + info!("resolved {processed_series_info_count}/{series_info_count} series info"); // content_wal contains the provider_id and series_info with episode listing // record_wal contains provider_id and timestamp if content_updated { @@ -112,8 +126,8 @@ async fn playlist_resolve_series_info(cfg: &AppConfig, client: &reqwest::Client, 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); - drop(wal_content_file); drop(record_writer); + drop(wal_content_file); drop(wal_record_file); handle_error!(xtream_update_input_info_file(cfg, fpl.input, &wal_content_path, XtreamCluster::Series).await, |err| errors.push(err)); @@ -142,6 +156,8 @@ async fn process_series_info( return result; }; + let mut write_counter = 0usize; + let _file_lock = app_config.file_locks.read_lock(&info_path).await; // Contains the Series Info with episode listing @@ -184,6 +200,13 @@ async fn process_series_info( handle_error!(write_series_episode_record_to_wal_file(&mut wal_writer, *provider_id, episode), |err| errors.push(info_err!(format!("Failed to write to series episode wal file: {err}")))); } + write_counter +=1; + // periodic flush to bound BufWriter memory + if write_counter.is_multiple_of(FLUSH_INTERVAL) { + if let Err(err) = wal_writer.flush() { + errors.push(notify_err!(format!("Failed periodic flush of wal content writer {err}"))); + } + } group_series.extend(series.into_iter().map(|(_, pli)| pli)); } Ok(None) => {} @@ -205,8 +228,9 @@ async fn process_series_info( } } - handle_error!(wal_writer.flush(), - |err| errors.push(notify_err!(format!("Failed to resolve series episodes, could not write to wal file {err}")))); + handle_error!(wal_writer.flush(), |err| errors.push(notify_err!(format!("Failed to resolve series episodes, could not write to wal file {err}")))); + handle_error!(wal_writer.get_ref().sync_all(), |err| errors.push(notify_err!(format!("Failed to sync series info to wal file {err}")))); + drop(wal_writer); drop(wal_file); handle_error!(xtream_update_input_series_episodes_record_from_wal_file(app_config, input, &wal_path).await, diff --git a/backend/src/repository/bplustree.rs b/backend/src/repository/bplustree.rs index c5576d1bf..5b15f4eef 100644 --- a/backend/src/repository/bplustree.rs +++ b/backend/src/repository/bplustree.rs @@ -1,7 +1,7 @@ use crate::utils; use log::error; use ruzstd::decoding::StreamingDecoder; -use ruzstd::encoding::{compress_to_vec, CompressionLevel}; +use ruzstd::encoding::CompressionLevel; use serde::{Deserialize, Serialize}; use shared::error::{str_to_io_error, to_io_error}; use std::fs::File; @@ -40,6 +40,7 @@ fn u32_from_bytes(bytes: &[u8]) -> io::Result { } +#[inline] fn get_entry_index_upper_bound(keys: &[K], key: &K) -> usize where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone, @@ -58,15 +59,14 @@ where } -fn query_tree_le(file: &mut R, key: &K) -> Option +fn query_tree_le(file: &mut R, buffer: &mut Vec, key: &K) -> Option where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone, V: Serialize + for<'de> Deserialize<'de> + Clone, { let mut offset = 0; - let mut buffer = vec![0u8; BLOCK_SIZE]; loop { - match BPlusTreeNode::::deserialize_from_block(file, &mut buffer, offset, false) { + match BPlusTreeNode::::deserialize_from_block(file, buffer, offset, false) { Ok((node, pointers)) => { if node.is_leaf { let idx = get_entry_index_upper_bound::(&node.keys, key); @@ -272,131 +272,156 @@ where self.children.iter().for_each(|child| child.traverse(visit)); } - fn serialize_to_block(&self, file: &mut W, buffer: &mut Vec, offset: u64) -> io::Result { + fn serialize_to_block( + &self, + file: &mut W, + buffer: &mut Vec, // must be length BLOCK_SIZE + offset: u64, + ) -> io::Result { + // Keep backward-compatible on-disk layout with minimal allocations. let mut current_offset = offset; + + // zero the working block and set node type buffer.fill(0_u8); let buffer_slice = &mut buffer[..]; - - // Write node type (leaf or internal) buffer_slice[0] = u8::from(self.is_leaf); let mut write_pos = FLAG_SIZE; - // Serialize and write keys + // ---- Write keys (length + bytes) into the first block ---- let keys_encoded = bincode_serialize(&self.keys)?; - let keys_bytes_len = keys_encoded.len(); - buffer_slice[write_pos..write_pos + LEN_SIZE].copy_from_slice(&(u32::try_from(keys_bytes_len).map_err(to_io_error)?).to_le_bytes()); + let keys_len = keys_encoded.len(); + buffer_slice[write_pos..write_pos + LEN_SIZE] + .copy_from_slice(&u32::try_from(keys_len).map_err(to_io_error)?.to_le_bytes()); write_pos += LEN_SIZE; - buffer_slice[write_pos..write_pos + keys_bytes_len].copy_from_slice(&keys_encoded); - write_pos += keys_bytes_len; + // NOTE: By design of the legacy layout, keys are expected to fit into the first block. + buffer_slice[write_pos..write_pos + keys_len].copy_from_slice(&keys_encoded); + write_pos += keys_len; + drop(keys_encoded); - let mut remaining: Option> = None; - // If leaf, serialize and write values + // Prepare pointer offset for internal nodes (must be within first block) + let pointer_offset_within_first_block = if self.is_leaf { 0u64 } else { offset + write_pos as u64 }; + + // ---- Leaf values (optional) ---- if self.is_leaf { + // Encode values and decide compression exactly like the old layout let values_encoded = bincode_serialize(&self.values)?; let use_compression = values_encoded.len() + write_pos >= BLOCK_SIZE; - let compression_byte = if use_compression { 1u8.to_le_bytes() } else { 0u8.to_le_bytes() }; - buffer_slice[write_pos..=write_pos].copy_from_slice(&compression_byte); + + // Compression flag + buffer_slice[write_pos] = u8::from(use_compression); write_pos += FLAG_SIZE; + + // Content bytes (possibly compressed) let content_bytes = if use_compression { - { - // use ruzstd::io::Read; - compress_to_vec(&*values_encoded, CompressionLevel::Fastest) - } + // ruzstd expects a Read implementor; wrap the slice in a Cursor + ruzstd::encoding::compress_to_vec(std::io::Cursor::new(values_encoded.as_slice()), CompressionLevel::Fastest) } else { values_encoded }; - let values_bytes_len = content_bytes.len(); - buffer_slice[write_pos..write_pos + LEN_SIZE].copy_from_slice(&(u32::try_from(values_bytes_len).map_err(to_io_error)?).to_le_bytes()); + + // Write content length + let content_len = content_bytes.len(); + buffer_slice[write_pos..write_pos + LEN_SIZE] + .copy_from_slice(&u32::try_from(content_len).map_err(to_io_error)?.to_le_bytes()); write_pos += LEN_SIZE; - let mut value_bytes_to_write = values_bytes_len; - let bytes_to_write = write_pos + values_bytes_len; - if bytes_to_write > BLOCK_SIZE { - value_bytes_to_write = BLOCK_SIZE - write_pos; - remaining = Some(content_bytes[value_bytes_to_write..].to_vec()); + // Copy as many content bytes as fit into the first block + let space_left = BLOCK_SIZE.saturating_sub(write_pos); + let first_copy = std::cmp::min(space_left, content_len); + buffer_slice[write_pos..write_pos + first_copy] + .copy_from_slice(&content_bytes[..first_copy]); + + // Write the full first block + file.seek(SeekFrom::Start(offset))?; + file.write_all(&buffer_slice[..BLOCK_SIZE])?; + current_offset += BLOCK_SIZE as u64; + + // Stream remaining content (if any) block-by-block without extra allocations + let mut pos = first_copy; + while pos < content_len { + let remaining = content_len - pos; + let chunk = std::cmp::min(remaining, BLOCK_SIZE); + // Copy the next chunk + buffer_slice[..chunk].copy_from_slice(&content_bytes[pos..pos + chunk]); + // If this is a partial block, zero the tail only once + if chunk < BLOCK_SIZE { + buffer_slice[chunk..BLOCK_SIZE].fill(0u8); + } + file.write_all(&buffer_slice[..BLOCK_SIZE])?; // always write full block (pad with zeros) + current_offset += BLOCK_SIZE as u64; + pos += chunk; } - buffer_slice[write_pos..write_pos + value_bytes_to_write].copy_from_slice(&content_bytes[..value_bytes_to_write]); - // write_pos += value_bytes_to_write; + + // Free content buffer ASAP + drop(content_bytes); + } else { + // Internal node: write out the full first block now + file.seek(SeekFrom::Start(offset))?; + file.write_all(&buffer_slice[..BLOCK_SIZE])?; + current_offset += BLOCK_SIZE as u64; } - // Write the complete buffer to file, do not optimize to real filled size, - file.seek(SeekFrom::Start(offset))?; - file.write_all(&buffer_slice[..BLOCK_SIZE])?; // use BLOCK_SIZE - current_offset += BLOCK_SIZE as u64; - - if let Some(mut left_over) = remaining { - // we calculate the needed blocks - let left_over_len = left_over.len(); - if !left_over_len.is_multiple_of(BLOCK_SIZE) { - let padding = BLOCK_SIZE - (left_over_len % BLOCK_SIZE); - left_over.extend(vec![0u8; padding]); - } - let total_len = left_over.len(); - file.write_all(&left_over[..total_len])?; - current_offset += total_len as u64; - } - - // we need to write the pointers after the children are persisted, - // because we need the offsets, which we get after persisting. + // ---- Internal node pointers: serialize children, then write pointer list into first block ---- if !self.is_leaf { - let pointer_offset = offset + write_pos as u64; - let mut pointer = Vec::with_capacity(self.children.len()); + let pointer_offset = pointer_offset_within_first_block; // inside the node's first block + let mut pointer_vec: Vec = Vec::with_capacity(self.children.len()); for child in &self.children { - pointer.push(current_offset); + pointer_vec.push(current_offset); current_offset = child.serialize_to_block(file, buffer, current_offset)?; } - - let pointer_encoded = bincode_serialize(&pointer)?; - let pointer_bytes_len = u32::try_from(pointer_encoded.len()).map_err(to_io_error)?; - + let pointer_encoded = bincode_serialize(&pointer_vec)?; + let pointer_len = u32::try_from(pointer_encoded.len()).map_err(to_io_error)?; file.seek(SeekFrom::Start(pointer_offset))?; - file.write_all(&pointer_bytes_len.to_le_bytes())?; + file.write_all(&pointer_len.to_le_bytes())?; file.write_all(&pointer_encoded)?; } Ok(current_offset) } - fn deserialize_from_block(file: &mut R, buffer: &mut Vec, offset: u64, nested: bool) -> io::Result<(Self, Option>)> { + fn deserialize_from_block( + file: &mut R, + buffer: &mut Vec, + offset: u64, + nested: bool, + ) -> io::Result<(Self, Option>)> { + // Read the full first block into buffer (always aligned on-disk) file.seek(SeekFrom::Start(offset))?; file.read_exact(buffer)?; - // Read the node type directly from buffer + // Node type let is_leaf = buffer[0] == 1u8; let mut read_pos = FLAG_SIZE; - // Deserialize keys + // ---- Keys ---- let keys_length = u32_from_bytes(&buffer[read_pos..read_pos + LEN_SIZE])? as usize; read_pos += LEN_SIZE; let keys: Vec = bincode_deserialize(&buffer[read_pos..read_pos + keys_length])?; read_pos += keys_length; - // Deserialize values if leaf node - let values = if is_leaf { - let use_compression = u8::from_le_bytes(buffer[read_pos..=read_pos].try_into().unwrap_or([0u8])) == 1; + // ---- Values for leaf nodes ---- + let values: Vec = if is_leaf { + // compression flag + let use_compression = buffer[read_pos] == 1u8; read_pos += FLAG_SIZE; + + // content length let values_length = u32_from_bytes(&buffer[read_pos..read_pos + LEN_SIZE])? as usize; read_pos += LEN_SIZE; + // Fast path: everything is inside the first block let bytes_available_on_block = BLOCK_SIZE - read_pos; - let content_bytes = if values_length > bytes_available_on_block { - let mut left_over_bytes = values_length; - let mut content_chunk = Vec::from(&buffer[read_pos..read_pos + bytes_available_on_block]); - left_over_bytes -= bytes_available_on_block; - while left_over_bytes > 0 { - file.read_exact(buffer)?; - let bytes_to_read = if left_over_bytes >= BLOCK_SIZE { - BLOCK_SIZE - } else { - left_over_bytes - }; - content_chunk.extend(&buffer[..bytes_to_read]); - left_over_bytes -= bytes_to_read; - } - content_chunk[..values_length].to_vec() + let mut content_bytes: Vec = vec![0u8; values_length]; + if values_length <= bytes_available_on_block { + content_bytes[..values_length] + .copy_from_slice(&buffer[read_pos..read_pos + values_length]); } else { - buffer[read_pos..read_pos + values_length].to_vec() - }; + // Copy what's left in the first block + content_bytes[..bytes_available_on_block] + .copy_from_slice(&buffer[read_pos..read_pos + bytes_available_on_block]); + // Read the remaining bytes directly into the final buffer + file.read_exact(&mut content_bytes[bytes_available_on_block..])?; + } let values_bytes = if use_compression { decode_content(&content_bytes).unwrap_or(content_bytes) @@ -404,32 +429,27 @@ where content_bytes }; - let values: Vec = bincode_deserialize(&values_bytes)?; - read_pos += values_length; - values + bincode_deserialize(&values_bytes)? } else { - vec![] + Vec::new() }; - // Deserialize children indices if internal node - let (children, children_pointer) = if is_leaf { - (vec![], None) + // ---- Pointers for internal nodes ---- + let (children, children_pointer): (Vec, Option>) = if is_leaf { + (Vec::new(), None) } else { let pointers_length = u32_from_bytes(&buffer[read_pos..read_pos + LEN_SIZE])? as usize; read_pos += LEN_SIZE; let pointers: Vec = bincode_deserialize(&buffer[read_pos..read_pos + pointers_length])?; if nested { - let nodes: Result, io::Error> = pointers - .iter() - .map(|pointer| { - Self::deserialize_from_block(file, buffer, *pointer, nested) - .map(|(node, _)| node) - }) - .collect(); - - (nodes?, None) + let mut nodes = Vec::with_capacity(pointers.len()); + for &ptr in &pointers { + let (child, _) = Self::deserialize_from_block(file, buffer, ptr, nested)?; + nodes.push(child); + } + (nodes, None) } else { - (vec![], Some(pointers)) + (Vec::new(), Some(pointers)) } }; @@ -437,8 +457,8 @@ where } } -fn decode_content(content_bytes: &Vec) -> Option> { - if let Ok(mut decoder) = StreamingDecoder::new(&**content_bytes) { +fn decode_content(content_bytes: &[u8]) -> Option> { + if let Ok(mut decoder) = StreamingDecoder::new(content_bytes) { let mut result = Vec::with_capacity(content_bytes.len()); if decoder.read_to_end(&mut result).is_ok() { return Some(result); @@ -592,15 +612,14 @@ where } } -fn query_tree(file: &mut R, key: &K) -> Option +fn query_tree(file: &mut R, buffer: &mut Vec, key: &K) -> Option where K: Ord + Serialize + for<'de> Deserialize<'de> + Clone, V: Serialize + for<'de> Deserialize<'de> + Clone, { let mut offset = 0; - let mut buffer = vec![0u8; BLOCK_SIZE]; loop { - match BPlusTreeNode::::deserialize_from_block(file, &mut buffer, offset, false) { + match BPlusTreeNode::::deserialize_from_block(file, buffer, offset, false) { Ok((node, pointers)) => { if node.is_leaf { return match node.keys.binary_search(key) { @@ -655,6 +674,7 @@ where /// pub struct BPlusTreeQuery { file: BufReader, + buffer: Vec, _marker_k: PhantomData, _marker_v: PhantomData, } @@ -668,6 +688,7 @@ where let file = is_file_valid(file)?; Ok(Self { file: utils::file_reader(file), + buffer: vec![0u8; BLOCK_SIZE], _marker_k: PhantomData, _marker_v: PhantomData, }) @@ -679,7 +700,7 @@ where } pub fn query(&mut self, key: &K) -> Option { - query_tree(&mut self.file, key) + query_tree(&mut self.file, &mut self.buffer, key) } /// On-disk: find largest key <= `key` and return owned V (cloned/deserialized) @@ -692,7 +713,7 @@ where // if seek fails, still try to query — but bail out with None return None; } - query_tree_le(file, key) + query_tree_le(file, &mut self.buffer, key) } // pub fn traverse(&mut self, mut visit: F) @@ -705,6 +726,8 @@ where pub struct BPlusTreeUpdate { file: File, + read_buffer: Vec, + write_buffer: Vec, _marker_k: PhantomData, _marker_v: PhantomData, } @@ -721,6 +744,8 @@ where let file = is_file_valid(utils::open_read_write_file(filepath)?)?; Ok(Self { file, + read_buffer: vec![0u8; BLOCK_SIZE], + write_buffer: vec![0u8; BLOCK_SIZE], _marker_k: PhantomData, _marker_v: PhantomData, }) @@ -728,22 +753,20 @@ where pub fn query(&mut self, key: &K) -> Option { let mut reader = utils::file_reader(&mut self.file); - query_tree(&mut reader, key) + query_tree(&mut reader, &mut self.read_buffer, key) } fn serialize_node(&mut self, offset: u64, node: &BPlusTreeNode) -> io::Result { - let mut buffer = vec![0u8; BLOCK_SIZE]; - let result = node.serialize_to_block(&mut self.file, &mut buffer, offset); + let result = node.serialize_to_block(&mut self.file, &mut self.write_buffer, offset); self.file.flush()?; result } pub fn update(&mut self, key: &K, value: V) -> io::Result { let mut offset = 0; - let mut buffer = vec![0u8; BLOCK_SIZE]; let mut reader = utils::file_reader(&mut self.file); loop { - match BPlusTreeNode::::deserialize_from_block(&mut reader, &mut buffer, offset, false) { + match BPlusTreeNode::::deserialize_from_block(&mut reader, &mut self.read_buffer, offset, false) { Ok((mut node, pointers)) => { if node.is_leaf { return match node.keys.binary_search(key) { @@ -776,7 +799,7 @@ where /// On-disk update helper: find largest key <= `key`. pub fn query_le(&mut self, key: &K) -> Option { let mut reader = utils::file_reader(&mut self.file); - query_tree_le(&mut reader, key) + query_tree_le(&mut reader, &mut self.read_buffer, key) } } diff --git a/backend/src/repository/epg_repository.rs b/backend/src/repository/epg_repository.rs index a1872015b..baeabe6ae 100644 --- a/backend/src/repository/epg_repository.rs +++ b/backend/src/repository/epg_repository.rs @@ -34,8 +34,6 @@ pub async fn epg_write_file(target: &ConfigTarget, epg: &Epg, path: &Path) -> Re // EPG Content epg.write_to_async(&mut writer).await.map_err(|e| notify_err!(format!("failed to write epg: {}", e)))?; - let inner = writer.get_mut(); // Zugriff auf den BufWriter - inner.flush().await.map_err(|e| notify_err!(format!("failed to flush epg: {}", e)))?; debug_if_enabled!("Epg for target {} written to {}", target.name, path.to_str().unwrap_or("?")); Ok(()) diff --git a/backend/src/utils/file/file_utils.rs b/backend/src/utils/file/file_utils.rs index d59fc75c0..f53477cd4 100644 --- a/backend/src/utils/file/file_utils.rs +++ b/backend/src/utils/file/file_utils.rs @@ -1,10 +1,9 @@ use std::borrow::Cow; use std::collections::HashSet; use std::fs::{File, OpenOptions}; -use std::io::{BufReader, BufWriter, Read, Write}; use std::path::{Path, PathBuf}; use std::{env, fs}; - +use std::io::Read as IORead; 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}; @@ -12,20 +11,30 @@ use log::{debug, error}; use path_clean::PathClean; use tokio::fs as tokio_fs; -pub fn file_writer(w: W) -> BufWriter +const WRITER_BUFFER_SIZE: usize = 131_072; // 128kb + +pub fn file_writer(w: W) -> std::io::BufWriter where - W: Write, + W: std::io::Write, { - BufWriter::with_capacity(131_072, w) + std::io::BufWriter::with_capacity(WRITER_BUFFER_SIZE, w) } -pub fn file_reader(r: R) -> BufReader +pub fn file_reader(r: R) -> std::io::BufReader where - R: Read, + R: std::io::Read, { - BufReader::with_capacity(131_072, r) + std::io::BufReader::with_capacity(WRITER_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) +} + + pub fn get_exe_path() -> PathBuf { let default_path = std::path::PathBuf::from("./"); let current_exe = std::env::current_exe();