#[cfg(not(target_arch = "wasm32"))]
use std::{
collections::VecDeque,
sync::{
Arc, Condvar, Mutex, MutexGuard, OnceLock, PoisonError,
atomic::{AtomicBool, AtomicU64, Ordering},
mpsc::{self, Receiver, Sender},
},
};
#[cfg(not(target_arch = "wasm32"))]
#[derive(Default)]
pub(crate) struct Landing {
queued: AtomicU64,
built: AtomicU64,
awaited: AtomicBool,
landed: AtomicBool,
wake: OnceLock<Box<dyn Fn() + Send + Sync>>,
}
#[cfg(not(target_arch = "wasm32"))]
impl Landing {
pub(crate) fn built(&self) -> u64 {
self.built.load(Ordering::Acquire)
}
pub(crate) fn in_flight(&self) -> bool {
self.queued.load(Ordering::Acquire) != self.built()
}
pub(crate) fn await_since(&self, since: u64) {
self.awaited.store(true, Ordering::Release);
if self.built() != since && self.awaited.swap(false, Ordering::AcqRel) {
self.landed.store(true, Ordering::Release);
}
}
pub(crate) fn take_landed(&self) -> bool {
self.landed.swap(false, Ordering::AcqRel)
}
pub(crate) fn wake_with(&self, wake: Box<dyn Fn() + Send + Sync>) {
let _ = self.wake.set(wake);
}
pub(crate) fn landed(&self) -> bool {
self.landed.load(Ordering::Acquire)
}
pub(crate) fn has_wake(&self) -> bool {
self.wake.get().is_some()
}
fn after_build(&self) {
self.built.fetch_add(1, Ordering::AcqRel);
if self.awaited.swap(false, Ordering::AcqRel) {
self.landed.store(true, Ordering::Release);
if let Some(wake) = self.wake.get() {
wake();
}
}
}
}
#[cfg(not(target_arch = "wasm32"))]
const LANE_FINISH_WAIT: web_time::Duration = web_time::Duration::from_secs(5);
#[cfg(not(target_arch = "wasm32"))]
fn warm_up_threads() -> usize {
std::thread::available_parallelism().map_or(1, |cores| (cores.get() / 2).clamp(1, 3))
}
#[cfg(not(target_arch = "wasm32"))]
type Job = Box<dyn FnOnce() + Send + 'static>;
#[cfg(not(target_arch = "wasm32"))]
pub(crate) trait CompilerSend: Send {}
#[cfg(not(target_arch = "wasm32"))]
impl<T: Send + ?Sized> CompilerSend for T {}
#[cfg(target_arch = "wasm32")]
pub(crate) trait CompilerSend {}
#[cfg(target_arch = "wasm32")]
impl<T: ?Sized> CompilerSend for T {}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) trait CompilerSync: Sync {}
#[cfg(not(target_arch = "wasm32"))]
impl<T: Sync + ?Sized> CompilerSync for T {}
#[cfg(target_arch = "wasm32")]
pub(crate) trait CompilerSync {}
#[cfg(target_arch = "wasm32")]
impl<T: ?Sized> CompilerSync for T {}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) enum CompileLane {
Demanded,
WarmUp,
}
#[derive(Clone, Default)]
pub(crate) struct PipelineCompiler {
#[cfg(not(target_arch = "wasm32"))]
workers: Option<Arc<Workers>>,
}
#[cfg(not(target_arch = "wasm32"))]
#[derive(Default)]
struct Queues {
demanded: VecDeque<Job>,
warm_up: VecDeque<Job>,
closed: bool,
}
#[cfg(not(target_arch = "wasm32"))]
impl Queues {
fn next(&mut self, warm_ups: bool, spread: bool) -> Option<Job> {
let demanded = if warm_ups && !spread {
None
} else {
self.demanded.pop_front()
};
demanded.or_else(|| warm_ups.then(|| self.warm_up.pop_front()).flatten())
}
}
#[cfg(not(target_arch = "wasm32"))]
#[derive(Default)]
struct Pool {
queues: Mutex<Queues>,
queued: Condvar,
spread: AtomicBool,
}
#[cfg(not(target_arch = "wasm32"))]
impl Pool {
fn queues(&self) -> MutexGuard<'_, Queues> {
self.queues.lock().unwrap_or_else(PoisonError::into_inner)
}
fn take(&self, warm_ups: bool) -> Option<Job> {
let mut queues = self.queues();
loop {
if queues.closed {
return None;
}
if let Some(job) = queues.next(warm_ups, self.spread.load(Ordering::Acquire)) {
return Some(job);
}
queues = self
.queued
.wait(queues)
.unwrap_or_else(PoisonError::into_inner);
}
}
}
#[cfg(not(target_arch = "wasm32"))]
struct Workers {
pool: Arc<Pool>,
finished: Mutex<Receiver<()>>,
threads: usize,
landing: Arc<Landing>,
}
#[cfg(not(target_arch = "wasm32"))]
impl Drop for Workers {
fn drop(&mut self) {
let skipped = {
let mut queues = self.pool.queues();
queues.closed = true;
(
std::mem::take(&mut queues.demanded),
std::mem::take(&mut queues.warm_up),
)
};
drop(skipped);
self.pool.queued.notify_all();
let finished = self
.finished
.get_mut()
.unwrap_or_else(PoisonError::into_inner);
let deadline = web_time::Instant::now() + LANE_FINISH_WAIT;
for _ in 0..self.threads {
let left = deadline.saturating_duration_since(web_time::Instant::now());
if finished.recv_timeout(left).is_err() {
log::warn!(
"[gpu-pipeline] a compile outlived its renderer by {LANE_FINISH_WAIT:?}"
);
break;
}
}
}
}
#[cfg(not(target_arch = "wasm32"))]
fn spawn_thread(
name: &str,
warm_ups: bool,
pool: &Arc<Pool>,
finished: &Sender<()>,
landing: &Arc<Landing>,
) -> std::io::Result<()> {
let pool = Arc::clone(pool);
let finished = finished.clone();
let landing = Arc::clone(landing);
std::thread::Builder::new()
.name(name.into())
.spawn(move || {
crate::render::mark_thread_off_frame();
while let Some(job) = pool.take(warm_ups) {
job();
landing.after_build();
}
let _ = finished.send(());
})
.map(drop)
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) enum PipelineCompilation {
Background,
Inline,
}
impl PipelineCompiler {
pub(crate) fn for_compilation(compilation: PipelineCompilation) -> Self {
match compilation {
PipelineCompilation::Background => Self::spawn(),
PipelineCompilation::Inline => Self::inactive(),
}
}
pub(crate) fn inactive() -> Self {
Self::default()
}
pub(crate) fn spawn() -> Self {
#[cfg(not(target_arch = "wasm32"))]
{
static BACKGROUND_PIPELINES: crate::debug_toggles::DebugToggle =
crate::debug_toggles::DebugToggle::new("CRANPOSE_BACKGROUND_PIPELINES");
if BACKGROUND_PIPELINES.equals("0") {
return Self::inactive();
}
let (finished_tx, finished) = mpsc::channel();
let mut workers = Workers {
pool: Arc::default(),
finished: Mutex::new(finished),
threads: 0,
landing: Arc::default(),
};
let spawned = std::iter::once(("cranpose-pipelines", false))
.chain(std::iter::repeat_n(
("cranpose-warm-up", true),
warm_up_threads(),
))
.try_for_each(|(name, warm_ups)| -> std::io::Result<()> {
spawn_thread(
name,
warm_ups,
&workers.pool,
&finished_tx,
&workers.landing,
)?;
workers.threads += 1;
Ok(())
});
match spawned {
Ok(()) => Self {
workers: Some(Arc::new(workers)),
},
Err(error) => {
log::error!("[gpu-pipeline] background compiler failed to spawn: {error}");
Self::inactive()
}
}
}
#[cfg(target_arch = "wasm32")]
{
Self::inactive()
}
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) fn landing(&self) -> Option<Arc<Landing>> {
self.workers
.as_ref()
.map(|workers| Arc::clone(&workers.landing))
}
pub(crate) fn is_active(&self) -> bool {
#[cfg(not(target_arch = "wasm32"))]
{
self.workers.is_some()
}
#[cfg(target_arch = "wasm32")]
{
false
}
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) fn spread_demand(&self, spread: bool) {
if let Some(workers) = self.workers.as_ref()
&& workers.pool.spread.swap(spread, Ordering::AcqRel) != spread
&& spread
{
drop(workers.pool.queues());
workers.pool.queued.notify_all();
}
}
pub(crate) fn enqueue(&self, lane: CompileLane, job: impl FnOnce() + CompilerSend + 'static) {
#[cfg(not(target_arch = "wasm32"))]
if let Some(workers) = self.workers.as_ref() {
workers.landing.queued.fetch_add(1, Ordering::AcqRel);
let mut queues = workers.pool.queues();
match lane {
CompileLane::Demanded => queues.demanded.push_back(Box::new(job)),
CompileLane::WarmUp => queues.warm_up.push_back(Box::new(job)),
}
drop(queues);
workers.pool.queued.notify_all();
}
#[cfg(target_arch = "wasm32")]
drop((lane, job));
}
}
#[cfg(all(test, not(target_arch = "wasm32")))]
#[path = "tests/pipeline_compiler_tests.rs"]
pub(crate) mod tests;