From ecd8e5657f76c876501029a09db46dfce1df8650 Mon Sep 17 00:00:00 2001 From: DarkBreakpoint <243206744+DarkBreakpoint@users.noreply.github.com> Date: Wed, 19 Nov 2025 14:13:54 -0600 Subject: [PATCH 1/6] async cleanup --- CHANGELOG.md | 5 ++ backend/src/api/endpoints/xmltv_api.rs | 6 +-- backend/src/main.rs | 48 ++++++++++---------- backend/src/processing/processor/playlist.rs | 39 ++++++---------- 4 files changed, 46 insertions(+), 52 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 34c06c500..738c1794d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -8,6 +8,11 @@ - Shared stream burst buffer zero copy data buffer to reduce memory usage. - Added detailed shared-stream/buffer/provider logging to trace lag, cache persistence, and session/provider lifecycle events. - Connection registration failures now trigger an explicit disconnect so zombie sockets don’t linger. +- Playlist updates now use Tokio tasks instead of spinning up per-source threads/runtimes, reducing CPU and memory overhead during large syncs. +- XMLTV timeshift responses stream asynchronously end-to-end to keep the Axum runtime responsive. +- `main` now uses `#[tokio::main]`, removing the manual runtime boilerplate and keeping every branch async end-to-end. +- XMLTV timeshift responses stream asynchronously end-to-end to keep the Axum runtime responsive. +- Healthcheck CLI path now uses the async reqwest client so startup checks no longer block a dedicated thread. - Shared stream shutdown now drops registry locks before releasing provider handles to prevent cross-lock stalls. diff --git a/backend/src/api/endpoints/xmltv_api.rs b/backend/src/api/endpoints/xmltv_api.rs index 2529bb459..1787597b4 100644 --- a/backend/src/api/endpoints/xmltv_api.rs +++ b/backend/src/api/endpoints/xmltv_api.rs @@ -149,16 +149,16 @@ async fn serve_epg( match tokio::fs::File::open(epg_path).await { Ok(epg_file) => match parse_timeshift(user.epg_timeshift.as_ref()) { None => serve_file(epg_path, mime::TEXT_XML).await.into_response(), - Some(duration) => serve_epg_with_timeshift(epg_file, duration).into_response(), + Some(duration) => serve_epg_with_timeshift(epg_file, duration).await, }, Err(_) => get_empty_epg_response().into_response(), } } -fn serve_epg_with_timeshift( +async fn serve_epg_with_timeshift( epg_file: tokio::fs::File, offset_minutes: i32, -) -> impl axum::response::IntoResponse + Send { +) -> axum::response::Response { let reader = tokio::io::BufReader::new(epg_file); let (tx, rx) = tokio::io::duplex(8192); tokio::spawn(async move { diff --git a/backend/src/main.rs b/backend/src/main.rs index 13a79a69f..59931a19f 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -82,7 +82,8 @@ const BUILD_TIMESTAMP: &str = env!("VERGEN_BUILD_TIMESTAMP"); // #[export_name = "malloc_conf"] // pub static malloc_conf: &[u8] = b"lg_prof_interval:25,prof:true,prof_leak:true,prof_active:true,prof_prefix:/tmp/jeprof\0"; -fn main() { +#[tokio::main] +async fn main() { let args = Args::parse(); if args.genpwd { @@ -98,29 +99,25 @@ fn main() { init_logger(args.log_level.as_ref(), config_paths.config_file_path.as_str()); if args.healthcheck { - healthcheck(config_paths.config_file_path.as_str()); - return; + let healthy = healthcheck(config_paths.config_file_path.as_str()).await; + std::process::exit(if healthy { 0 } else { 1 }); } info!("Version: {VERSION}"); if let Some(bts) = BUILD_TIMESTAMP.to_string().parse::>().ok().map(|datetime| datetime.format("%Y-%m-%d %H:%M:%S %Z").to_string()) { info!("Build time: {bts}"); } - let rt = tokio::runtime::Runtime::new().unwrap(); - let () = rt.block_on(async { + let app_config = utils::read_initial_app_config(&mut config_paths, true, true, args.server).await.unwrap_or_else(|err| exit!("{}", err)); + print_info(&app_config); - let app_config = utils::read_initial_app_config(&mut config_paths, true, true, args.server).await.unwrap_or_else(|err| exit!("{}", err)); - print_info(&app_config); + let sources = > as Access>::load(&app_config.sources); + let targets = sources.validate_targets(args.target.as_ref()).unwrap_or_else(|err| exit!("{}", err)); - let sources = > as Access>::load(&app_config.sources); - let targets = sources.validate_targets(args.target.as_ref()).unwrap_or_else(|err| exit!("{}", err)); - - if args.server { - start_in_server_mode(Arc::new(app_config), Arc::new(targets)).await; - } else { - start_in_cli_mode(Arc::new(app_config), Arc::new(targets)).await; - } - }); + if args.server { + start_in_server_mode(Arc::new(app_config), Arc::new(targets)).await; + } else { + start_in_cli_mode(Arc::new(app_config), Arc::new(targets)).await; + } } fn print_info(app_config: &AppConfig) { @@ -176,17 +173,20 @@ async fn start_in_server_mode(cfg: Arc, targets: Arc) } } -fn healthcheck(config_file: &str) { +async fn healthcheck(config_file: &str) -> bool { let path = std::path::PathBuf::from(config_file); let file = File::open(path).expect("Failed to open config file"); let config: HealthcheckConfig = serde_yaml::from_reader(config_file_reader(file, true)).expect("Failed to parse config file"); - if let Ok(response) = reqwest::blocking::get(format!("http://localhost:{}/healthcheck", config.api.port)) { - if let Ok(check) = response.json::() { - if check.status == "ok" { - std::process::exit(0); - } - } + match reqwest::Client::new() + .get(format!("http://localhost:{}/healthcheck", config.api.port)) + .send() + .await + { + Ok(response) => match response.json::().await { + Ok(check) if check.status == "ok" => true, + _ => false, + }, + Err(_) => false, } - std::process::exit(1); } diff --git a/backend/src/processing/processor/playlist.rs b/backend/src/processing/processor/playlist.rs index 55df09775..03429e5ec 100644 --- a/backend/src/processing/processor/playlist.rs +++ b/backend/src/processing/processor/playlist.rs @@ -6,7 +6,7 @@ use crate::Config; use std::collections::{HashMap, HashSet}; use std::path::PathBuf; use std::sync::Arc; -use std::thread; +use tokio::task::JoinSet; use tokio::sync::Mutex; use crate::api::model::{EventManager, EventMessage, PlaylistStorageState}; @@ -385,7 +385,7 @@ fn create_input_stat(group_count: usize, channel_count: usize, error_count: usiz async fn process_sources(client: Arc, config: &Arc, user_targets: Arc, event_manager: Option>, playlist_state: Option<&Arc>, ) -> (Vec, Vec) { - let mut handle_list = vec![]; + let mut async_tasks = JoinSet::new(); let thread_num = config.config.load().threads; let sources = config.sources.load(); let process_parallel = thread_num > 1 && sources.sources.len() > 1; @@ -410,27 +410,14 @@ async fn process_sources(client: Arc, config: &Arc, if process_parallel { let http_client = Arc::clone(&client); let playlist_state = playlist_state.cloned(); - let handles = &mut handle_list; - let process = move || { - // TODO better way ? - match tokio::runtime::Runtime::new() { - Ok(rt) => { - rt.block_on(async { - let (input_stats, target_stats, mut res_errors) = - process_source(Arc::clone(&http_client), cfg, index, usr_trgts, event_manager, playlist_state.as_ref()).await; - shared_errors.lock().await.append(&mut res_errors); - if let Some(process_stats) = SourceStats::try_new(input_stats, target_stats) { - shared_stats.lock().await.push(process_stats); - } - }); - } - Err(err) => error!("Could not create runtime !!! {err}"), + async_tasks.spawn(async move { + let (input_stats, target_stats, mut res_errors) = + process_source(Arc::clone(&http_client), cfg, index, usr_trgts, event_manager, playlist_state.as_ref()).await; + shared_errors.lock().await.append(&mut res_errors); + if let Some(process_stats) = SourceStats::try_new(input_stats, target_stats) { + shared_stats.lock().await.push(process_stats); } - }; - handles.push(thread::spawn(process)); - if handles.len() >= thread_num as usize { - handles.drain(..).for_each(|handle| { let _ = handle.join(); }); - } + }); } else { let (input_stats, target_stats, mut res_errors) = process_source(Arc::clone(&client), cfg, index, usr_trgts, event_manager, playlist_state).await; @@ -441,8 +428,10 @@ async fn process_sources(client: Arc, config: &Arc, } drop(update_lock); } - for handle in handle_list { - let _ = handle.join(); + while let Some(result) = async_tasks.join_next().await { + if let Err(err) = result { + error!("Playlist processing task failed: {err:?}"); + } } if let (Ok(s), Ok(e)) = (Arc::try_unwrap(stats), Arc::try_unwrap(errors)) { (s.into_inner(), e.into_inner()) @@ -700,4 +689,4 @@ pub async fn exec_processing(client: Arc, app_config: Arc Date: Wed, 19 Nov 2025 14:41:57 -0600 Subject: [PATCH 2/6] Video download async cleanup --- CHANGELOG.md | 2 +- backend/src/api/endpoints/download_api.rs | 14 +++++++------- 2 files changed, 8 insertions(+), 8 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 738c1794d..3fc8ee6b8 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -8,10 +8,10 @@ - Shared stream burst buffer zero copy data buffer to reduce memory usage. - Added detailed shared-stream/buffer/provider logging to trace lag, cache persistence, and session/provider lifecycle events. - Connection registration failures now trigger an explicit disconnect so zombie sockets don’t linger. +- Video download queue now uses async file I/O to keep the runtime responsive during large transfers. - Playlist updates now use Tokio tasks instead of spinning up per-source threads/runtimes, reducing CPU and memory overhead during large syncs. - XMLTV timeshift responses stream asynchronously end-to-end to keep the Axum runtime responsive. - `main` now uses `#[tokio::main]`, removing the manual runtime boilerplate and keeping every branch async end-to-end. -- XMLTV timeshift responses stream asynchronously end-to-end to keep the Axum runtime responsive. - Healthcheck CLI path now uses the async reqwest client so startup checks no longer block a dedicated thread. - Shared stream shutdown now drops registry locks before releasing provider handles to prevent cross-lock stalls. diff --git a/backend/src/api/endpoints/download_api.rs b/backend/src/api/endpoints/download_api.rs index bfc18180b..3ec401d1f 100644 --- a/backend/src/api/endpoints/download_api.rs +++ b/backend/src/api/endpoints/download_api.rs @@ -6,11 +6,11 @@ use tokio::sync::RwLock; use futures::stream::TryStreamExt; use log::info; use serde_json::{json, Value}; -use std::fs::File; -use std::io::{Write}; +use tokio::fs::File; +use tokio::io::AsyncWriteExt; use std::ops::Deref; use std::sync::Arc; -use std::{fs}; +use tokio::fs; use axum::response::IntoResponse; use shared::utils::bytes_to_megabytes; use shared::error::to_io_error; @@ -20,11 +20,11 @@ async fn download_file(active: Arc>>, client: &reqwe if let Some(file_download) = active.read().await.as_ref().as_ref() { match client.get(file_download.url.clone()).send().await { Ok(response) => { - match fs::create_dir_all(&file_download.file_dir) { + match fs::create_dir_all(&file_download.file_dir).await { Ok(()) => { if let Some(file_path_str) = file_download.file_path.to_str() { info!("Downloading {file_path_str}"); - match File::create(&file_download.file_path) { + match File::create(&file_download.file_path).await { Ok(mut file) => { let mut downloaded: u64 = 0; let mut stream = response.bytes_stream().map_err(to_io_error); @@ -32,7 +32,7 @@ async fn download_file(active: Arc>>, client: &reqwe match stream.try_next().await { Ok(item) => { if let Some(chunk) = item { - match file.write_all(&chunk) { + match file.write_all(&chunk).await { Ok(()) => { downloaded += chunk.len() as u64; if let Some(lock) = active.write().await.as_mut() { @@ -168,4 +168,4 @@ pub async fn download_file_info( })), |file_download| axum::Json(json!({ "completed": false, "downloads": finished_list, "active": download_info!(file_download) }))) -} \ No newline at end of file +} From 8f6c49c73f006c2d6192f0ee1efb28f9647a679f Mon Sep 17 00:00:00 2001 From: DarkBreakpoint <243206744+DarkBreakpoint@users.noreply.github.com> Date: Wed, 19 Nov 2025 14:59:52 -0600 Subject: [PATCH 3/6] Async Json writer --- CHANGELOG.md | 1 + backend/src/repository/user_repository.rs | 12 ++++++------ backend/src/repository/xtream_repository.rs | 2 +- backend/src/utils/json_utils.rs | 16 +++++++++------- 4 files changed, 17 insertions(+), 14 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 3fc8ee6b8..4a3ad0df6 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,7 @@ - Added detailed shared-stream/buffer/provider logging to trace lag, cache persistence, and session/provider lifecycle events. - Connection registration failures now trigger an explicit disconnect so zombie sockets don’t linger. - Video download queue now uses async file I/O to keep the runtime responsive during large transfers. +- JSON playlist/category writers (xtream collections, user bouquets) now stream through Tokio I/O so persisting those files no longer blocks the runtime. - Playlist updates now use Tokio tasks instead of spinning up per-source threads/runtimes, reducing CPU and memory overhead during large syncs. - XMLTV timeshift responses stream asynchronously end-to-end to keep the Axum runtime responsive. - `main` now uses `#[tokio::main]`, removing the manual runtime boilerplate and keeping every branch async end-to-end. diff --git a/backend/src/repository/user_repository.rs b/backend/src/repository/user_repository.rs index 20f1a52a7..645ec42da 100644 --- a/backend/src/repository/user_repository.rs +++ b/backend/src/repository/user_repository.rs @@ -243,7 +243,7 @@ async fn save_xtream_user_bouquet_for_target(config: &Config, target_name: &str, if let Some(bouquet_categories) = bouquet { if let Some(xtream_categories) = xtream_get_playlist_categories(config, target_name, cluster).await { let filtered: Vec<&PlaylistXtreamCategory> = xtream_categories.iter().filter(|p| bouquet_categories.contains(&p.name)).collect(); - return json_write_documents_to_file(&bouquet_path, &filtered); + return json_write_documents_to_file(&bouquet_path, &filtered).await; } } @@ -253,7 +253,7 @@ async fn save_xtream_user_bouquet_for_target(config: &Config, target_name: &str, Ok(()) } -fn save_m3u_user_bouquet_for_target(storage_path: &Path, target: TargetType, cluster: XtreamCluster, bouquet: Option<&Vec>) -> Result<(), Error> { +async fn save_m3u_user_bouquet_for_target(storage_path: &Path, target: TargetType, cluster: XtreamCluster, bouquet: Option<&Vec>) -> Result<(), Error> { let bouquet_path = match cluster { XtreamCluster::Live => user_get_live_bouquet_path(storage_path, target), XtreamCluster::Video => user_get_vod_bouquet_path(storage_path, target), @@ -261,7 +261,7 @@ fn save_m3u_user_bouquet_for_target(storage_path: &Path, target: TargetType, clu }; match bouquet { Some(bouquet_categories) => { - json_write_documents_to_file(&bouquet_path, bouquet_categories)?; + json_write_documents_to_file(&bouquet_path, bouquet_categories).await?; } None => if bouquet_path.exists() { std::fs::remove_file(bouquet_path)?; @@ -277,9 +277,9 @@ async fn save_user_bouquet_for_target(config: &Config, target_name: &str, storag save_xtream_user_bouquet_for_target(config, target_name, storage_path, XtreamCluster::Video, bouquet.vod.as_ref()).await?; save_xtream_user_bouquet_for_target(config, target_name, storage_path, XtreamCluster::Series, bouquet.series.as_ref()).await?; } else { - save_m3u_user_bouquet_for_target(storage_path, target, XtreamCluster::Live, bouquet.live.as_ref())?; - save_m3u_user_bouquet_for_target(storage_path, target, XtreamCluster::Video, bouquet.vod.as_ref())?; - save_m3u_user_bouquet_for_target(storage_path, target, XtreamCluster::Series, bouquet.series.as_ref())?; + save_m3u_user_bouquet_for_target(storage_path, target, XtreamCluster::Live, bouquet.live.as_ref()).await?; + save_m3u_user_bouquet_for_target(storage_path, target, XtreamCluster::Video, bouquet.vod.as_ref()).await?; + save_m3u_user_bouquet_for_target(storage_path, target, XtreamCluster::Series, bouquet.series.as_ref()).await?; } Ok(()) } diff --git a/backend/src/repository/xtream_repository.rs b/backend/src/repository/xtream_repository.rs index ad741a087..a8b4374b2 100644 --- a/backend/src/repository/xtream_repository.rs +++ b/backend/src/repository/xtream_repository.rs @@ -257,7 +257,7 @@ pub async fn xtream_write_playlist( (get_collection_path(&path, storage_const::COL_CAT_VOD), &cat_vod_col), (get_collection_path(&path, storage_const::COL_CAT_SERIES), &cat_series_col), ] { - match json_write_documents_to_file(&col_path, data) { + match json_write_documents_to_file(&col_path, data).await { Ok(()) => {} Err(err) => { errors.push(format!("Persisting collection failed: {}: {err}", col_path.display())); diff --git a/backend/src/utils/json_utils.rs b/backend/src/utils/json_utils.rs index 2c57bf350..24f2037ad 100644 --- a/backend/src/utils/json_utils.rs +++ b/backend/src/utils/json_utils.rs @@ -1,11 +1,13 @@ use std::collections::{HashMap, HashSet}; use std::fs::File; -use std::io::{BufReader, Error, Write}; +use std::io::{BufReader, Error, ErrorKind}; use std::path::Path; use serde::Serialize; use serde_json::Value; use shared::utils::json_iter_array; -use crate::utils::{file_reader, file_writer}; +use crate::utils::file_reader; +use tokio::fs; +use tokio::io::AsyncWriteExt; pub fn json_filter_file(file_path: &Path, filter: &HashMap<&str, HashSet, S>) -> Vec { let mut filtered: Vec = Vec::with_capacity(1024); @@ -35,12 +37,12 @@ pub fn json_filter_file(file_path: &Path, filter: & filtered } -pub fn json_write_documents_to_file(file: &Path, value: &T) -> Result<(), Error> +pub async fn json_write_documents_to_file(file: &Path, value: &T) -> Result<(), Error> where T: ?Sized + Serialize, { - let file = File::create(file)?; - let mut writer = file_writer(&file); - serde_json::to_writer(&mut writer, value)?; - writer.flush() + let mut file = fs::File::create(file).await?; + let payload = serde_json::to_vec(value).map_err(|err| Error::new(ErrorKind::Other, err))?; + file.write_all(&payload).await?; + file.flush().await } From 000df50b76dc11754e37b4c032f36af310a12d5e Mon Sep 17 00:00:00 2001 From: DarkBreakpoint <243206744+DarkBreakpoint@users.noreply.github.com> Date: Wed, 19 Nov 2025 15:07:30 -0600 Subject: [PATCH 4/6] EPG export async --- CHANGELOG.md | 1 + backend/src/repository/epg_repository.rs | 17 +++++++++-------- backend/src/repository/playlist_repository.rs | 4 ++-- backend/src/utils/network/request.rs | 12 ++++++------ 4 files changed, 18 insertions(+), 16 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 4a3ad0df6..2f67e07f5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -10,6 +10,7 @@ - Connection registration failures now trigger an explicit disconnect so zombie sockets don’t linger. - Video download queue now uses async file I/O to keep the runtime responsive during large transfers. - JSON playlist/category writers (xtream collections, user bouquets) now stream through Tokio I/O so persisting those files no longer blocks the runtime. +- Playlist EPG exports now write via async file handles to avoid blocking the runtime during XML serialization. - Playlist updates now use Tokio tasks instead of spinning up per-source threads/runtimes, reducing CPU and memory overhead during large syncs. - XMLTV timeshift responses stream asynchronously end-to-end to keep the Axum runtime responsive. - `main` now uses `#[tokio::main]`, removing the manual runtime boilerplate and keeping every branch async end-to-end. diff --git a/backend/src/repository/epg_repository.rs b/backend/src/repository/epg_repository.rs index 74ae3c513..b1dce2d48 100644 --- a/backend/src/repository/epg_repository.rs +++ b/backend/src/repository/epg_repository.rs @@ -5,22 +5,23 @@ 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 quick_xml::Writer; -use std::fs::File; use std::io::{Cursor, Write}; use std::path::Path; +use tokio::fs::File; +use tokio::io::AsyncWriteExt; -fn epg_write_file(target: &ConfigTarget, epg: &Epg, path: &Path) -> Result<(), TuliproxError> { +async fn epg_write_file(target: &ConfigTarget, epg: &Epg, path: &Path) -> Result<(), TuliproxError> { let mut writer = Writer::new(Cursor::new(vec![])); match epg.write_to(&mut writer) { Ok(()) => { let result = writer.into_inner().into_inner(); - match File::create(path) { + match File::create(path).await { Ok(mut epg_file) => { - match epg_file.write_all("".as_bytes()) { + match epg_file.write_all("".as_bytes()).await { Ok(()) => {} Err(err) => return Err(notify_err!(format!("failed to write epg: {} - {}", path.to_str().unwrap_or("?"), err))), } - match epg_file.write_all(&result) { + match epg_file.write_all(&result).await { Ok(()) => { debug_if_enabled!("Epg for target {} written to {}", target.name, path.to_str().unwrap_or("?")); } @@ -35,7 +36,7 @@ fn epg_write_file(target: &ConfigTarget, epg: &Epg, path: &Path) -> Result<(), T Ok(()) } -pub fn epg_write(cfg: &Config, target: &ConfigTarget, target_path: &Path, epg: Option<&Epg>, output: &TargetOutput) -> Result<(), TuliproxError> { +pub async fn epg_write(cfg: &Config, target: &ConfigTarget, target_path: &Path, epg: Option<&Epg>, output: &TargetOutput) -> Result<(), TuliproxError> { if let Some(epg_data) = epg { match output { TargetOutput::Xtream(_) => { @@ -43,7 +44,7 @@ pub fn epg_write(cfg: &Config, target: &ConfigTarget, target_path: &Path, epg: O Some(path) => { let epg_path = xtream_get_epg_file_path(&path); debug_if_enabled!("writing xtream epg to {}", epg_path.to_str().unwrap_or("?")); - epg_write_file(target, epg_data, &epg_path)?; + 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))), } @@ -51,7 +52,7 @@ pub fn epg_write(cfg: &Config, target: &ConfigTarget, target_path: &Path, epg: O TargetOutput::M3u(_) => { let path = m3u_get_epg_file_path(target_path); debug_if_enabled!("writing m3u epg to {}", path.to_str().unwrap_or("?")); - epg_write_file(target, epg_data, &path)?; + epg_write_file(target, epg_data, &path).await?; } TargetOutput::Strm(_) | TargetOutput::HdHomeRun(_) => {} } diff --git a/backend/src/repository/playlist_repository.rs b/backend/src/repository/playlist_repository.rs index eaeb154d3..5708113ae 100644 --- a/backend/src/repository/playlist_repository.rs +++ b/backend/src/repository/playlist_repository.rs @@ -79,7 +79,7 @@ pub async fn persist_playlist(app_config: &AppConfig, playlist: &mut [PlaylistGr match result { Ok(()) => { if !playlist.is_empty() { - if let Err(err) = epg_write(config, target, &target_path, epg, output) { + if let Err(err) = epg_write(config, target, &target_path, epg, output).await { errors.push(err); } } @@ -220,4 +220,4 @@ pub async fn load_target_into_memory_cache(app_state: &AppState, target: &Arc, input: &Config Ok(response) => { if response.status().is_success() { // Open a file in write mode - let mut file = BufWriter::with_capacity(8192, File::create(file_path)?); + let mut file = BufWriter::with_capacity(8192, File::create(file_path).await?); // 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)?; + file.write_all(&bytes).await?; } Err(err) => { return Err(str_to_io_error(&format!("Failed to read chunk: {err}"))); @@ -258,7 +258,7 @@ async fn get_remote_content_as_file(client: Arc, input: &Config } } - file.flush()?; + file.flush().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()) From 0876ebc43b2e78a8e74374533b4143ffb5dc8c0a Mon Sep 17 00:00:00 2001 From: DarkBreakpoint <243206744+DarkBreakpoint@users.noreply.github.com> Date: Wed, 19 Nov 2025 15:13:33 -0600 Subject: [PATCH 5/6] async Config --- CHANGELOG.md | 1 + backend/src/api/endpoints/v1_api_config.rs | 12 ++++++------ backend/src/model/config/api_proxy.rs | 4 ++-- backend/src/utils/file/config_reader.rs | 22 +++++++++++++--------- 4 files changed, 22 insertions(+), 17 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 2f67e07f5..f90eb7e37 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -11,6 +11,7 @@ - Video download queue now uses async file I/O to keep the runtime responsive during large transfers. - JSON playlist/category writers (xtream collections, user bouquets) now stream through Tokio I/O so persisting those files no longer blocks the runtime. - Playlist EPG exports now write via async file handles to avoid blocking the runtime during XML serialization. +- Config and API proxy save endpoints now serialize via Tokio I/O, so editing configs through the API no longer blocks the runtime threads. - Playlist updates now use Tokio tasks instead of spinning up per-source threads/runtimes, reducing CPU and memory overhead during large syncs. - XMLTV timeshift responses stream asynchronously end-to-end to keep the Axum runtime responsive. - `main` now uses `#[tokio::main]`, removing the manual runtime boilerplate and keeping every branch async end-to-end. diff --git a/backend/src/api/endpoints/v1_api_config.rs b/backend/src/api/endpoints/v1_api_config.rs index 5873383ff..3bfb32e52 100644 --- a/backend/src/api/endpoints/v1_api_config.rs +++ b/backend/src/api/endpoints/v1_api_config.rs @@ -12,8 +12,8 @@ use crate::{utils}; use crate::utils::{prepare_sources_batch, prepare_users}; use crate::utils::request::download_text_content; -pub(in crate::api::endpoints) fn intern_save_config_api_proxy(backup_dir: &str, api_proxy: &ApiProxyConfigDto, file_path: &str) -> Option { - match utils::save_api_proxy(file_path, backup_dir, api_proxy) { +pub(in crate::api::endpoints) async fn intern_save_config_api_proxy(backup_dir: &str, api_proxy: &ApiProxyConfigDto, file_path: &str) -> Option { + match utils::save_api_proxy(file_path, backup_dir, api_proxy).await { Ok(()) => {} Err(err) => { error!("Failed to save api_proxy.yml {err}"); @@ -23,8 +23,8 @@ pub(in crate::api::endpoints) fn intern_save_config_api_proxy(backup_dir: &str, None } -fn intern_save_config_main(file_path: &str, backup_dir: &str, cfg: &ConfigDto) -> Option { - match utils::save_main_config(file_path, backup_dir, cfg) { +async fn intern_save_config_main(file_path: &str, backup_dir: &str, cfg: &ConfigDto) -> Option { + match utils::save_main_config(file_path, backup_dir, cfg).await { Ok(()) => {} Err(err) => { error!("Failed to save config.yml {err}"); @@ -43,7 +43,7 @@ async fn save_config_main( let file_path = paths.config_file_path.as_str(); let config = app_state.app_config.config.load(); let backup_dir = config.get_backup_dir(); - if let Some(err) = intern_save_config_main(file_path, backup_dir.as_ref(), &cfg) { + if let Some(err) = intern_save_config_main(file_path, backup_dir.as_ref(), &cfg).await { return (axum::http::StatusCode::INTERNAL_SERVER_ERROR, axum::Json(json!({"error": err.to_string()}))).into_response(); } axum::http::StatusCode::OK.into_response() @@ -75,7 +75,7 @@ async fn save_config_api_proxy_config( let backup_dir = config.get_backup_dir(); let paths = app_state.app_config.paths.load(); - if let Some(err) = intern_save_config_api_proxy(backup_dir.as_ref(), &ApiProxyConfigDto::from(&updated_api_proxy), paths.api_proxy_file_path.as_str()) { + if let Some(err) = intern_save_config_api_proxy(backup_dir.as_ref(), &ApiProxyConfigDto::from(&updated_api_proxy), paths.api_proxy_file_path.as_str()).await { return (axum::http::StatusCode::INTERNAL_SERVER_ERROR, axum::Json(json!({"error": err.to_string()}))).into_response(); } // Persist succeeded — now update in‑memory state diff --git a/backend/src/model/config/api_proxy.rs b/backend/src/model/config/api_proxy.rs index 108d3e43a..47c130cab 100644 --- a/backend/src/model/config/api_proxy.rs +++ b/backend/src/model/config/api_proxy.rs @@ -113,7 +113,7 @@ impl ApiProxyConfig { let config = > as Access>::load(&cfg.config); let backup_dir = config.get_backup_dir(); self.user = vec![]; - if let Err(err) = utils::save_api_proxy(api_proxy_file, backup_dir.as_ref(), &ApiProxyConfigDto::from(&*self)) { + if let Err(err) = utils::save_api_proxy(api_proxy_file, backup_dir.as_ref(), &ApiProxyConfigDto::from(&*self)).await { errors.push(format!("Error saving api proxy file: {err}")); } } @@ -148,7 +148,7 @@ impl ApiProxyConfig { let config = > as Access>::load(&cfg.config); let backup_dir = config.get_backup_dir(); - if let Err(err) = save_api_proxy(api_proxy_file, backup_dir.as_ref(), &ApiProxyConfigDto::from(&*self)) { + if let Err(err) = save_api_proxy(api_proxy_file, backup_dir.as_ref(), &ApiProxyConfigDto::from(&*self)).await { errors.push(format!("Error saving api proxy file: {err}")); } else { backup_api_user_db_file(cfg, &user_db_path).await; diff --git a/backend/src/utils/file/config_reader.rs b/backend/src/utils/file/config_reader.rs index e7dc8855f..02630eb59 100644 --- a/backend/src/utils/file/config_reader.rs +++ b/backend/src/utils/file/config_reader.rs @@ -18,6 +18,7 @@ use std::path::PathBuf; use std::sync::Arc; use arc_swap::access::{Access}; use crate::repository::user_repository::{get_api_user_db_path, load_api_user}; +use tokio::fs; enum EitherReader { Left(L), @@ -301,7 +302,7 @@ pub async fn read_api_proxy(config: &AppConfig, resolve_env: bool) -> Option(file_path: &str, backup_dir: &str, config: &T, default_name: &str) -> Result<(), TuliproxError> +async fn write_config_file(file_path: &str, backup_dir: &str, config: &T, default_name: &str) -> Result<(), TuliproxError> where T: ?Sized + Serialize, { @@ -310,23 +311,26 @@ where let backup_path = PathBuf::from(backup_dir).join(format!("{filename}_{}", Local::now().format("%Y%m%d_%H%M%S"))); - match std::fs::copy(&path, &backup_path) { + match fs::copy(&path, &backup_path).await { Ok(_) => {} Err(err) => { error!("Could not backup file {}:{}", &backup_path.to_str().unwrap_or("?"), err) } } info!("Saving file to {}", &path.to_str().unwrap_or("?")); - File::create(&path) - .and_then(|f| serde_yaml::to_writer(f, &config).map_err(to_io_error)) + let serialized = serde_yaml::to_string(config) + .map_err(|err| create_tuliprox_error!(TuliproxErrorKind::Info, "Could not serialize file {}: {}", &path.to_str().unwrap_or("?"), err))?; + + fs::write(&path, serialized) + .await .map_err(|err| create_tuliprox_error!(TuliproxErrorKind::Info, "Could not write file {}: {}", &path.to_str().unwrap_or("?"), err)) } -pub fn save_api_proxy(file_path: &str, backup_dir: &str, config: &ApiProxyConfigDto) -> Result<(), TuliproxError> { - write_config_file(file_path, backup_dir, config, "api-proxy.yml") +pub async fn save_api_proxy(file_path: &str, backup_dir: &str, config: &ApiProxyConfigDto) -> Result<(), TuliproxError> { + write_config_file(file_path, backup_dir, config, "api-proxy.yml").await } -pub fn save_main_config(file_path: &str, backup_dir: &str, config: &ConfigDto) -> Result<(), TuliproxError> { - write_config_file(file_path, backup_dir, config, "config.yml") +pub async fn save_main_config(file_path: &str, backup_dir: &str, config: &ConfigDto) -> Result<(), TuliproxError> { + write_config_file(file_path, backup_dir, config, "config.yml").await } pub fn resolve_env_var(value: &str) -> String { @@ -351,4 +355,4 @@ mod tests { let resolved = resolve_env_var("${env:HOME}"); assert_eq!(resolved, std::env::var("HOME").unwrap()); } -} \ No newline at end of file +} From 34a6246cbc969d9451a99f6815782d4a2d651c63 Mon Sep 17 00:00:00 2001 From: DarkBreakpoint <243206744+DarkBreakpoint@users.noreply.github.com> Date: Wed, 19 Nov 2025 15:35:34 -0600 Subject: [PATCH 6/6] async API user --- CHANGELOG.md | 1 + backend/src/repository/user_repository.rs | 42 ++++++++++++++++------- 2 files changed, 30 insertions(+), 13 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index f90eb7e37..5b50a3233 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,7 @@ - Added detailed shared-stream/buffer/provider logging to trace lag, cache persistence, and session/provider lifecycle events. - Connection registration failures now trigger an explicit disconnect so zombie sockets don’t linger. - Video download queue now uses async file I/O to keep the runtime responsive during large transfers. +- API user DB persistence (merge/backup/store) now executes through async Tokio I/O so user-management APIs stay responsive without extra blocking hops. - JSON playlist/category writers (xtream collections, user bouquets) now stream through Tokio I/O so persisting those files no longer blocks the runtime. - Playlist EPG exports now write via async file handles to avoid blocking the runtime during XML serialization. - Config and API proxy save endpoints now serialize via Tokio I/O, so editing configs through the API no longer blocks the runtime threads. diff --git a/backend/src/repository/user_repository.rs b/backend/src/repository/user_repository.rs index 645ec42da..e167a4b36 100644 --- a/backend/src/repository/user_repository.rs +++ b/backend/src/repository/user_repository.rs @@ -9,9 +9,10 @@ use crate::utils::json_write_documents_to_file; use chrono::Local; use log::error; use std::collections::{HashMap, HashSet}; -use std::io::Error; +use std::io::{Error, ErrorKind}; use std::path::{Path, PathBuf}; use crate::utils; +use tokio::task; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] struct StoredProxyUserCredentialsDeprecated { @@ -122,12 +123,24 @@ fn add_target_user_to_user_tree(target_users: &[TargetUser], user_tree: &mut BPl pub async fn merge_api_user(cfg: &AppConfig, target_users: &[TargetUser]) -> Result { let path = get_api_user_db_path(cfg); - let lock = cfg.file_locks.read_lock(&path).await; - let mut user_tree: BPlusTree = BPlusTree::load(&path).unwrap_or_else(|_| BPlusTree::new()); - drop(lock); + let read_lock = cfg.file_locks.read_lock(&path).await; + let mut user_tree: BPlusTree = task::spawn_blocking({ + let path = path.clone(); + move || BPlusTree::load(&path).unwrap_or_else(|_| BPlusTree::new()) + }) + .await + .map_err(|err| Error::new(ErrorKind::Other, format!("Failed to load user db: {err}")))?; + drop(read_lock); add_target_user_to_user_tree(target_users, &mut user_tree); - let _lock = cfg.file_locks.write_lock(&path).await; - user_tree.store(&path) + let write_lock = cfg.file_locks.write_lock(&path).await; + let result = task::spawn_blocking({ + let path = path.clone(); + move || user_tree.store(&path) + }) + .await + .map_err(|err| Error::new(ErrorKind::Other, format!("Failed to store user db: {err}")))?; + drop(write_lock); + result } /// # Panics @@ -136,10 +149,11 @@ pub async fn merge_api_user(cfg: &AppConfig, target_users: &[TargetUser]) -> Res pub async fn backup_api_user_db_file(cfg: &AppConfig, path: &Path) { if let Some(backup_dir) = cfg.config.load().backup_dir.as_ref() { let backup_path = PathBuf::from(backup_dir).join(format!("{}_{}", storage_const::API_USER_DB_FILE, Local::now().format("%Y%m%d_%H%M%S"))); - let _lock = cfg.file_locks.read_lock(path).await; - match std::fs::copy(path, &backup_path) { - Ok(_) => {} - Err(err) => { error!("Could not backup file {}:{}", &backup_path.to_str().unwrap_or("?"), err) } + let lock = cfg.file_locks.read_lock(path).await; + let copy_result = tokio::fs::copy(path, &backup_path).await; + drop(lock); + if let Err(err) = copy_result { + error!("Could not backup file {}:{}", &backup_path.to_str().unwrap_or("?"), err); } } } @@ -149,8 +163,10 @@ pub async fn store_api_user(cfg: &AppConfig, target_users: &[TargetUser]) -> Res add_target_user_to_user_tree(target_users, &mut user_tree); let path = get_api_user_db_path(cfg); backup_api_user_db_file(cfg, &path).await; - let _lock = cfg.file_locks.write_lock(&path).await; - user_tree.store(&path) + let write_lock = cfg.file_locks.write_lock(&path).await; + let result = user_tree.store(&path); + drop(write_lock); + result } // TODO remove me if we get stable on user_db @@ -492,4 +508,4 @@ mod tests { assert_eq!(user_list.as_ref().unwrap().len(), 1); assert_eq!(user_list.as_ref().unwrap().first().unwrap().credentials.len(), 4); } -} \ No newline at end of file +}