use std::cell::{Cell, RefCell};
use std::path::{Path, PathBuf};
use std::rc::Rc;
use deno_core::{Extension, JsRuntime, OpDecl, OpState, RuntimeOptions, op2};
use deno_error::JsErrorBox;
use tokio::sync::Mutex;
use tokio::sync::Notify;
use tokio::sync::mpsc::UnboundedReceiver;
use crate::animations::AnimationCommand;
pub use crate::host::{FlushInfo, HostSenders};
use crate::message::ReactMessage;
use crate::protocol::{op::OpBatch, outbound::Outbound};
use crate::request::RawRequest;
#[derive(Clone)]
struct EventLoop {
outbound: Rc<Mutex<UnboundedReceiver<Outbound>>>,
reload: Rc<Mutex<UnboundedReceiver<()>>>,
reload_flag: Rc<Cell<bool>>,
reload_notify: Rc<Notify>,
}
#[op2(fast)]
fn op_flush(state: &mut OpState, #[string] json: &str, devtools: bool) -> Result<(), JsErrorBox> {
if crate::ext::with_thread_registry(|r| r.is_none()) {
let registry = state.borrow::<ExtSlot>().0.wait();
crate::ext::set_thread_registry(registry);
}
let ops: OpBatch =
serde_json::from_str(json).map_err(|e| JsErrorBox::type_error(e.to_string()))?;
let senders = state.borrow::<HostSenders>();
let _ = senders.flush.send(FlushInfo {
sent: Some(std::time::Instant::now()),
devtools,
});
let _ = senders.ops.send(ops.0);
Ok(())
}
#[op2]
#[serde]
fn op_take_decode_warnings() -> Vec<crate::diag::Warning> {
crate::diag::take_decode_warnings()
}
#[op2]
fn op_emit(state: &mut OpState, #[string] name: String, #[serde] value: serde_json::Value) {
let _ = state
.borrow::<HostSenders>()
.emit
.send(ReactMessage { name, value });
}
#[op2(fast)]
fn op_log(#[string] level: String, #[string] msg: String) {
match level.as_str() {
"error" => tracing::error!(target: "bevy_react_core::js", "{msg}"),
"warn" => tracing::warn!(target: "bevy_react_core::js", "{msg}"),
"debug" => tracing::debug!(target: "bevy_react_core::js", "{msg}"),
_ => tracing::info!(target: "bevy_react_core::js", "{msg}"),
}
crate::console_log::push(
crate::console_log::Source::Js,
crate::console_log::Level::from_js(&level),
&msg,
);
}
#[op2]
fn op_animate(state: &mut OpState, #[serde] cmd: AnimationCommand) {
let _ = state.borrow::<HostSenders>().anim.send(cmd);
}
#[op2]
fn op_request(
state: &mut OpState,
#[bigint] id: u64,
#[string] name: String,
#[serde] value: serde_json::Value,
) {
let _ = state
.borrow::<HostSenders>()
.request
.send(RawRequest { id, name, value });
}
#[op2]
#[serde]
async fn op_next_event(state: Rc<RefCell<OpState>>) -> Option<Outbound> {
let handles = state.borrow().borrow::<EventLoop>().clone();
let mut events = handles.outbound.lock().await;
let mut reload = handles.reload.lock().await;
tokio::select! {
ev = events.recv() => ev, r = reload.recv() => match r {
Some(()) => {
handles.reload_flag.set(true);
handles.reload_notify.notify_one();
Some(Outbound::Reload)
}
None => None, }
}
}
#[op2]
async fn op_sleep(ms: f64) {
let ms = ms.max(0.0) as u64;
tokio::time::sleep(std::time::Duration::from_millis(ms)).await;
}
#[op2(fast)]
fn op_now() -> f64 {
static ORIGIN: std::sync::LazyLock<std::time::Instant> =
std::sync::LazyLock::new(std::time::Instant::now);
ORIGIN.elapsed().as_secs_f64() * 1000.0
}
const PRELUDE: &str = r#"
let __nextTimer = 1;
const __cancelled = new Set();
globalThis.setTimeout = (cb, ms = 0, ...args) => {
const id = __nextTimer++;
const delay = Math.max(0, +ms || 0);
const run = () => { if (!__cancelled.delete(id)) cb(...args); };
if (delay === 0) Promise.resolve().then(run);
else Deno.core.ops.op_sleep(delay).then(run);
return id;
};
globalThis.clearTimeout = (id) => { if (id != null) __cancelled.add(id); };
globalThis.setInterval = (cb, ms = 0, ...args) => {
const id = __nextTimer++;
const delay = Math.max(0, +ms || 0);
(async () => {
while (!__cancelled.has(id)) {
await Deno.core.ops.op_sleep(delay);
if (__cancelled.has(id)) break;
cb(...args);
}
__cancelled.delete(id);
})();
return id;
};
globalThis.clearInterval = (id) => { if (id != null) __cancelled.add(id); };
globalThis.queueMicrotask = globalThis.queueMicrotask || ((cb) => { Promise.resolve().then(cb); });
if (!globalThis.performance) {
globalThis.performance = { now: () => Deno.core.ops.op_now(), timeOrigin: Date.now() };
}
const __fmtArg = (a) => {
if (typeof a === "string") return a;
if (a instanceof Error) return a.stack || (a.name + ": " + a.message);
try { return JSON.stringify(a); } catch { return String(a); }
};
const __log = (level) => (...args) =>
Deno.core.ops.op_log(level, args.map(__fmtArg).join(" "));
globalThis.console = {
log: __log("info"),
info: __log("info"),
debug: __log("debug"),
trace: __log("debug"),
warn: __log("warn"),
error: __log("error"),
dir: __log("info"),
table: __log("info"),
// No-op fallbacks so libraries that probe these never throw:
group: () => {}, groupCollapsed: () => {}, groupEnd: () => {}, assert: () => {},
};
// Unhandled promise rejections: log (→ op_log → the devtools console) and
// swallow. Returning true suppresses op_dispatch_exception — a deliberate
// behavior change: previously a rejection errored the event loop, which
// `pump` treats as a reload and re-executes the whole app bundle. A logged
// rejection with a live app beats a silent restart.
Deno.core.setUnhandledPromiseRejectionHandler((_promise, reason) => {
console.error("[js] unhandled promise rejection:", __fmtArg(reason));
return true;
});
"#;
enum Pumped {
Reload,
Shutdown,
}
struct ExtSlot(crate::ext::ExtRegistrySlot);
pub fn spawn_js_thread(
ext: crate::ext::ExtRegistrySlot,
vendor_path: PathBuf,
app_path: PathBuf,
senders: HostSenders,
outbound_rx: UnboundedReceiver<Outbound>,
reload_rx: UnboundedReceiver<()>,
) {
std::thread::Builder::new()
.name("js-runtime".to_string())
.spawn(move || {
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("build current-thread tokio runtime");
rt.block_on(async move {
let event_loop = EventLoop {
outbound: Rc::new(Mutex::new(outbound_rx)),
reload: Rc::new(Mutex::new(reload_rx)),
reload_flag: Rc::new(Cell::new(false)),
reload_notify: Rc::new(Notify::new()),
};
let mut last_good_app = match read_app(&app_path) {
Ok(code) => code,
Err(e) => {
tracing::error!(target: "bevy_react_core::js", "reading app failed: {e:?}");
return;
}
};
let mut runtime = match build_runtime(
&vendor_path,
&last_good_app,
&ext,
&senders,
&event_loop,
) {
Ok(rt) => rt,
Err(e) => {
tracing::error!(target: "bevy_react_core::js", "initial runtime build failed: {e:?}");
return;
}
};
loop {
event_loop.reload_flag.set(false);
match pump(&mut runtime, &event_loop).await {
Pumped::Shutdown => break,
Pumped::Reload => {
let new_code = match read_app(&app_path) {
Ok(code) => code,
Err(e) => {
tracing::warn!(target: "bevy_react_core::js", "reading rebuilt app failed ({e}); keeping the previous working version");
continue;
}
};
match runtime.execute_script("[app-update]", new_code.clone()) {
Ok(_) => last_good_app = new_code,
Err(e) => {
tracing::warn!(target: "bevy_react_core::js", "update rejected ({e}); keeping the previous working version");
crate::console_log::push(
crate::console_log::Source::Rust,
crate::console_log::Level::Error,
&format!("hot reload rejected ({e}); keeping the previous working version"),
);
if let Err(e) = runtime
.execute_script("[app-restore]", last_good_app.clone())
{
tracing::error!(target: "bevy_react_core::js", "restoring previous app failed: {e:?}");
crate::console_log::push(
crate::console_log::Source::Rust,
crate::console_log::Level::Error,
&format!("restoring previous app failed: {e:?}"),
);
}
}
}
}
}
}
});
})
.expect("spawn js-runtime thread");
}
fn build_runtime(
vendor_path: &Path,
app_code: &str,
registry: &crate::ext::ExtRegistrySlot,
senders: &HostSenders,
event_loop: &EventLoop,
) -> Result<JsRuntime, Box<dyn std::error::Error>> {
const FLUSH: OpDecl = op_flush();
const TAKE_WARNINGS: OpDecl = op_take_decode_warnings();
const EMIT: OpDecl = op_emit();
const REQUEST: OpDecl = op_request();
const ANIMATE: OpDecl = op_animate();
const NEXT: OpDecl = op_next_event();
const SLEEP: OpDecl = op_sleep();
const LOG: OpDecl = op_log();
const NOW: OpDecl = op_now();
let ext = Extension {
name: "bevy_react_bridge",
ops: std::borrow::Cow::Borrowed(&[
FLUSH,
TAKE_WARNINGS,
EMIT,
REQUEST,
ANIMATE,
NEXT,
SLEEP,
LOG,
NOW,
]),
..Default::default()
};
let mut runtime = JsRuntime::new(RuntimeOptions {
extensions: vec![ext],
..Default::default()
});
{
let op_state = runtime.op_state();
let mut op_state = op_state.borrow_mut();
op_state.put(ExtSlot(registry.clone()));
op_state.put(senders.clone());
op_state.put(event_loop.clone());
}
runtime.execute_script("[prelude]", PRELUDE)?;
let vendor_code = std::fs::read_to_string(vendor_path)
.map_err(|e| format!("reading vendor {}: {e}", vendor_path.display()))?;
runtime.execute_script("[vendor]", vendor_code)?;
runtime.execute_script("[app]", app_code.to_owned())?;
Ok(runtime)
}
async fn pump(runtime: &mut JsRuntime, event_loop: &EventLoop) -> Pumped {
let reload_flag = &event_loop.reload_flag;
loop {
let loop_result = tokio::select! {
biased;
res = runtime.run_event_loop(Default::default()) => Some(res),
_ = event_loop.reload_notify.notified() => None,
};
match loop_result {
None => {
if reload_flag.get() {
return Pumped::Reload;
}
}
Some(Err(e)) => {
tracing::error!(target: "bevy_react_core::js", "event loop error: {e}");
crate::console_log::push(
crate::console_log::Source::Rust,
crate::console_log::Level::Error,
&format!("JS event loop error: {e}"),
);
return Pumped::Reload;
}
Some(Ok(())) => {
return if reload_flag.get() {
Pumped::Reload
} else {
Pumped::Shutdown
};
}
}
}
}
fn read_app(app_path: &Path) -> Result<String, String> {
std::fs::read_to_string(app_path)
.map_err(|e| format!("reading app {}: {e}", app_path.display()))
}