Create worker module to contain all the feature flag chaos

This commit is contained in:
Héctor Ramón Jiménez 2025-10-26 21:48:10 +01:00
parent b408961d77
commit 22488c537c
No known key found for this signature in database
GPG key ID: 7CC46565708259A7
3 changed files with 220 additions and 201 deletions

View file

@ -2,11 +2,11 @@ use crate::core::{self, Size};
use crate::graphics::Shell; use crate::graphics::Shell;
use crate::image::atlas::{self, Atlas}; use crate::image::atlas::{self, Atlas};
#[cfg(feature = "image")] #[cfg(all(feature = "image", not(target_arch = "wasm32")))]
use std::collections::HashMap; use worker::Worker;
#[cfg(feature = "image")] #[cfg(feature = "image")]
use std::sync::mpsc; use std::collections::HashMap;
use std::sync::Arc; use std::sync::Arc;
@ -16,14 +16,8 @@ pub struct Cache {
raster: Raster, raster: Raster,
#[cfg(feature = "svg")] #[cfg(feature = "svg")]
vector: crate::image::vector::Cache, vector: crate::image::vector::Cache,
#[cfg(feature = "image")]
jobs: mpsc::SyncSender<Job>,
#[cfg(feature = "image")]
quit: mpsc::SyncSender<()>,
#[cfg(feature = "image")]
work: mpsc::Receiver<Work>,
#[cfg(all(feature = "image", not(target_arch = "wasm32")))] #[cfg(all(feature = "image", not(target_arch = "wasm32")))]
worker_: Option<std::thread::JoinHandle<()>>, worker: Worker,
} }
impl Cache { impl Cache {
@ -35,33 +29,20 @@ impl Cache {
_shell: &Shell, _shell: &Shell,
) -> Self { ) -> Self {
#[cfg(all(feature = "image", not(target_arch = "wasm32")))] #[cfg(all(feature = "image", not(target_arch = "wasm32")))]
let (worker, jobs, quit, work) = let worker =
Worker::new(device, _queue, backend, layout.clone(), _shell); Worker::new(device, _queue, backend, layout.clone(), _shell);
#[cfg(all(feature = "image", target_arch = "wasm32"))]
let (jobs, work) = (mpsc::sync_channel(0).0, mpsc::sync_channel(0).1);
#[cfg(all(feature = "image", not(target_arch = "wasm32")))]
let handle = std::thread::spawn(move || worker.run());
Self { Self {
atlas: Atlas::new(device, backend, layout), atlas: Atlas::new(device, backend, layout),
#[cfg(feature = "image")] #[cfg(feature = "image")]
raster: Raster { raster: Raster {
cache: crate::image::raster::Cache::default(), cache: crate::image::raster::Cache::default(),
pending: HashMap::new(), pending: HashMap::new(),
jobs: jobs.clone(),
}, },
#[cfg(feature = "svg")] #[cfg(feature = "svg")]
vector: crate::image::vector::Cache::default(), vector: crate::image::vector::Cache::default(),
#[cfg(feature = "image")]
jobs,
#[cfg(feature = "image")]
quit,
#[cfg(feature = "image")]
work,
#[cfg(all(feature = "image", not(target_arch = "wasm32")))] #[cfg(all(feature = "image", not(target_arch = "wasm32")))]
worker_: Some(handle), worker,
} }
} }
@ -100,7 +81,7 @@ impl Cache {
} }
let _ = self.raster.pending.insert(handle.id(), vec![callback]); let _ = self.raster.pending.insert(handle.id(), vec![callback]);
let _ = self.raster.jobs.send(Job::Load(handle.clone())); self.worker.load(handle);
} }
#[cfg(feature = "image")] #[cfg(feature = "image")]
@ -110,7 +91,7 @@ impl Cache {
if let Some(memory) = load_image( if let Some(memory) = load_image(
&mut self.raster.cache, &mut self.raster.cache,
&mut self.raster.pending, &mut self.raster.pending,
&mut self.raster.jobs, &self.worker,
handle, handle,
None, None,
) { ) {
@ -141,7 +122,7 @@ impl Cache {
let memory = load_image( let memory = load_image(
&mut self.raster.cache, &mut self.raster.cache,
&mut self.raster.pending, &mut self.raster.pending,
&mut self.raster.jobs, &self.worker,
handle, handle,
None, None,
)?; )?;
@ -183,14 +164,8 @@ impl Cache {
} }
if !self.raster.pending.contains_key(&handle.id()) { if !self.raster.pending.contains_key(&handle.id()) {
let _ = self.jobs.send(Job::Upload {
handle: handle.clone(),
rgba: image.clone().into_raw(),
width: image.width(),
height: image.height(),
});
let _ = self.raster.pending.insert(handle.id(), Vec::new()); let _ = self.raster.pending.insert(handle.id(), Vec::new());
self.worker.upload(handle, image);
} }
None None
@ -225,7 +200,7 @@ impl Cache {
pub fn trim(&mut self) { pub fn trim(&mut self) {
#[cfg(feature = "image")] #[cfg(feature = "image")]
self.raster.cache.trim(&mut self.atlas, |bind_group| { self.raster.cache.trim(&mut self.atlas, |bind_group| {
let _ = self.jobs.send(Job::Drop(bind_group)); self.worker.drop(bind_group);
}); });
#[cfg(feature = "svg")] #[cfg(feature = "svg")]
@ -236,9 +211,9 @@ impl Cache {
fn receive(&mut self) { fn receive(&mut self) {
use crate::image::raster::Memory; use crate::image::raster::Memory;
while let Ok(work) = self.work.try_recv() { while let Ok(work) = self.worker.try_recv() {
match work { match work {
Work::Upload { worker::Work::Upload {
handle, handle,
entry, entry,
bind_group, bind_group,
@ -270,7 +245,7 @@ impl Cache {
}, },
); );
} }
Work::Error { handle, error } => { worker::Work::Error { handle, error } => {
self.raster.cache.insert(&handle, Memory::error(error)); self.raster.cache.insert(&handle, Memory::error(error));
} }
} }
@ -281,9 +256,7 @@ impl Cache {
#[cfg(all(feature = "image", not(target_arch = "wasm32")))] #[cfg(all(feature = "image", not(target_arch = "wasm32")))]
impl Drop for Cache { impl Drop for Cache {
fn drop(&mut self) { fn drop(&mut self) {
let _ = self.quit.try_send(()); self.worker.quit();
let _ = self.jobs.send(Job::Quit);
let _ = self.worker_.take().unwrap().join();
} }
} }
@ -291,7 +264,6 @@ impl Drop for Cache {
struct Raster { struct Raster {
cache: crate::image::raster::Cache, cache: crate::image::raster::Cache,
pending: HashMap<core::image::Id, Vec<Callback>>, pending: HashMap<core::image::Id, Vec<Callback>>,
jobs: mpsc::SyncSender<Job>,
} }
#[cfg(feature = "image")] #[cfg(feature = "image")]
@ -301,7 +273,7 @@ type Callback = Box<dyn FnOnce(core::image::Allocation) + Send>;
fn load_image<'a>( fn load_image<'a>(
cache: &'a mut crate::image::raster::Cache, cache: &'a mut crate::image::raster::Cache,
pending: &mut HashMap<core::image::Id, Vec<Callback>>, pending: &mut HashMap<core::image::Id, Vec<Callback>>,
jobs: &mut mpsc::SyncSender<Job>, worker: &Worker,
handle: &core::image::Handle, handle: &core::image::Handle,
callback: Option<Callback>, callback: Option<Callback>,
) -> Option<&'a mut crate::image::raster::Memory> { ) -> Option<&'a mut crate::image::raster::Memory> {
@ -315,74 +287,45 @@ fn load_image<'a>(
// Load RGBA handles synchronously, since it's very cheap // Load RGBA handles synchronously, since it's very cheap
cache.insert(handle, Memory::load(handle)); cache.insert(handle, Memory::load(handle));
} else if !pending.contains_key(&handle.id()) { } else if !pending.contains_key(&handle.id()) {
let _ = jobs.send(Job::Load(handle.clone()));
let _ = pending.insert(handle.id(), Vec::from_iter(callback)); let _ = pending.insert(handle.id(), Vec::from_iter(callback));
worker.load(handle);
} }
} }
cache.get_mut(handle) cache.get_mut(handle)
} }
#[cfg(feature = "image")]
#[derive(Debug)]
enum Job {
Load(core::image::Handle),
Upload {
handle: core::image::Handle,
rgba: core::image::Bytes,
width: u32,
height: u32,
},
Drop(Arc<wgpu::BindGroup>),
Quit,
}
#[cfg(feature = "image")]
enum Work {
Upload {
handle: core::image::Handle,
entry: atlas::Entry,
bind_group: Arc<wgpu::BindGroup>,
},
Error {
handle: core::image::Handle,
error: crate::graphics::image::image_rs::error::ImageError,
},
}
#[cfg(feature = "image")]
struct Worker {
device: wgpu::Device,
queue: wgpu::Queue,
backend: wgpu::Backend,
texture_layout: wgpu::BindGroupLayout,
shell: Shell,
belt: wgpu::util::StagingBelt,
jobs: mpsc::Receiver<Job>,
output: mpsc::SyncSender<Work>,
quit: mpsc::Receiver<()>,
}
#[cfg(all(feature = "image", not(target_arch = "wasm32")))] #[cfg(all(feature = "image", not(target_arch = "wasm32")))]
impl Worker { mod worker {
fn new( use crate::core::image;
device: &wgpu::Device, use crate::graphics::Shell;
queue: &wgpu::Queue, use crate::image::atlas::{self, Atlas};
backend: wgpu::Backend, use crate::image::raster;
texture_layout: wgpu::BindGroupLayout,
shell: &Shell,
) -> (
Self,
mpsc::SyncSender<Job>,
mpsc::SyncSender<()>,
mpsc::Receiver<Work>,
) {
let (jobs_sender, jobs_receiver) = mpsc::sync_channel(1_000);
let (quit_sender, quit_receiver) = mpsc::sync_channel(1);
let (work_sender, work_receiver) = mpsc::sync_channel(1_000);
( use std::sync::Arc;
Self { use std::sync::mpsc;
use std::thread;
pub struct Worker {
jobs: mpsc::SyncSender<Job>,
quit: mpsc::SyncSender<()>,
work: mpsc::Receiver<Work>,
handle: Option<std::thread::JoinHandle<()>>,
}
impl Worker {
pub fn new(
device: &wgpu::Device,
queue: &wgpu::Queue,
backend: wgpu::Backend,
texture_layout: wgpu::BindGroupLayout,
shell: &Shell,
) -> Self {
let (jobs_sender, jobs_receiver) = mpsc::sync_channel(1_000);
let (quit_sender, quit_receiver) = mpsc::sync_channel(1);
let (work_sender, work_receiver) = mpsc::sync_channel(1_000);
let instance = Instance {
device: device.clone(), device: device.clone(),
queue: queue.clone(), queue: queue.clone(),
backend, backend,
@ -392,114 +335,190 @@ impl Worker {
jobs: jobs_receiver, jobs: jobs_receiver,
output: work_sender, output: work_sender,
quit: quit_receiver, quit: quit_receiver,
},
jobs_sender,
quit_sender,
work_receiver,
)
}
fn run(mut self) {
loop {
if self.quit.try_recv().is_ok() {
return;
}
let Ok(job) = self.jobs.recv() else {
return;
}; };
match job { let handle = thread::spawn(move || instance.run());
Job::Load(handle) => {
match crate::graphics::image::load(&handle) { Self {
Ok(image) => self.upload( jobs: jobs_sender,
handle, quit: quit_sender,
image.width(), work: work_receiver,
image.height(), handle: Some(handle),
image.into_raw(),
Shell::invalidate_layout,
),
Err(error) => {
let _ =
self.output.send(Work::Error { handle, error });
}
}
}
Job::Upload {
handle,
rgba,
width,
height,
} => {
self.upload(
handle,
width,
height,
rgba,
Shell::request_redraw,
);
}
Job::Drop(bind_group) => {
drop(bind_group);
}
Job::Quit => return,
} }
} }
pub fn load(&self, handle: &image::Handle) {
let _ = self.jobs.send(Job::Load(handle.clone()));
}
pub fn upload(&self, handle: &image::Handle, image: raster::Image) {
let _ = self.jobs.send(Job::Upload {
handle: handle.clone(),
width: image.width(),
height: image.height(),
rgba: image.into_raw(),
});
}
pub fn drop(&self, bind_group: Arc<wgpu::BindGroup>) {
let _ = self.jobs.send(Job::Drop(bind_group));
}
pub fn try_recv(&self) -> Result<Work, mpsc::TryRecvError> {
self.work.try_recv()
}
pub fn quit(&mut self) {
let _ = self.quit.try_send(());
let _ = self.jobs.send(Job::Quit);
let _ = self.handle.take().map(thread::JoinHandle::join);
}
} }
fn upload( pub struct Instance {
&mut self, device: wgpu::Device,
handle: core::image::Handle, queue: wgpu::Queue,
width: u32, backend: wgpu::Backend,
height: u32, texture_layout: wgpu::BindGroupLayout,
rgba: core::image::Bytes, shell: Shell,
callback: fn(&Shell), belt: wgpu::util::StagingBelt,
) { jobs: mpsc::Receiver<Job>,
let mut encoder = self.device.create_command_encoder( output: mpsc::SyncSender<Work>,
&wgpu::CommandEncoderDescriptor { quit: mpsc::Receiver<()>,
label: Some("raster image upload"), }
},
);
let mut atlas = Atlas::with_size( #[cfg(feature = "image")]
&self.device, #[derive(Debug)]
self.backend, enum Job {
self.texture_layout.clone(), Load(image::Handle),
width.max(height), Upload {
); handle: image::Handle,
rgba: image::Bytes,
width: u32,
height: u32,
},
Drop(Arc<wgpu::BindGroup>),
Quit,
}
let Some(entry) = atlas.upload( #[cfg(feature = "image")]
&self.device, pub enum Work {
&mut encoder, Upload {
&mut self.belt, handle: image::Handle,
width, entry: atlas::Entry,
height, bind_group: Arc<wgpu::BindGroup>,
&rgba, },
) else { Error {
return; handle: image::Handle,
}; error: crate::graphics::image::image_rs::error::ImageError,
},
}
let output = self.output.clone(); #[cfg(all(feature = "image", not(target_arch = "wasm32")))]
let shell = self.shell.clone(); impl Instance {
fn run(mut self) {
loop {
if self.quit.try_recv().is_ok() {
return;
}
self.belt.finish(); let Ok(job) = self.jobs.recv() else {
let submission = self.queue.submit([encoder.finish()]); return;
self.belt.recall(); };
let bind_group = atlas.bind_group().clone(); match job {
Job::Load(handle) => {
match crate::graphics::image::load(&handle) {
Ok(image) => self.upload(
handle,
image.width(),
image.height(),
image.into_raw(),
Shell::invalidate_layout,
),
Err(error) => {
let _ = self
.output
.send(Work::Error { handle, error });
}
}
}
Job::Upload {
handle,
rgba,
width,
height,
} => {
self.upload(
handle,
width,
height,
rgba,
Shell::request_redraw,
);
}
Job::Drop(bind_group) => {
drop(bind_group);
}
Job::Quit => return,
}
}
}
self.queue.on_submitted_work_done(move || { fn upload(
let _ = output.send(Work::Upload { &mut self,
handle, handle: image::Handle,
entry, width: u32,
bind_group, height: u32,
rgba: image::Bytes,
callback: fn(&Shell),
) {
let mut encoder = self.device.create_command_encoder(
&wgpu::CommandEncoderDescriptor {
label: Some("raster image upload"),
},
);
let mut atlas = Atlas::with_size(
&self.device,
self.backend,
self.texture_layout.clone(),
width.max(height),
);
let Some(entry) = atlas.upload(
&self.device,
&mut encoder,
&mut self.belt,
width,
height,
&rgba,
) else {
return;
};
let output = self.output.clone();
let shell = self.shell.clone();
self.belt.finish();
let submission = self.queue.submit([encoder.finish()]);
self.belt.recall();
let bind_group = atlas.bind_group().clone();
self.queue.on_submitted_work_done(move || {
let _ = output.send(Work::Upload {
handle,
entry,
bind_group,
});
callback(&shell);
}); });
callback(&shell); let _ = self
}); .device
.poll(wgpu::PollType::WaitForSubmissionIndex(submission));
let _ = self }
.device
.poll(wgpu::PollType::WaitForSubmissionIndex(submission));
} }
} }

View file

@ -7,7 +7,7 @@ use crate::image::atlas::{self, Atlas};
use rustc_hash::{FxHashMap, FxHashSet}; use rustc_hash::{FxHashMap, FxHashSet};
use std::sync::{Arc, Weak}; use std::sync::{Arc, Weak};
type Image = image_rs::ImageBuffer<image_rs::Rgba<u8>, image::Bytes>; pub type Image = image_rs::ImageBuffer<image_rs::Rgba<u8>, image::Bytes>;
/// Entry in cache corresponding to an image handle /// Entry in cache corresponding to an image handle
#[derive(Debug)] #[derive(Debug)]

View file

@ -82,9 +82,9 @@ fn fs_main(input: VertexOutput) -> @location(0) vec4<f32> {
2.0 * (fragment - position - scale / 2.0), 2.0 * (fragment - position - scale / 2.0),
scale, scale,
input.border_radius * 2.0, input.border_radius * 2.0,
); ) / 2.0;
let antialias: f32 = clamp(0.5 - d, 0.0, 1.0); let antialias: f32 = clamp(1.0 - d, 0.0, 1.0);
return textureSample(u_texture, u_sampler, input.uv, i32(input.layer)) * vec4<f32>(1.0, 1.0, 1.0, antialias * input.opacity); return textureSample(u_texture, u_sampler, input.uv, i32(input.layer)) * vec4<f32>(1.0, 1.0, 1.0, antialias * input.opacity);
} }