diff --git a/src/pipeline/physics_pipeline/mod.rs b/src/pipeline/physics_pipeline/mod.rs index a6f3cfae4..9a2a8daa1 100644 --- a/src/pipeline/physics_pipeline/mod.rs +++ b/src/pipeline/physics_pipeline/mod.rs @@ -23,6 +23,59 @@ mod test; #[cfg(test)] mod test_staged; +/// A deferred optimization whose input stays available until a Rayon worker +/// starts it. The joining thread can run the still-pending work itself instead +/// of blocking on a task queued behind a parked worker. +#[cfg(feature = "parallel")] +struct DeferredBvhOptimizeJob { + pending: std::sync::Arc>>, + receiver: std::sync::mpsc::Receiver, + sender: std::sync::mpsc::Sender, + run: fn(&mut T), +} + +#[cfg(feature = "parallel")] +impl DeferredBvhOptimizeJob { + fn spawn(task: T, run: fn(&mut T)) -> Self { + let pending = std::sync::Arc::new(std::sync::Mutex::new(Some(task))); + let worker_pending = std::sync::Arc::clone(&pending); + let (sender, receiver) = std::sync::mpsc::channel(); + let worker_sender = sender.clone(); + + rayon::spawn(move || { + let Some(mut task) = worker_pending.lock().unwrap().take() else { + return; + }; + run(&mut task); + let _ = worker_sender.send(task); + }); + + Self { + pending, + receiver, + sender, + run, + } + } + + fn join(self) -> T { + // If the Rayon job is still queued, take over its work here before + // blocking on the result channel. This avoids depending on another + // worker waking up to run it. + if let Some(mut task) = self.pending.lock().unwrap().take() { + (self.run)(&mut task); + let _ = self.sender.send(task); + } + + // Drop the helper sender before waiting so a worker panic still + // disconnects the channel instead of leaving recv() blocked forever. + drop(self.sender); + self.receiver + .recv() + .expect("the deferred BVH optimization task died") + } +} + /// The main physics simulation engine that runs your physics world forward in time. /// /// Think of this as the "game loop" for your physics simulation. Each frame, you call @@ -72,12 +125,12 @@ pub struct PhysicsPipeline { /// staged workers. On a non-parallel (or wasm) build it runs with one worker /// inline on the calling thread. staged_solver: crate::dynamics::StagedIslandSolver, - /// Handle on the BVH optimization pass running concurrently with the narrow - /// phase and solver (the `Mutex` only exists to keep the pipeline `Sync`; it is - /// never contended). + /// Handle on the BVH optimization pass scheduled alongside the narrow phase + /// and solver. It retains the task until a worker claims it, so the join point + /// can run queued work itself (the `Mutex` keeps the pipeline `Sync`). #[cfg(feature = "parallel")] deferred_bvh: - std::sync::Mutex>>, + std::sync::Mutex>>, /// Deferred BVH optimization that had no spare worker to run on (single-threaded /// pool, or `parallel` off): run inline by `join_deferred_bvh_optimize`, i.e. at /// the same point of the step where the concurrent one is joined. @@ -130,13 +183,13 @@ impl PhysicsPipeline { /// any) and puts the optimized tree back into the broad-phase. Must be called before /// anything uses the broad-phase tree again. /// - /// Waits for the concurrent pass when one was spawned; otherwise runs it here. Both + /// Runs a still-queued pass here or waits for a worker that already claimed it. Both /// paths leave the same tree behind, so the build and the pool size don't change what /// the rest of the step sees. fn join_deferred_bvh_optimize(&mut self, broad_phase: &mut BroadPhaseBvh) { #[cfg(feature = "parallel")] - if let Some(rx) = self.deferred_bvh.get_mut().unwrap().take() { - let task = rx.recv().expect("the deferred BVH optimization task died"); + if let Some(job) = self.deferred_bvh.get_mut().unwrap().take() { + let task = job.join(); broad_phase.finish_deferred_optimize(task); return; } diff --git a/src/pipeline/physics_pipeline/solve.rs b/src/pipeline/physics_pipeline/solve.rs index 7ba32026d..11797e85a 100644 --- a/src/pipeline/physics_pipeline/solve.rs +++ b/src/pipeline/physics_pipeline/solve.rs @@ -85,17 +85,15 @@ impl PhysicsPipeline { // walked must stay un-optimized until the join, in every build — but the // *execution* needs a spare worker. `step` itself runs inside the pool (see // `PhysicsPipeline::step`), so on a single-worker pool a detached task would - // queue behind the `recv` waiting for it: deadlock. Hand those to the join - // point instead, which runs them inline. + // queue behind the join waiting for it: deadlock. Hand those to the join + // point instead, which runs them inline. With multiple workers, the joining + // thread can also take over if the Rayon job is still queued. #[cfg(feature = "parallel")] if rayon::current_num_threads() > 1 { - let (tx, rx) = std::sync::mpsc::channel(); - let mut task = task; - rayon::spawn(move || { - task.run(); - let _ = tx.send(task); - }); - *self.deferred_bvh.get_mut().unwrap() = Some(rx); + *self.deferred_bvh.get_mut().unwrap() = Some(super::DeferredBvhOptimizeJob::spawn( + task, + crate::geometry::DeferredBvhOptimize::run, + )); } else { self.deferred_bvh_inline = Some(task); } diff --git a/src/pipeline/physics_pipeline/test.rs b/src/pipeline/physics_pipeline/test.rs index d69a5ac34..a6d826441 100644 --- a/src/pipeline/physics_pipeline/test.rs +++ b/src/pipeline/physics_pipeline/test.rs @@ -738,3 +738,56 @@ fn contact_force_events_follow_runtime_active_events_flips() { } assert_eq!(events.0.load(Ordering::Relaxed), after_enable); } + +#[cfg(feature = "parallel")] +#[test] +fn deferred_bvh_join_runs_pending_work_when_other_worker_is_blocked() { + use std::sync::mpsc; + use std::time::Duration; + + let pool = rayon::ThreadPoolBuilder::new() + .num_threads(2) + .build() + .unwrap(); + let (blocked_tx, blocked_rx) = mpsc::channel(); + let (resume_tx, resume_rx) = mpsc::channel(); + let (done_tx, done_rx) = mpsc::channel(); + let worker_resume_tx = resume_tx.clone(); + + let worker = std::thread::spawn(move || { + let result = pool.install(move || { + // Occupy the second worker so the deferred optimization job stays + // queued. The joining worker must take and run that pending work. + rayon::spawn(move || { + blocked_tx.send(()).unwrap(); + let _ = resume_rx.recv(); + }); + blocked_rx + .recv_timeout(Duration::from_secs(2)) + .expect("the second Rayon worker did not start"); + + let job = super::DeferredBvhOptimizeJob::spawn(41usize, |value| *value += 1); + let result = job.join(); + let _ = worker_resume_tx.send(()); + result + }); + let _ = done_tx.send(result); + }); + + match done_rx.recv_timeout(Duration::from_secs(2)) { + Ok(result) => { + assert_eq!(result, 42); + worker.join().expect("deferred-work test worker panicked"); + } + Err(mpsc::RecvTimeoutError::Timeout) => { + // Release the parked worker so an incorrect blocking join can exit + // and the test reports a failure instead of wedging the harness. + let _ = resume_tx.send(()); + let _ = done_rx.recv_timeout(Duration::from_secs(2)); + panic!("joining deferred work blocked while the other worker was busy"); + } + Err(mpsc::RecvTimeoutError::Disconnected) => { + panic!("deferred-work test worker died before reporting a result"); + } + } +}