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] 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