use std::cell::RefCell;
use std::collections::VecDeque;
use std::rc::Rc;
use brep_render::brep_kernel::HistoryRequest;
use brep_render::runner::{
Command, HistoryRunner, MeasureQuery, MeasureReply, MeshImportReply, MeshImportRequest, Reply,
RunReply,
};
use wasm_bindgen::prelude::*;
use wasm_bindgen::JsCast;
use web_sys::MessageEvent;
pub struct WorkerRunner {
worker: web_sys::Worker,
inbox: Rc<RefCell<VecDeque<Reply>>>,
_onmessage: Closure<dyn FnMut(MessageEvent)>,
run_buf: VecDeque<RunReply>,
query_buf: VecDeque<MeasureReply>,
mesh_import_buf: VecDeque<MeshImportReply>,
in_flight: bool,
pending_run: Option<(HistoryRequest, u64, u64)>,
sent_library_revision: Option<u64>,
library_requested: bool,
}
impl WorkerRunner {
pub fn new() -> Self {
let options = web_sys::WorkerOptions::new();
options.set_type(web_sys::WorkerType::Module);
let worker = web_sys::Worker::new_with_options("./worker.js", &options)
.unwrap_or_else(|e| panic!("failed to spawn history worker (./worker.js): {e:?}"));
let inbox: Rc<RefCell<VecDeque<Reply>>> = Rc::new(RefCell::new(VecDeque::new()));
let inbox_cb = inbox.clone();
let onmessage = Closure::wrap(Box::new(move |event: MessageEvent| {
if let Some(text) = event.data().as_string() {
match serde_json::from_str::<Reply>(&text) {
Ok(reply) => inbox_cb.borrow_mut().push_back(reply),
Err(error) => log::error!("worker reply parse failed: {error}"),
}
}
}) as Box<dyn FnMut(MessageEvent)>);
worker.set_onmessage(Some(onmessage.as_ref().unchecked_ref()));
Self {
worker,
inbox,
_onmessage: onmessage,
run_buf: VecDeque::new(),
query_buf: VecDeque::new(),
mesh_import_buf: VecDeque::new(),
in_flight: false,
pending_run: None,
sent_library_revision: None,
library_requested: false,
}
}
fn post(&self, command: &Command) -> Result<(), JsValue> {
let json = serde_json::to_string(command).expect("serialize worker command");
self.worker.post_message(&JsValue::from_str(&json))
}
fn post_run(&mut self, request: HistoryRequest, generation: u64, library_revision: u64) {
self.post(&Command::Run {
request,
generation,
parts_library_revision: library_revision,
})
.unwrap_or_else(|e| panic!("post run to history worker failed: {e:?}"));
self.in_flight = true;
}
fn drain(&mut self) {
loop {
let next = self.inbox.borrow_mut().pop_front();
let Some(reply) = next else { break };
match reply {
Reply::Run(run) => {
self.run_buf.push_back(run);
self.in_flight = false;
if let Some((request, generation, revision)) = self.pending_run.take() {
self.post_run(request, generation, revision);
}
}
Reply::Query(query) => self.query_buf.push_back(query),
Reply::MeshImport(reply) => self.mesh_import_buf.push_back(reply),
Reply::NeedPartsLibrary => {
self.in_flight = false;
self.sent_library_revision = None;
self.library_requested = true;
self.pending_run = None;
}
}
}
}
}
impl Default for WorkerRunner {
fn default() -> Self {
Self::new()
}
}
impl HistoryRunner for WorkerRunner {
fn submit_run(&mut self, request: HistoryRequest, generation: u64) {
let revision = self.sent_library_revision.unwrap_or(0);
if self.in_flight {
self.pending_run = Some((request, generation, revision));
} else {
self.post_run(request, generation, revision);
}
}
fn sync_parts_library(
&mut self,
revision: u64,
fetch: &mut dyn FnMut() -> brep_render::brep_kernel::PartsLibraryMap,
) {
if self.sent_library_revision == Some(revision) {
return;
}
if self
.post(&Command::SetPartsLibrary {
library: fetch(),
revision,
})
.is_ok()
{
self.sent_library_revision = Some(revision);
}
}
fn poll_library_request(&mut self) -> bool {
self.drain();
std::mem::take(&mut self.library_requested)
}
fn poll_run(&mut self) -> Option<RunReply> {
self.drain();
self.run_buf.pop_front()
}
fn submit_query(&mut self, query: MeasureQuery) {
let _ = self.post(&Command::Query(query));
}
fn poll_query(&mut self) -> Option<MeasureReply> {
self.drain();
self.query_buf.pop_front()
}
fn submit_mesh_import(&mut self, request: MeshImportRequest) {
let _ = self.post(&Command::MeshImport(request));
}
fn poll_mesh_import(&mut self) -> Option<MeshImportReply> {
self.drain();
self.mesh_import_buf.pop_front()
}
fn reset(&mut self) {
let _ = self.post(&Command::Reset);
self.run_buf.clear();
self.query_buf.clear();
self.mesh_import_buf.clear();
self.pending_run = None;
self.in_flight = false;
self.sent_library_revision = None;
self.library_requested = false;
}
}
#[wasm_bindgen]
pub fn worker_entry() {
console_error_panic_hook::set_once();
let scope: web_sys::DedicatedWorkerGlobalScope = js_sys::global().unchecked_into();
let scope_reply = scope.clone();
let runner = Rc::new(RefCell::new(brep_render::pipeline::SceneRunner::new()));
let onmessage = Closure::wrap(Box::new(move |event: MessageEvent| {
let Some(text) = event.data().as_string() else {
return;
};
let command = match serde_json::from_str::<Command>(&text) {
Ok(command) => command,
Err(error) => {
log::error!("worker command parse failed: {error}");
return;
}
};
let reply = {
let mut resident = runner.borrow_mut();
brep_render::runner::process_command(&mut resident, command)
};
if let Some(reply) = reply {
let json = serde_json::to_string(&reply).expect("serialize worker reply");
let _ = scope_reply.post_message(&JsValue::from_str(&json));
}
}) as Box<dyn FnMut(MessageEvent)>);
scope.set_onmessage(Some(onmessage.as_ref().unchecked_ref()));
onmessage.forget();
}