Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
67 changes: 60 additions & 7 deletions src/pipeline/physics_pipeline/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<T> {
pending: std::sync::Arc<std::sync::Mutex<Option<T>>>,
receiver: std::sync::mpsc::Receiver<T>,
sender: std::sync::mpsc::Sender<T>,
run: fn(&mut T),
}

#[cfg(feature = "parallel")]
impl<T: Send + 'static> DeferredBvhOptimizeJob<T> {
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
Expand Down Expand Up @@ -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<Option<std::sync::mpsc::Receiver<crate::geometry::DeferredBvhOptimize>>>,
std::sync::Mutex<Option<DeferredBvhOptimizeJob<crate::geometry::DeferredBvhOptimize>>>,
/// 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.
Expand Down Expand Up @@ -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;
}
Expand Down
16 changes: 7 additions & 9 deletions src/pipeline/physics_pipeline/solve.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
Expand Down
53 changes: 53 additions & 0 deletions src/pipeline/physics_pipeline/test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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");
}
}
}
Loading