diff --git a/backend/src/processing/processor/mod.rs b/backend/src/processing/processor/mod.rs index 760e45bea..3f8e54cc3 100644 --- a/backend/src/processing/processor/mod.rs +++ b/backend/src/processing/processor/mod.rs @@ -30,6 +30,9 @@ macro_rules! handle_error_and_return { use handle_error_and_return; +// +// fn get_resolve__options(target: &ConfigTarget, fpl: &FetchedPlaylist) -> (bool, u16) +// #[macro_export] macro_rules! create_resolve_options_function_for_xtream_target { ($cluster:ident) => { diff --git a/backend/src/processing/processor/xtream.rs b/backend/src/processing/processor/xtream.rs index 9d18c423d..fdaccaa0f 100644 --- a/backend/src/processing/processor/xtream.rs +++ b/backend/src/processing/processor/xtream.rs @@ -7,7 +7,6 @@ use crate::model::normalize_release_date; use crate::repository::storage::get_input_storage_path; use serde::{Deserialize, Serialize}; use std::collections::HashMap; -use std::fs::File; use std::path::PathBuf; use crate::repository::bplustree::BPlusTree; use crate::repository::storage_const; @@ -47,18 +46,18 @@ pub(in crate::processing) fn normalize_json_content(content: String) -> String { } } -pub(in crate::processing) fn create_resolve_episode_wal_files(cfg: &Config, input: &ConfigInput) -> Option<(File, PathBuf)> { +pub(in crate::processing) async fn create_resolve_episode_wal_files(cfg: &Config, input: &ConfigInput) -> Option<(tokio::fs::File, PathBuf)> { match get_input_storage_path(&input.name, &cfg.working_dir) { Ok(storage_path) => { let info_path = storage_path.join(format!("{}.{}", crate::model::XC_FILE_SERIES_EPISODE_RECORD, storage_const::FILE_SUFFIX_WAL)); - let info_file = utils::append_or_crate_file(&info_path).ok()?; + let info_file = utils::async_append_or_create_file(&info_path).await.ok()?; Some((info_file, info_path)) } Err(_) => None } } -pub(in crate::processing) fn create_resolve_info_wal_files(cfg: &Config, input: &ConfigInput, cluster: XtreamCluster) -> Option<(File, File, PathBuf, PathBuf)> { +pub(in crate::processing) async fn create_resolve_info_wal_files(cfg: &Config, input: &ConfigInput, cluster: XtreamCluster) -> Option<(tokio::fs::File, tokio::fs::File, PathBuf, PathBuf)> { match get_input_storage_path(&input.name, &cfg.working_dir) { Ok(storage_path) => { if let Some(file_prefix) = match cluster { @@ -68,8 +67,8 @@ pub(in crate::processing) fn create_resolve_info_wal_files(cfg: &Config, input: } { let content_path = storage_path.join(format!("{file_prefix}_content.{}", storage_const::FILE_SUFFIX_WAL)); let info_path = storage_path.join(format!("{file_prefix}_record.{}", storage_const::FILE_SUFFIX_WAL)); - let content_file = utils::append_or_crate_file(&content_path).ok()?; - let info_file = utils::append_or_crate_file(&info_path).ok()?; + let content_file = utils::async_append_or_create_file(&content_path).await.ok()?; + let info_file = utils::async_append_or_create_file(&info_path).await.ok()?; return Some((content_file, info_file, content_path, info_path)); } None diff --git a/backend/src/processing/processor/xtream_series.rs b/backend/src/processing/processor/xtream_series.rs index e6fba65c0..ba13f5ce1 100644 --- a/backend/src/processing/processor/xtream_series.rs +++ b/backend/src/processing/processor/xtream_series.rs @@ -1,25 +1,24 @@ -use shared::model::InputType; -use shared::error::{TuliproxError}; use crate::model::{AppConfig, ConfigTarget}; -use crate::model::{FetchedPlaylist}; -use shared::model::{PlaylistGroup, PlaylistItem, PlaylistItemType, XtreamCluster}; -use crate::processing::processor::playlist::ProcessingPipe; +use crate::model::FetchedPlaylist; +use crate::model::{XtreamSeriesEpisode, XtreamSeriesInfoEpisode}; use crate::processing::parser::xtream::parse_xtream_series_info; +use crate::processing::processor::playlist::ProcessingPipe; +use crate::processing::processor::xtream::normalize_json_content; use crate::processing::processor::xtream::{create_resolve_episode_wal_files, create_resolve_info_wal_files, playlist_resolve_download_playlist_item, read_processed_info_ids, should_update_info}; +use crate::processing::processor::{create_resolve_options_function_for_xtream_target, handle_error, handle_error_and_return}; use crate::repository::storage::get_input_storage_path; use crate::repository::xtream_repository::{write_series_info_to_wal_file, xtream_get_info_file_paths, xtream_update_input_info_file, xtream_update_input_series_episodes_record_from_wal_file, xtream_update_input_series_record_from_wal_file}; use crate::repository::IndexedDocumentReader; -use shared::error::{notify_err, info_err}; -use crate::processing::processor::{handle_error, handle_error_and_return, create_resolve_options_function_for_xtream_target}; -use std::collections::{HashMap, HashSet}; -use std::fs::File; -use std::io::{BufWriter, Write}; -use std::time::Instant; -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, IO_BUFFER_SIZE}; +use log::{error, info, log_enabled, warn, Level}; +use shared::error::{info_err, notify_err}; +use shared::error::TuliproxError; +use shared::model::InputType; +use shared::model::{PlaylistGroup, PlaylistItem, PlaylistItemType, XtreamCluster}; +use std::collections::{HashMap, HashSet}; +use std::time::Instant; +use tokio::io::AsyncWriteExt; create_resolve_options_function_for_xtream_target!(series); @@ -27,19 +26,19 @@ async fn read_processed_series_info_ids(cfg: &AppConfig, errors: &mut Vec, +async fn write_series_episode_record_to_wal_file( + writer: &mut tokio::io::BufWriter<&mut tokio::fs::File>, provider_id: u32, episode: &XtreamSeriesInfoEpisode, ) -> std::io::Result { let series_episode = XtreamSeriesEpisode::from(episode); if let Ok(content_bytes) = bincode_serialize(&series_episode) { - writer.write_all(&provider_id.to_le_bytes())?; + writer.write_all(&provider_id.to_le_bytes()).await?; let content_len = content_bytes.len(); - if let Ok(len) = u32::try_from(content_len) { - writer.write_all(&len.to_le_bytes())?; - writer.write_all(&content_bytes)?; - return Ok(content_len + 4usize) + if let Ok(len) = u32::try_from(content_len) { + writer.write_all(&len.to_le_bytes()).await?; + writer.write_all(&content_bytes).await?; + return Ok(content_len + 4usize); } error!("Cant write to WAL file, content length exceeds u32"); } @@ -52,16 +51,20 @@ fn should_update_series_info(pli: &mut PlaylistItem, processed_provider_ids: &Ha async fn playlist_resolve_series_info(cfg: &AppConfig, client: &reqwest::Client, errors: &mut Vec, fpl: &mut FetchedPlaylist<'_>, resolve_delay: u16) -> bool { + // TODO read existing WAL File and import it to avoid duplicate requests + + 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. + // we can't 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. // We collect the content into a wal file and write it once we collected everything. - let Some((wal_content_file, wal_record_file, wal_content_path, wal_record_path)) = create_resolve_info_wal_files(&cfg.config.load(), fpl.input, XtreamCluster::Series) + let Some((mut wal_content_file, mut wal_record_file, wal_content_path, wal_record_path)) + = create_resolve_info_wal_files(&cfg.config.load(), fpl.input, XtreamCluster::Series).await else { return !processed_info_ids.is_empty(); }; - let mut content_writer = utils::file_writer(&wal_content_file); - let mut record_writer = utils::file_writer(&wal_record_file); + let mut content_writer = utils::async_file_writer(&mut wal_content_file); + let mut record_writer = utils::async_file_writer(&mut wal_record_file); let mut content_updated = false; let series_info_count = fpl.playlistgroups.iter() @@ -88,7 +91,7 @@ async fn playlist_resolve_series_info(cfg: &AppConfig, client: &reqwest::Client, 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), + handle_error_and_return!(write_series_info_to_wal_file(provider_id, ts, normalized_str, &mut content_writer, &mut record_writer).await, |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; @@ -97,10 +100,10 @@ async fn playlist_resolve_series_info(cfg: &AppConfig, client: &reqwest::Client, // periodic flush to bound BufWriter memory if write_counter >= IO_BUFFER_SIZE { write_counter = 0; - if let Err(err) = content_writer.flush() { + if let Err(err) = content_writer.flush().await { errors.push(notify_err!(format!("Failed periodic flush of wal content writer {err}"))); } - if let Err(err) = record_writer.flush() { + if let Err(err) = record_writer.flush().await { errors.push(notify_err!(format!("Failed periodic flush of wal record writer {err}"))); } } @@ -119,12 +122,12 @@ async fn playlist_resolve_series_info(cfg: &AppConfig, client: &reqwest::Client, // content_wal contains the provider_id and series_info with episode listing // record_wal contains provider_id and timestamp if content_updated { - handle_error!(content_writer.flush(), + handle_error!(content_writer.flush().await, |err| errors.push(notify_err!(format!("Failed to resolve series, could not write to wal file {err}")))); - handle_error!(record_writer.flush(), + handle_error!(record_writer.flush().await, |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}")))); + handle_error!(content_writer.get_ref().sync_all().await, |err| errors.push(notify_err!(format!("Failed to sync series info to wal file {err}")))); + handle_error!(record_writer.get_ref().sync_all().await, |err| errors.push(notify_err!(format!("Failed to sync series info record to wal file {err}")))); drop(content_writer); drop(record_writer); drop(wal_content_file); @@ -135,7 +138,6 @@ async fn playlist_resolve_series_info(cfg: &AppConfig, client: &reqwest::Client, |err| errors.push(err)); } - // TODO better approach for transactional updates is multiplexed WAL file. // we updated now // - series_info.db which contains the original series_info json // - series_record.db which contains the series_info provider_id and timestamp @@ -163,11 +165,11 @@ async fn process_series_info( // Contains the Series Info with episode listing let Ok(mut info_reader) = IndexedDocumentReader::::new(&info_path, &idx_path) else { return result; }; - let Some((wal_file, wal_path)) = create_resolve_episode_wal_files(&config, input) else { + let Some((mut wal_file, wal_path)) = create_resolve_episode_wal_files(&config, input).await else { errors.push(notify_err!("Could not create wal file for series episodes record".to_string())); return result; }; - let mut wal_writer = utils::file_writer(&wal_file); + let mut wal_writer = utils::async_file_writer(&mut wal_file); for plg in fpl .playlistgroups @@ -191,24 +193,24 @@ async fn process_series_info( Ok(series_content) => { let (group, series_name) = { let header = &pli.header; - (header.group.clone(), if header.name.is_empty() {header.title.clone()} else { header.name.clone()}) + (header.group.clone(), if header.name.is_empty() { header.title.clone() } else { header.name.clone() }) }; match parse_xtream_series_info(&series_content, &group, &series_name, input) { Ok(Some(mut series)) => { for (episode, pli_episode) in &mut series { let Some(provider_id) = &pli_episode.header.get_provider_id() else { continue; }; - match write_series_episode_record_to_wal_file(&mut wal_writer, *provider_id, episode) { + match write_series_episode_record_to_wal_file(&mut wal_writer, *provider_id, episode).await { Ok(written_bytes) => { write_counter += written_bytes; // periodic flush to bound BufWriter memory if write_counter >= IO_BUFFER_SIZE { write_counter = 0; - if let Err(err) = wal_writer.flush() { + if let Err(err) = wal_writer.flush().await { errors.push(notify_err!(format!("Failed periodic flush of wal content writer {err}"))); } } } - Err(err) => {errors.push(info_err!(format!("Failed to write to series episode wal file: {err}"))) } + Err(err) => { errors.push(info_err!(format!("Failed to write to series episode wal file: {err}"))) } } } group_series.extend(series.into_iter().map(|(_, pli)| pli)); @@ -232,8 +234,8 @@ 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.get_ref().sync_all(), |err| errors.push(notify_err!(format!("Failed to sync series info to wal file {err}")))); + handle_error!(wal_writer.flush().await, |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().await, |err| errors.push(notify_err!(format!("Failed to sync series info to wal file {err}")))); drop(wal_writer); drop(wal_file); diff --git a/backend/src/processing/processor/xtream_vod.rs b/backend/src/processing/processor/xtream_vod.rs index 18fdee3de..2d0a3e4cb 100644 --- a/backend/src/processing/processor/xtream_vod.rs +++ b/backend/src/processing/processor/xtream_vod.rs @@ -15,8 +15,8 @@ use shared::model::{InputType, PlaylistEntry}; use shared::model::{PlaylistItem, PlaylistItemType, XtreamCluster}; use shared::utils::{get_string_from_serde_value, get_u32_from_serde_value, get_u64_from_serde_value}; use std::collections::{HashMap, HashSet}; -use std::io::Write; use std::time::Instant; +use tokio::io::AsyncWriteExt; create_resolve_options_function_for_xtream_target!(vod); @@ -75,13 +75,14 @@ pub async fn playlist_resolve_vod(app_config: &AppConfig, client: &reqwest::Clie // All readers would be waiting for the lock and the app would be unresponsive. // We collect the content into a wal file and write it once we collected everything. let config = app_config.config.load(); - let Some((wal_content_file, wal_record_file, wal_content_path, wal_record_path)) = create_resolve_info_wal_files(&config, fpl.input, XtreamCluster::Video) + let Some((mut wal_content_file, mut wal_record_file, wal_content_path, wal_record_path)) + = create_resolve_info_wal_files(&config, fpl.input, XtreamCluster::Video).await else { return; }; let mut processed_info_ids: HashMap = read_processed_vod_info_ids(app_config, errors, fpl).await; let mut fetched_in_run: HashSet = HashSet::new(); - let mut content_writer = utils::file_writer(&wal_content_file); - let mut record_writer = utils::file_writer(&wal_record_file); + let mut content_writer = utils::async_file_writer(&mut wal_content_file); + let mut record_writer = utils::async_file_writer(&mut wal_record_file); let mut content_updated = false; // TODO merge both filters to one @@ -110,7 +111,7 @@ pub async fn playlist_resolve_vod(app_config: &AppConfig, client: &reqwest::Clie let normalized_str: &str = &normalized_content; if let Some((provider_id, info_record)) = extract_info_record_from_vod_info(normalized_str) { let ts = info_record.ts; - handle_error_and_return!(write_vod_info_to_wal_file(provider_id, normalized_str, &info_record, &mut content_writer, &mut record_writer), + handle_error_and_return!(write_vod_info_to_wal_file(provider_id, normalized_str, &info_record, &mut content_writer, &mut record_writer).await, |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; @@ -118,10 +119,10 @@ pub async fn playlist_resolve_vod(app_config: &AppConfig, client: &reqwest::Clie // periodic flush to bound BufWriter memory if write_counter >= IO_BUFFER_SIZE { write_counter = 0; - if let Err(err) = content_writer.flush() { + if let Err(err) = content_writer.flush().await { errors.push(notify_err!(format!("Failed periodic flush of wal content writer {err}"))); } - if let Err(err) = record_writer.flush() { + if let Err(err) = record_writer.flush().await { errors.push(notify_err!(format!("Failed periodic flush of wal record writer {err}"))); } } @@ -151,12 +152,12 @@ pub async fn playlist_resolve_vod(app_config: &AppConfig, client: &reqwest::Clie if content_updated { // TODO better approach for transactional updates is multiplexed WAL file. // final flush & sync with proper error handling - handle_error!(content_writer.flush(), + handle_error!(content_writer.flush().await, |err| errors.push(notify_err!(format!("Failed to resolve vod, could not write to wal file {err}")))); - handle_error!(record_writer.flush(), + handle_error!(record_writer.flush().await, |err| errors.push(notify_err!(format!("Failed to resolve vod 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 vod info to wal file {err}")))); - handle_error!(record_writer.get_ref().sync_all(), |err| errors.push(notify_err!(format!("Failed to sync vod info record to wal file {err}")))); + handle_error!(content_writer.get_ref().sync_all().await, |err| errors.push(notify_err!(format!("Failed to sync vod info to wal file {err}")))); + handle_error!(record_writer.get_ref().sync_all().await, |err| errors.push(notify_err!(format!("Failed to sync vod info record to wal file {err}")))); // drop writers and files to release handles drop(content_writer); drop(record_writer); diff --git a/backend/src/repository/xtream_repository.rs b/backend/src/repository/xtream_repository.rs index bee1153e0..b740e398e 100644 --- a/backend/src/repository/xtream_repository.rs +++ b/backend/src/repository/xtream_repository.rs @@ -10,9 +10,8 @@ use crate::repository::storage::{get_input_storage_path, get_target_id_mapping_f use crate::repository::storage_const; use crate::repository::target_id_mapping::VirtualIdRecord; use crate::repository::xtream_playlist_iterator::XtreamPlaylistJsonIterator; -use crate::utils::file_reader; +use crate::utils::{async_file_reader, async_open_readonly_file, file_reader}; use crate::utils::json_write_documents_to_file; -use crate::utils::open_readonly_file; use crate::utils::{bincode_deserialize, FileReadGuard}; use bytes::Bytes; use futures::{stream, Stream, StreamExt}; @@ -25,9 +24,10 @@ use shared::utils::{generate_playlist_uuid, get_u32_from_serde_value, hex_encode use std::collections::HashMap; use std::fs; use std::fs::File; -use std::io::{BufWriter, Error, ErrorKind, Read, Write}; +use std::io::{Error, ErrorKind}; use std::path::{Path, PathBuf}; use std::sync::Arc; +use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio::task; macro_rules! cant_write_result { @@ -868,20 +868,20 @@ pub async fn xtream_update_input_info_file( Ok(Some((info_path, idx_path))) => { { let _file_lock = cfg.file_locks.write_lock(&info_path).await; - let mut reader = file_reader(open_readonly_file(wal_path).map_err(|err| notify_err!(format!("Could not read {cluster} info {err}")))?); + let mut reader = async_file_reader(async_open_readonly_file(wal_path).await.map_err(|err| notify_err!(format!("Could not read {cluster} info {err}")))?); match IndexedDocumentWriter::::new_append(info_path.clone(), idx_path) { Ok(mut writer) => { let mut provider_id_bytes = [0u8; 4]; let mut length_bytes = [0u8; 4]; loop { - if reader.read_exact(&mut provider_id_bytes).is_err() { + if reader.read_exact(&mut provider_id_bytes).await.is_err() { break; // End of file } let provider_id = u32::from_le_bytes(provider_id_bytes); - reader.read_exact(&mut length_bytes).map_err(|err| notify_err!(format!("Could not read temporary {cluster} info {err}")))?; + reader.read_exact(&mut length_bytes).await.map_err(|err| notify_err!(format!("Could not read temporary {cluster} info {err}")))?; let length = u32::from_le_bytes(length_bytes) as usize; let mut buffer = vec![0u8; length]; - reader.read_exact(&mut buffer).map_err(|err| notify_err!(format!("Could not read temporary {cluster} info {err}")))?; + reader.read_exact(&mut buffer).await.map_err(|err| notify_err!(format!("Could not read temporary {cluster} info {err}")))?; if let Ok(content) = String::from_utf8(buffer) { let _ = writer.write_doc(provider_id, &content); } @@ -914,25 +914,25 @@ pub async fn xtream_update_input_vod_record_from_wal_file( { let _file_lock = cfg.file_locks.write_lock(&record_path).await; - let mut reader = file_reader(open_readonly_file(wal_path).map_err(|err| notify_err!(format!("Could not read vod wal info {err}")))?); + let mut reader = async_file_reader(async_open_readonly_file(wal_path).await.map_err(|err| notify_err!(format!("Could not read vod wal info {err}")))?); let mut provider_id_bytes = [0u8; 4]; let mut tmdb_id_bytes = [0u8; 4]; let mut ts_bytes = [0u8; 8]; let mut tree_record_index: BPlusTree = BPlusTree::load(&record_path).unwrap_or_else(|_| BPlusTree::new()); loop { - if reader.read_exact(&mut provider_id_bytes).is_err() { + if reader.read_exact(&mut provider_id_bytes).await.is_err() { break; // End of file } let provider_id = u32::from_le_bytes(provider_id_bytes); - if reader.read_exact(&mut tmdb_id_bytes).is_err() { + if reader.read_exact(&mut tmdb_id_bytes).await.is_err() { error!("Unexpected EOF after reading provider_id {provider_id} for VOD record."); break; } let tmdb_id = u32::from_le_bytes(tmdb_id_bytes); - if reader.read_exact(&mut ts_bytes).is_err() { + if reader.read_exact(&mut ts_bytes).await.is_err() { error!("Unexpected EOF after reading tmdb_id for VOD record with provider_id {provider_id}."); break; } @@ -940,7 +940,7 @@ pub async fn xtream_update_input_vod_record_from_wal_file( // Read the date string length as a 4-byte u32. let mut len_bytes = [0u8; 4]; - if reader.read_exact(&mut len_bytes).is_err() { + if reader.read_exact(&mut len_bytes).await.is_err() { error!("Unexpected EOF when reading release_date length for VOD record with provider_id {provider_id}."); break; } @@ -948,7 +948,7 @@ pub async fn xtream_update_input_vod_record_from_wal_file( let release_date = if len > 0 { let mut date_buffer = vec![0u8; len]; - if reader.read_exact(&mut date_buffer).is_err() { + if reader.read_exact(&mut date_buffer).await.is_err() { error!("Unexpected EOF when reading release_date string for VOD record with provider_id {provider_id}."); break; } @@ -964,7 +964,7 @@ pub async fn xtream_update_input_vod_record_from_wal_file( tree_record_index.store(&record_path).map_err(|err| notify_err!(format!("Could not store vod record info {err}")))?; drop(reader); - if let Err(err) = fs::remove_file(wal_path) { + if let Err(err) = tokio::fs::remove_file(wal_path).await { error!("Failed to delete record WAL file for vod {err}"); } Ok(()) @@ -982,16 +982,16 @@ pub async fn xtream_update_input_series_record_from_wal_file( .and_then(|opt| opt.ok_or_else(|| notify_err!(format!("Error accessing storage path for input: {}", &input.name))))?; { let _file_lock = cfg.file_locks.write_lock(&record_path).await; - let mut reader = file_reader(open_readonly_file(wal_path).map_err(|err| notify_err!(format!("Could not read series wal info {err}")))?); + let mut reader = async_file_reader(async_open_readonly_file(wal_path).await.map_err(|err| notify_err!(format!("Could not read series wal info {err}")))?); let mut provider_id_bytes = [0u8; 4]; let mut ts_bytes = [0u8; 8]; let mut tree_record_index: BPlusTree = BPlusTree::load(&record_path).unwrap_or_else(|_| BPlusTree::new()); loop { - if reader.read_exact(&mut provider_id_bytes).is_err() { + if reader.read_exact(&mut provider_id_bytes).await.is_err() { break; // End of file } let provider_id = u32::from_le_bytes(provider_id_bytes); - if reader.read_exact(&mut ts_bytes).is_err() { + if reader.read_exact(&mut ts_bytes).await.is_err() { break; // End of file } let ts = u64::from_le_bytes(ts_bytes); @@ -999,7 +999,7 @@ pub async fn xtream_update_input_series_record_from_wal_file( } tree_record_index.store(&record_path).map_err(|err| notify_err!(format!("Could not store series record info {err}")))?; drop(reader); - if let Err(err) = fs::remove_file(wal_path) { + if let Err(err) = tokio::fs::remove_file(wal_path).await { error!("Failed to delete record WAL file for series {err}"); } Ok(()) @@ -1017,17 +1017,17 @@ pub async fn xtream_update_input_series_episodes_record_from_wal_file( .and_then(|opt| opt.ok_or_else(|| notify_err!(format!("Error accessing storage path for input: {}", &input.name))))?; { let _file_lock = cfg.file_locks.write_lock(&record_path).await; - let mut reader = file_reader(open_readonly_file(wal_path).map_err(|err| notify_err!(format!("Could not read series episode wal info {err}")))?); + let mut reader = async_file_reader(async_open_readonly_file(wal_path).await.map_err(|err| notify_err!(format!("Could not read series episode wal info {err}")))?); let mut provider_id_bytes = [0u8; 4]; let mut len_bytes = [0u8; 4]; let mut tree_record_index: BPlusTree = BPlusTree::load(&record_path).unwrap_or_else(|_| BPlusTree::new()); let mut buffer = vec![0u8; 4096]; loop { - if reader.read_exact(&mut provider_id_bytes).is_err() { + if reader.read_exact(&mut provider_id_bytes).await.is_err() { break; // End of file } let provider_id = u32::from_le_bytes(provider_id_bytes); - if reader.read_exact(&mut len_bytes).is_err() { + if reader.read_exact(&mut len_bytes).await.is_err() { break; // End of file } let len = usize::try_from(u32::from_le_bytes(len_bytes)).unwrap_or(0); @@ -1037,7 +1037,7 @@ pub async fn xtream_update_input_series_episodes_record_from_wal_file( if len > buffer.len() { buffer = vec![0u8; len]; } - if reader.read_exact(&mut buffer[0..len]).is_err() { + if reader.read_exact(&mut buffer[0..len]).await.is_err() { break; } match bincode_deserialize(&buffer[0..len]) { @@ -1051,7 +1051,7 @@ pub async fn xtream_update_input_series_episodes_record_from_wal_file( } tree_record_index.store(&record_path).map_err(|err| notify_err!(format!("Could not store series episode record info {err}")))?; drop(reader); - if let Err(err) = fs::remove_file(wal_path) { + if let Err(err) = tokio::fs::remove_file(wal_path).await { error!("Failed to delete record WAL file for series episode {err}"); } Ok(()) @@ -1115,20 +1115,23 @@ pub(crate) async fn xtream_get_playlist_categories(config: &Config, target_name: None } -pub fn write_series_info_to_wal_file(provider_id: u32, ts: u64, content: &str, content_write: &mut BufWriter<&File>, record_writer: &mut BufWriter<&File>) -> std::io::Result<()> { +pub async fn write_series_info_to_wal_file(provider_id: u32, ts: u64, content: &str, + content_write: &mut tokio::io::BufWriter<&mut tokio::fs::File>, + record_writer: &mut tokio::io::BufWriter<&mut tokio::fs::File>) -> std::io::Result<()> { let encoded_content = encode_info_content_for_wal_file(provider_id, content)?; let encoded_record = encode_series_info_record_for_wal_file(provider_id, ts); - content_write.write_all(&encoded_content)?; - record_writer.write_all(&encoded_record)?; + content_write.write_all(&encoded_content).await?; + record_writer.write_all(&encoded_record).await?; Ok(()) } - -pub fn write_vod_info_to_wal_file(provider_id: u32, content: &str, info_record: &InputVodInfoRecord, content_write: &mut BufWriter<&File>, record_writer: &mut BufWriter<&File>) -> std::io::Result<()> { +pub async fn write_vod_info_to_wal_file(provider_id: u32, content: &str, info_record: &InputVodInfoRecord, + content_write: &mut tokio::io::BufWriter<&mut tokio::fs::File>, + record_writer: &mut tokio::io::BufWriter<&mut tokio::fs::File>) -> std::io::Result<()> { let encoded_content = encode_info_content_for_wal_file(provider_id, content)?; let encoded_record = encode_vod_info_record_for_wal_file(provider_id, info_record)?; - content_write.write_all(&encoded_content)?; - record_writer.write_all(&encoded_record)?; + content_write.write_all(&encoded_content).await?; + record_writer.write_all(&encoded_record).await?; Ok(()) } diff --git a/backend/src/utils/file/file_utils.rs b/backend/src/utils/file/file_utils.rs index 2341f2022..b7298f5d0 100644 --- a/backend/src/utils/file/file_utils.rs +++ b/backend/src/utils/file/file_utils.rs @@ -215,10 +215,15 @@ pub fn sanitize_filename(file_name: &str) -> String { } #[inline] -pub fn append_or_crate_file(path: &Path) -> std::io::Result { +pub fn append_or_create_file(path: &Path) -> std::io::Result { OpenOptions::new().create(true).append(true).open(path) } +#[inline] +pub async fn async_append_or_create_file(path: &Path) -> std::io::Result { + tokio::fs::OpenOptions::new().create(true).append(true).open(path).await +} + #[inline] pub async fn create_new_file_for_write(path: &Path) -> tokio::io::Result { tokio_fs::OpenOptions::new() @@ -244,6 +249,11 @@ pub fn open_readonly_file(path: &Path) -> std::io::Result { OpenOptions::new().read(true).write(false).truncate(false).create(false).open(path) } +#[inline] +pub async fn async_open_readonly_file(path: &Path) -> std::io::Result { + tokio::fs::OpenOptions::new().read(true).write(false).truncate(false).create(false).open(path).await +} + pub fn rename_or_copy(src: &Path, dest: &Path, remove_old: bool) -> std::io::Result<()> { // Try to rename the file if fs::rename(src, dest).is_err() {