From 0d9cc1b9ad274029c9589c5d0c6c2ac0cffe5ecc Mon Sep 17 00:00:00 2001 From: HampusM Date: Sat, 10 Oct 2026 14:47:41 +0200 Subject: perf(engine): use thread pool for asset importing --- engine/src/asset.rs | 115 ++++++++++++++++++++++++++++------------------------ 1 file changed, 63 insertions(+), 52 deletions(-) (limited to 'engine/src/asset.rs') diff --git a/engine/src/asset.rs b/engine/src/asset.rs index 8f7829c..66efe35 100644 --- a/engine/src/asset.rs +++ b/engine/src/asset.rs @@ -4,22 +4,24 @@ use std::cell::RefCell; use std::convert::Infallible; use std::fmt::{Debug, Display}; use std::hash::Hash; -use std::hint::cold_path; use std::marker::PhantomData; +use std::panic::{catch_unwind, RefUnwindSafe, UnwindSafe}; use std::path::{Path, PathBuf}; use std::sync::mpsc::{ channel as mpsc_channel, Receiver as MpscReceiver, Sender as MpscSender, }; +use std::thread::available_parallelism; -use ecs::actions::Actions; +use rayon::{ThreadPool, ThreadPoolBuilder}; use crate::ecs::pair::ChildOf; use crate::ecs::phase::{Phase, PRE_UPDATE as PRE_UPDATE_PHASE}; use crate::ecs::sole::Single; use crate::ecs::{declare_entity, pair, Sole}; -use crate::work_queue::{Work, WorkQueue}; + +const BACKUP_IMPORT_WORK_THREAD_CNT: usize = 2; declare_entity! { pub HANDLE_ASSETS_PHASE: (Phase, pair!(ChildOf, { *PRE_UPDATE_PHASE })); @@ -106,7 +108,7 @@ pub struct Assets store: store::Store, metadata_lut: RefCell>, importers: hashbrown::HashMap, - import_work_queue: WorkQueue, + import_work_thread_pool: ThreadPool, import_work_msg_receiver: MpscReceiver, import_work_msg_sender: MpscSender, curr_tick_events: Vec, @@ -115,16 +117,23 @@ pub struct Assets impl Assets { #[must_use] - pub fn with_capacity(capacity: usize) -> Self + #[tracing::instrument(skip_all)] + pub(crate) fn with_capacity(capacity: usize) -> Self { let (import_work_msg_sender, import_work_msg_receiver) = mpsc_channel::(); + let import_work_thread_cnt = calc_import_work_thread_cnt(); + Self { store: store::Store::with_capacity(capacity), metadata_lut: RefCell::new(hashbrown::HashMap::with_capacity(capacity)), importers: hashbrown::HashMap::new(), - import_work_queue: WorkQueue::new("asset_importing_work_queue"), + import_work_thread_pool: ThreadPoolBuilder::new() + .thread_name(|thread_index| format!("asset_import_worker_{thread_index}")) + .num_threads(import_work_thread_cnt) + .build() + .expect("Unable to create thread pool"), import_work_msg_receiver, import_work_msg_sender, curr_tick_events: Vec::with_capacity(capacity), @@ -136,7 +145,7 @@ impl Assets func: impl Fn(&mut Submitter<'_>, &Path, Option<&AssetSettings>) -> Result<(), Err>, ) where AssetT: Asset, - AssetSettings: 'static, + AssetSettings: UnwindSafe + RefUnwindSafe + 'static, Err: std::error::Error + Send + Sync + 'static, { self.importers @@ -284,7 +293,7 @@ impl Assets ) -> Handle where AssetT: Asset, - AssetSettings: Send + Sync + Debug + 'static, + AssetSettings: UnwindSafe + RefUnwindSafe + Send + Sync + Debug + 'static, { let label = label.into(); @@ -517,7 +526,7 @@ impl Assets asset_settings: Option, ) -> IdValid where - AssetSettings: Send + Sync + Debug + 'static, + AssetSettings: UnwindSafe + RefUnwindSafe + Send + Sync + Debug + 'static, { let id = match self.get_asset_by_label(&label) { Some((Some(_), id)) => { @@ -565,24 +574,25 @@ impl Assets label: &Label<'_>, asset_settings: Option, ) where - AssetSettings: Any + Send + Sync, + AssetSettings: UnwindSafe + RefUnwindSafe + Any + Send + Sync, { - let Some(importer) = self.importers.get(&asset_ty_id) else { + let Some(importer) = self.importers.get(&asset_ty_id).cloned() else { tracing::error!("No importer exists for asset"); return; }; - self.import_work_queue.add_work(Work { - func: |ImportWorkUserData { - import_work_msg_sender, - asset_path, - asset_settings, - importer, - }| { + let import_work_msg_sender = self.import_work_msg_sender.clone(); + + let asset_path = label.path.to_path_buf(); + + self.import_work_thread_pool.spawn(move || { + let result = catch_unwind(|| { if let Err(err) = importer.call( import_work_msg_sender, asset_path.as_path(), - asset_settings.as_deref(), + asset_settings + .as_ref() + .map(|asset_settings| asset_settings as &(dyn Any + Send + Sync)), ) { tracing::error!( "Failed to import asset {}: {:#}", @@ -590,15 +600,22 @@ impl Assets crate::Error::new(err) ); } - }, - user_data: ImportWorkUserData { - import_work_msg_sender: self.import_work_msg_sender.clone(), - asset_path: label.path.to_path_buf(), - asset_settings: asset_settings.map(|asset_settings| { - Box::new(asset_settings) as Box - }), - importer: importer.clone(), - }, + }); + + if let Err(panic_msg) = result { + let panic_msg = if let Some(msg) = panic_msg.downcast_ref::<&str>() { + msg + } else if let Some(msg) = panic_msg.downcast_ref::() { + msg.as_str() + } else { + "(unknown panic payload type)" + }; + + tracing::error!( + "Failed to import asset {}: Thread panicked: {panic_msg}", + asset_path.display(), + ); + } }); } @@ -827,7 +844,7 @@ impl WrappedImporterFn where InnerFunc: Fn(&mut Submitter<'_>, &Path, Option<&AssetSettings>) -> Result<(), Err>, - AssetSettings: 'static, + AssetSettings: UnwindSafe + RefUnwindSafe + 'static, Err: std::error::Error + Send + Sync + 'static, { assert_eq!(size_of::(), 0); @@ -892,7 +909,6 @@ impl crate::ecs::extension::Extension for Extension collector.spawn_declared_entity(&HANDLE_ASSETS_PHASE); collector.add_system(*HANDLE_ASSETS_PHASE, add_received_assets); - collector.add_system(*HANDLE_ASSETS_PHASE, check_import_wq_thread_not_panicked); } } @@ -919,28 +935,6 @@ fn add_received_assets(mut assets: Single, mut events: Single) } } -fn check_import_wq_thread_not_panicked(assets: Single, mut actions: Actions<'_>) -{ - let Ok(assets) = assets.get() else { - unreachable!(); - }; - - if assets.import_work_queue.get_thread_panic().is_some() { - cold_path(); - - actions.stop(); - } -} - -#[derive(Debug)] -struct ImportWorkUserData -{ - import_work_msg_sender: MpscSender, - asset_path: PathBuf, - asset_settings: Option>, - importer: WrappedImporterFn, -} - #[derive(Debug)] enum ImportWorkMessage { @@ -1055,6 +1049,23 @@ impl hashbrown::Equivalent for Label<'_> } } +fn calc_import_work_thread_cnt() -> usize +{ + available_parallelism() + .map(|avail_parallelism| avail_parallelism.get().div_ceil(2)) + .inspect_err(|err| { + tracing::warn!( + concat!( + "Failed to get available parallelism for calculating import work ", + "thread count. {} threads will be used. Error: {}" + ), + BACKUP_IMPORT_WORK_THREAD_CNT, + err + ); + }) + .unwrap_or(BACKUP_IMPORT_WORK_THREAD_CNT) +} + mod store { use std::num::NonZero; -- cgit v1.2.3-18-g5258