use super::*;
use std::sync::Arc;
pub struct PluginHost {
engine: Mutex<Option<std::sync::Arc<oj_js::JsEngine>>>,
host_module: String,
hook_plan: std::sync::RwLock<BuildHookPlan>,
hook_plan_fetched: std::sync::atomic::AtomicBool,
hook_plan_prime_started: std::sync::atomic::AtomicBool,
ws_out: std::sync::OnceLock<tokio::sync::broadcast::Sender<String>>,
server_events: std::sync::OnceLock<tokio::sync::mpsc::UnboundedSender<serde_json::Value>>,
serve_info_push: tokio::sync::watch::Sender<Option<ServeInfo>>,
initialized: tokio::sync::watch::Sender<bool>,
pub(crate) host_gone: tokio::sync::watch::Sender<bool>,
spawned: Mutex<tokio::time::Instant>,
pub(crate) init_wait: std::time::Duration,
lazy: bool,
init_failed: tokio::sync::watch::Sender<bool>,
init_knob: &'static str,
resync_done: tokio::sync::watch::Sender<u64>,
init_progress_seen: tokio::sync::watch::Sender<u64>,
rpc_wait: std::time::Duration,
init_progress: Mutex<std::time::Instant>,
boot: BootContext,
pub(crate) revive: Mutex<ReviveState>,
shut_down: std::sync::atomic::AtomicBool,
self_ref: std::sync::OnceLock<std::sync::Weak<PluginHost>>,
}
pub(crate) static DEV_LISTENER: Mutex<Option<(u16, String)>> = Mutex::new(None);
pub(crate) static LIVE_HOSTS: Mutex<Vec<std::sync::Weak<PluginHost>>> = Mutex::new(Vec::new());
pub fn dev_listener_bound(port: u16, interface: &str) {
*DEV_LISTENER.lock().unwrap() = Some((port, interface.to_string()));
LIVE_HOSTS.lock().unwrap().retain(|w| match w.upgrade() {
Some(host) => {
host.announce_dev_listener(port, interface);
true
}
None => false,
});
}
pub(crate) struct BootContext {
boot_seed: String,
root: PathBuf,
stall_wait: std::time::Duration,
memory_limit_bytes: usize,
registry: Option<oj_js::EngineRegistry>,
import_meta_env: Arc<std::sync::OnceLock<Arc<oj_compiler::ImportMetaEnv>>>,
}
pub(crate) struct ReviveState {
pub(crate) generation: u64,
pub(crate) attempts: u32,
pub(crate) last: Option<std::time::Instant>,
pub(crate) pending_before: std::collections::HashSet<PathBuf>,
pub(crate) reported: bool,
}
pub(crate) fn first_death_report(revive: &mut ReviveState, generation: u64) -> bool {
if revive.generation != generation || revive.reported {
return false;
}
revive.reported = true;
true
}
pub(crate) const PLUGIN_HOST_RESPAWN_LIMIT: u32 = 3;
pub(crate) const PLUGIN_HOST_RESPAWN_SPACING: std::time::Duration =
std::time::Duration::from_secs(5);
pub(crate) fn plugin_host_memory_mb() -> usize {
let knobs = &oj_env::get().knobs;
if let Some(mb) = knobs.plugin_memory_mb {
return mb;
}
knobs
.node_options
.as_deref()
.and_then(max_old_space_mb)
.unwrap_or(4096)
}
pub(crate) fn max_old_space_mb(node_options: &str) -> Option<usize> {
let mut found = None;
for token in node_options.split_whitespace() {
let norm = token.replace('_', "-");
if let Some(v) = norm.strip_prefix("--max-old-space-size=") {
if let Some(mb) = v.parse::<usize>().ok().filter(|m| *m > 0) {
found = Some(mb);
}
}
}
found
}
pub(crate) static ADDON_KEEPER: Mutex<Option<std::sync::Arc<oj_js::JsEngine>>> = Mutex::new(None);
pub(crate) const ADDON_KEEPER_DEADLINE: std::time::Duration = std::time::Duration::from_secs(4);
pub(crate) async fn keep_addons_alive(
root: &Path,
addons: &[PathBuf],
registry: Option<oj_js::EngineRegistry>,
) -> Result<(), String> {
let engine = {
let mut keeper = ADDON_KEEPER.lock().unwrap();
match &*keeper {
Some(engine) => std::sync::Arc::clone(engine),
None => {
let engine = std::sync::Arc::new({
let mut config = oj_js::EngineConfig::new(root);
config.registry = registry;
oj_js::JsEngine::spawn(config, None, None)
.map_err(|e| format!("keeper engine failed to spawn: {e}"))?
});
*keeper = Some(std::sync::Arc::clone(&engine));
engine
}
}
};
let live: std::collections::HashSet<PathBuf> = oj_js::addons_with_live_registrations()
.into_iter()
.collect();
let paths = serde_json::to_string(
&addons
.iter()
.filter(|p| live.contains(*p))
.map(|p| p.to_string_lossy())
.collect::<Vec<_>>(),
)
.map_err(|e| e.to_string())?;
let anchor = serde_json::to_string(&root.join("package.json").to_string_lossy())
.map_err(|e| e.to_string())?;
let script = format!(
r#"import {{ createRequire }} from "node:module";
const req = createRequire({anchor});
globalThis.__ojAddonKeeper ??= [];
for (const p of {paths}) {{
try {{ globalThis.__ojAddonKeeper.push(req(p)); }} catch {{}}
}}
"#
);
engine
.eval(
oj_js::EvalInput::Source(script),
Some(ADDON_KEEPER_DEADLINE),
)
.await
.map(|_| ())
.map_err(|e| format!("keeper load failed: {e}"))
}
pub(crate) fn process_rss_mb() -> Option<u64> {
#[cfg(target_os = "linux")]
{
let status = std::fs::read_to_string("/proc/self/status").ok()?;
let line = status.lines().find(|l| l.starts_with("VmRSS:"))?;
let kb: u64 = line.split_whitespace().nth(1)?.parse().ok()?;
return Some(kb / 1024);
}
#[cfg(all(unix, not(target_os = "linux")))]
{
let mut usage = std::mem::MaybeUninit::<libc::rusage>::zeroed();
if unsafe { libc::getrusage(libc::RUSAGE_SELF, usage.as_mut_ptr()) } != 0 {
return None;
}
let max_rss = unsafe { usage.assume_init() }.ru_maxrss.max(0) as u64;
if cfg!(target_os = "macos") {
return Some(max_rss / (1024 * 1024));
}
return Some(max_rss / 1024);
}
#[allow(unreachable_code)]
None
}
pub(crate) fn ctx_rpc(
method: &str,
args: &[serde_json::Value],
resolver: &OjResolver,
root: &Path,
env: Option<&Arc<oj_compiler::ImportMetaEnv>>,
) -> Result<serde_json::Value, String> {
let arg = |i: usize| args.get(i).and_then(|v| v.as_str()).unwrap_or("");
let dir_of = |path: &Path| {
path.parent()
.map(Path::to_path_buf)
.unwrap_or_else(|| root.to_path_buf())
};
match method {
"resolve" => {
let (source, importer) = (arg(0), arg(1));
let dir = if importer.is_empty() {
root.to_path_buf()
} else {
dir_of(Path::new(importer))
};
Ok(match resolver.resolve(&dir, source) {
Ok(p) => serde_json::Value::String(p.display().to_string()),
Err(_) => serde_json::Value::Null,
})
}
"moduleInfo" => {
let id = arg(0);
let path = Path::new(id);
let Ok(src) = std::fs::read_to_string(path) else {
return Ok(serde_json::Value::Null);
};
let dir = dir_of(path);
let opts = oj_compiler::CompileOptions {
env: env.cloned(),
..oj_compiler::CompileOptions::prod()
};
let (code, imports) = match oj_compiler::compile(path, &src, &opts) {
Ok(out) => (out.code, out.imports),
Err(_) => (src, Vec::new()),
};
let imported_ids: Vec<String> = imports
.iter()
.map(|spec| {
resolver
.resolve(&dir, spec)
.map(|p| p.display().to_string())
.unwrap_or_else(|_| spec.clone())
})
.collect();
Ok(serde_json::json!({ "id": id, "code": code, "importedIds": imported_ids }))
}
other => Err(format!("unknown ctx method: {other}")),
}
}
pub fn plugin_rpc_timeout() -> std::time::Duration {
plugin_rpc_timeout_from(oj_env::get().knobs.plugin_timeout.as_deref())
}
pub(crate) fn plugin_rpc_timeout_from(raw: Option<&str>) -> std::time::Duration {
let secs = raw
.and_then(|v| v.trim().parse::<u64>().ok())
.filter(|s| *s > 0)
.unwrap_or(20);
std::time::Duration::from_secs(secs)
}
pub fn plugin_init_timeout() -> std::time::Duration {
plugin_init_timeout_from(oj_env::get().knobs.plugin_init_timeout.as_deref())
}
pub(crate) fn plugin_init_timeout_from(raw: Option<&str>) -> std::time::Duration {
let secs = raw
.and_then(|v| v.trim().parse::<u64>().ok())
.filter(|s| *s > 0)
.unwrap_or(300);
std::time::Duration::from_secs(secs)
}
fn ws_push_payload(ws: &serde_json::Value) -> String {
match ws.get("event").and_then(|e| e.as_str()) {
Some(event) => serde_json::json!({
"type": "custom",
"event": event,
"data": ws.get("data").cloned().unwrap_or(serde_json::Value::Null),
})
.to_string(),
None => ws
.get("data")
.filter(|d| d.is_object())
.map(|d| d.to_string())
.unwrap_or_default(),
}
}
struct AbandonOnDrop(Option<std::sync::Arc<oj_js::JsEngine>>);
impl Drop for AbandonOnDrop {
fn drop(&mut self) {
if let Some(engine) = self.0.take() {
engine.abandon();
}
}
}
impl std::fmt::Debug for PluginHost {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str("PluginHost")
}
}
impl Drop for PluginHost {
fn drop(&mut self) {
if let Some(engine) = self
.engine
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.take()
{
engine.abandon();
}
}
}
pub(crate) fn call_init_deadline(
lazy: bool,
spawned: tokio::time::Instant,
init_wait: std::time::Duration,
now: tokio::time::Instant,
) -> tokio::time::Instant {
if lazy {
now + init_wait
} else {
spawned + init_wait
}
}
pub(crate) fn init_wait_policy(lazy: bool) -> (std::time::Duration, &'static str) {
if lazy {
(plugin_rpc_timeout(), "OJ_PLUGIN_TIMEOUT")
} else {
(plugin_init_timeout(), "OJ_PLUGIN_INIT_TIMEOUT")
}
}
#[derive(Default)]
pub(crate) struct SpawnTimeouts {
pub(crate) init_wait: Option<std::time::Duration>,
pub(crate) stall: Option<std::time::Duration>,
pub(crate) rpc: Option<std::time::Duration>,
pub(crate) memory: Option<usize>,
}
impl PluginHost {
pub async fn spawn(
root: &Path,
plugins_file: &Path,
config_json: &str,
registry: Option<oj_js::EngineRegistry>,
) -> anyhow::Result<std::sync::Arc<PluginHost>> {
Self::spawn_with_policy(
root,
plugins_file,
config_json,
false,
SpawnTimeouts::default(),
registry,
)
.await
}
pub async fn spawn_lazy(
root: &Path,
plugins_file: &Path,
config_json: &str,
registry: Option<oj_js::EngineRegistry>,
) -> anyhow::Result<std::sync::Arc<PluginHost>> {
Self::spawn_with_policy(
root,
plugins_file,
config_json,
true,
SpawnTimeouts::default(),
registry,
)
.await
}
#[cfg(test)]
pub(crate) async fn spawn_lazy_with_wait(
root: &Path,
plugins_file: &Path,
config_json: &str,
init_wait: std::time::Duration,
) -> anyhow::Result<std::sync::Arc<PluginHost>> {
Self::spawn_with_policy(
root,
plugins_file,
config_json,
true,
SpawnTimeouts {
init_wait: Some(init_wait),
..Default::default()
},
None,
)
.await
}
#[cfg(test)]
pub(crate) async fn spawn_with_timeouts(
root: &Path,
plugins_file: &Path,
config_json: &str,
lazy: bool,
timeouts: SpawnTimeouts,
) -> anyhow::Result<std::sync::Arc<PluginHost>> {
Self::spawn_with_policy(root, plugins_file, config_json, lazy, timeouts, None).await
}
async fn spawn_with_policy(
root: &Path,
plugins_file: &Path,
config_json: &str,
lazy: bool,
timeouts: SpawnTimeouts,
registry: Option<oj_js::EngineRegistry>,
) -> anyhow::Result<std::sync::Arc<PluginHost>> {
let script = oj_cache::cache_root(root).join("plugin-host.mjs");
let (mut init_wait, init_knob) = init_wait_policy(lazy);
if let Some(w) = timeouts.init_wait {
init_wait = w;
}
let rpc_wait = timeouts.rpc.unwrap_or_else(plugin_rpc_timeout);
let stall_wait = timeouts.stall.unwrap_or(rpc_wait);
let boot_seed = serde_json::json!({
"pluginsPath": plugins_file,
"initialJson": config_json,
"cacheRoot": oj_cache::cache_root(root),
});
let host = std::sync::Arc::new(PluginHost {
engine: Mutex::new(None),
host_module: script.to_string_lossy().into_owned(),
hook_plan: std::sync::RwLock::new(BuildHookPlan::fail_open()),
hook_plan_fetched: std::sync::atomic::AtomicBool::new(false),
hook_plan_prime_started: std::sync::atomic::AtomicBool::new(false),
ws_out: std::sync::OnceLock::new(),
server_events: std::sync::OnceLock::new(),
serve_info_push: tokio::sync::watch::channel(None).0,
initialized: tokio::sync::watch::channel(false).0,
host_gone: tokio::sync::watch::channel(false).0,
spawned: Mutex::new(tokio::time::Instant::now()),
init_wait,
lazy,
init_failed: tokio::sync::watch::channel(false).0,
init_knob,
resync_done: tokio::sync::watch::channel(0).0,
init_progress_seen: tokio::sync::watch::channel(0).0,
rpc_wait,
init_progress: Mutex::new(std::time::Instant::now()),
boot: BootContext {
boot_seed: boot_seed.to_string(),
root: root.to_path_buf(),
stall_wait,
memory_limit_bytes: timeouts
.memory
.unwrap_or_else(|| plugin_host_memory_mb() * 1024 * 1024),
registry,
import_meta_env: Arc::default(),
},
revive: Mutex::new(ReviveState {
generation: 0,
attempts: 0,
last: None,
pending_before: std::collections::HashSet::new(),
reported: false,
}),
shut_down: std::sync::atomic::AtomicBool::new(false),
self_ref: std::sync::OnceLock::new(),
});
let _ = host.self_ref.set(std::sync::Arc::downgrade(&host));
LIVE_HOSTS
.lock()
.unwrap()
.push(std::sync::Arc::downgrade(&host));
Self::ignite(&host, 0).map_err(|e| anyhow::anyhow!("{e}"))?;
Ok(host)
}
fn ignite(host: &std::sync::Arc<PluginHost>, generation: u64) -> Result<(), String> {
host.write_host_module()?;
let (post_tx, post_rx) = tokio::sync::mpsc::unbounded_channel();
let engine = std::sync::Arc::new(host.spawn_engine(post_tx)?);
{
let revive = host.revive.lock().unwrap();
if revive.generation != generation
|| host.shut_down.load(std::sync::atomic::Ordering::SeqCst)
{
drop(revive);
engine.abandon();
return Err("superseded by a newer respawn or a shutdown".into());
}
*host.engine.lock().unwrap() = Some(std::sync::Arc::clone(&engine));
}
Self::spawn_boot_task(host, engine, generation);
Self::spawn_push_dispatcher(host, post_rx, generation);
Self::spawn_stall_monitor(host);
Ok(())
}
fn write_host_module(&self) -> Result<(), String> {
let script = PathBuf::from(&self.host_module);
if let Some(parent) = script.parent() {
std::fs::create_dir_all(parent).map_err(|e| e.to_string())?;
}
let _ = ensure_asset(
oj_cache::cache_root(&self.boot.root).as_path(),
"discovered-deps.mjs",
DISCOVERED_DEPS_JS,
);
if std::fs::read(&script).ok().as_deref() != Some(PLUGIN_HOST_JS.as_bytes()) {
let tmp = script.with_extension(format!("tmp-{}.mjs", std::process::id()));
std::fs::write(&tmp, PLUGIN_HOST_JS).map_err(|e| e.to_string())?;
std::fs::rename(&tmp, &script).map_err(|e| e.to_string())?;
}
Ok(())
}
fn spawn_engine(
&self,
post_tx: tokio::sync::mpsc::UnboundedSender<serde_json::Value>,
) -> Result<oj_js::JsEngine, String> {
let root = self.boot.root.clone();
let resolver = std::sync::Arc::new(OjResolver::new(&root));
let rpc_handler: oj_js::RpcHandler = {
let root = root.clone();
let env = Arc::clone(&self.boot.import_meta_env);
Box::new(move |method, args| ctx_rpc(method, args, &resolver, &root, env.get()))
};
let mut engine_config = oj_js::EngineConfig::new(&root);
engine_config.code_cache_dir = Some(crate::engine_code_cache_dir(&root));
engine_config.memory_limit_bytes = Some(self.boot.memory_limit_bytes);
engine_config.registry = self.boot.registry.clone();
oj_js::JsEngine::spawn(
engine_config,
None,
Some(oj_js::EngineHooks {
post: post_tx,
rpc: Some(rpc_handler),
}),
)
.map_err(|e| format!("cannot start the embedded plugin host: {e}"))
}
fn spawn_boot_task(
host: &std::sync::Arc<PluginHost>,
engine: std::sync::Arc<oj_js::JsEngine>,
generation: u64,
) {
let host = std::sync::Arc::clone(host);
let prelude = format!("globalThis.__ojPluginHost = {};", host.boot.boot_seed);
let host_module = host.host_module.clone();
tokio::spawn(async move {
if let Err(e) = engine.eval(oj_js::EvalInput::Source(prelude), None).await {
host.declare_gone(&format!("plugin host boot prelude failed: {e}"), generation);
return;
}
match engine
.call(host_module, "ojHostReady", Vec::new(), None)
.await
{
Ok(_) => {
let listener = DEV_LISTENER.lock().unwrap().clone();
if let Some((port, interface)) = listener {
host.announce_dev_listener(port, &interface);
}
}
Err(oj_js::EngineError::Closed) => {}
Err(e) => {
host.declare_gone(
&format!("plugin host failed to initialize: {e}"),
generation,
);
}
}
});
}
fn spawn_push_dispatcher(
host: &std::sync::Arc<PluginHost>,
mut post_rx: tokio::sync::mpsc::UnboundedReceiver<serde_json::Value>,
generation: u64,
) {
let host = std::sync::Arc::clone(host);
tokio::spawn(async move {
while let Some(msg) = post_rx.recv().await {
if *host.host_gone.borrow() || host.revive.lock().unwrap().generation != generation
{
continue;
}
host.dispatch_push(&msg);
}
let revive = host.revive.lock().unwrap();
if revive.generation == generation {
drop(revive);
let _ = host.host_gone.send_replace(true);
}
});
}
fn dispatch_push(&self, msg: &serde_json::Value) {
if let Some(info) = msg.get("ojServeInfo") {
self.serve_info_push
.send_replace(Some(ServeInfo::from_json(info)));
self.mark_initialized();
} else if msg.get("ojInit").is_some() {
self.mark_initialized();
} else if msg.get("ojResyncDone").is_some() {
self.resync_done.send_modify(|c| *c += 1);
} else if msg.get("ojInitProgress").is_some() {
self.init_progress_seen.send_modify(|c| *c += 1);
let _ = self.init_failed.send_replace(false);
} else if let Some(ev) = msg.get("ojServer") {
if let Some(tx) = self.server_events.get() {
let _ = tx.send(ev.clone());
}
} else if let Some(ws) = msg.get("ojWs") {
if let Some(tx) = self.ws_out.get() {
let payload = ws_push_payload(ws);
if !payload.is_empty() {
let _ = tx.send(payload);
}
}
}
}
fn mark_initialized(&self) {
let _ = self.initialized.send_replace(true);
let _ = self.init_failed.send_replace(false);
}
fn spawn_stall_monitor(host: &std::sync::Arc<PluginHost>) {
let host = std::sync::Arc::clone(host);
let stall_wait = host.boot.stall_wait;
tokio::spawn(async move {
let mut init_rx = host.initialized.subscribe();
let mut gone_rx = host.host_gone.subscribe();
let mut prog_rx = host.init_progress_seen.subscribe();
loop {
if *init_rx.borrow_and_update() || *gone_rx.borrow_and_update() {
return;
}
let _ = prog_rx.borrow_and_update();
let deadline = tokio::time::Instant::now() + stall_wait;
let mut stalled = false;
tokio::select! {
biased;
changed = init_rx.changed() => { if changed.is_err() { return; } }
changed = gone_rx.changed() => { if changed.is_err() { return; } }
changed = prog_rx.changed() => { if changed.is_err() { return; } }
_ = tokio::time::sleep_until(deadline) => { stalled = true; }
}
if stalled {
let _ = host.init_failed.send_replace(true);
loop {
tokio::select! {
biased;
changed = init_rx.changed() => {
if changed.is_err() || *init_rx.borrow() { return; }
}
changed = gone_rx.changed() => {
if changed.is_err() || *gone_rx.borrow() { return; }
}
changed = prog_rx.changed() => {
if changed.is_err() { return; }
break;
}
}
}
}
}
});
}
pub fn can_revive(&self) -> bool {
!self.shut_down.load(std::sync::atomic::Ordering::SeqCst)
&& self.revive.lock().unwrap().attempts < PLUGIN_HOST_RESPAWN_LIMIT
}
fn try_revive(&self) -> bool {
if self.shut_down.load(std::sync::atomic::Ordering::SeqCst) {
return false;
}
let Some(host) = self.self_ref.get().and_then(std::sync::Weak::upgrade) else {
return false;
};
let mut revive = self.revive.lock().unwrap();
if !*self.host_gone.borrow() {
return true;
}
if revive.attempts >= PLUGIN_HOST_RESPAWN_LIMIT {
return false;
}
if let Some(last) = revive.last {
if last.elapsed() < PLUGIN_HOST_RESPAWN_SPACING {
return false;
}
}
let orphaned = oj_js::addons_pending_unsafe_reregistration();
if let Some(addon) = orphaned
.iter()
.find(|p| !revive.pending_before.contains(*p))
{
revive.last = Some(std::time::Instant::now());
eprintln!(
"oj: not respawning the plugin host: native addon {} was torn down with it and re-registering can crash (napi-rs before 3.10); restart the dev server to recover",
addon.display()
);
return false;
}
revive.attempts += 1;
revive.last = Some(std::time::Instant::now());
revive.generation += 1;
revive.reported = false;
eprintln!(
"oj: respawning the plugin host (attempt {} of {PLUGIN_HOST_RESPAWN_LIMIT})",
revive.attempts
);
self.reset_generation_state();
let generation = revive.generation;
drop(revive);
match Self::ignite(&host, generation) {
Ok(()) => {
let _ = self.host_gone.send_replace(false);
if let Ok(handle) = tokio::runtime::Handle::try_current() {
let host = std::sync::Arc::clone(&host);
handle.spawn(async move { host.ensure_hook_plan().await });
}
true
}
Err(e) => {
eprintln!("oj: plugin host respawn failed: {e}");
false
}
}
}
fn reset_generation_state(&self) {
use std::sync::atomic::Ordering;
let _ = self.initialized.send_replace(false);
let _ = self.init_failed.send_replace(false);
self.serve_info_push.send_replace(None);
*self.hook_plan.write().unwrap() = BuildHookPlan::fail_open();
self.hook_plan_fetched.store(false, Ordering::Release);
self.hook_plan_prime_started.store(false, Ordering::Release);
*self.spawned.lock().unwrap() = tokio::time::Instant::now();
}
fn announce_dev_listener(self: &std::sync::Arc<Self>, port: u16, interface: &str) {
let host = std::sync::Arc::clone(self);
let interface = interface.to_string();
tokio::spawn(async move {
let _ = host
.call("serverListening", &[&port.to_string(), &interface])
.await;
});
}
async fn call(&self, hook: &str, args: &[&str]) -> Result<Option<String>, String> {
if *self.host_gone.borrow() && !self.try_revive() {
return Err("plugin host exited".into());
}
self.wait_for_init(hook).await?;
self.run_hook(hook, args).await
}
async fn wait_for_init(&self, hook: &str) -> Result<(), String> {
let mut init_rx = self.initialized.subscribe();
if *init_rx.borrow_and_update() {
return Ok(());
}
let deadline = call_init_deadline(
self.lazy,
*self.spawned.lock().unwrap(),
self.init_wait,
tokio::time::Instant::now(),
);
let mut host_gone_rx = self.host_gone.subscribe();
if *host_gone_rx.borrow_and_update() {
return Err("plugin host exited".into());
}
let mut progress = tokio::time::interval_at(
tokio::time::Instant::now() + std::time::Duration::from_secs(30),
std::time::Duration::from_secs(30),
);
loop {
tokio::select! {
biased;
changed = init_rx.changed() => {
if changed.is_err() || *init_rx.borrow() {
return Ok(());
}
}
changed = host_gone_rx.changed() => {
if changed.is_err() || *host_gone_rx.borrow() {
return Err("plugin host exited".into());
}
}
_ = tokio::time::sleep_until(deadline) => {
let _ = self.init_failed.send_replace(true);
return Err(format!(
"plugin host still initializing after {}s running {hook} (raise {} for slower boots)",
self.init_wait.as_secs(),
self.init_knob,
));
}
_ = progress.tick() => self.log_init_progress(),
}
}
}
fn log_init_progress(&self) {
let elapsed = self.spawned.lock().unwrap().elapsed().as_secs();
let mut last = self
.init_progress
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if last.elapsed().as_secs() >= 29 {
*last = std::time::Instant::now();
eprintln!("oj: plugin host still initializing ({elapsed}s)…");
}
}
async fn run_hook(&self, hook: &str, args: &[&str]) -> Result<Option<String>, String> {
let deadline = tokio::time::Instant::now() + self.rpc_wait;
let generation = self.revive.lock().unwrap().generation;
let engine = self
.engine
.lock()
.unwrap()
.clone()
.ok_or_else(|| "plugin host exited".to_string())?;
let call = engine.call(
self.host_module.clone(),
"ojRun",
vec![
serde_json::Value::String(hook.to_string()),
serde_json::Value::Array(
args.iter()
.map(|a| serde_json::Value::String((*a).to_string()))
.collect(),
),
],
Some(self.rpc_wait),
);
tokio::pin!(call);
let result = tokio::select! {
biased;
r = &mut call => r,
_ = self.host_gone_wait() => return Err("plugin host exited".into()),
_ = tokio::time::sleep_until(deadline + self.rpc_wait) => {
let msg = format!(
"plugin host unresponsive for {}s running {hook} (the engine stopped scheduling)",
2 * self.rpc_wait.as_secs()
);
self.declare_gone(&msg, generation);
return Err(msg);
}
};
self.finish_call(hook, generation, result)
}
fn finish_call(
&self,
hook: &str,
generation: u64,
result: Result<serde_json::Value, oj_js::EngineError>,
) -> Result<Option<String>, String> {
match result {
Ok(value) => {
if self.revive.lock().unwrap().generation == generation {
self.mark_initialized();
}
Ok(match value {
serde_json::Value::Null => None,
serde_json::Value::String(s) => Some(s),
other => Some(other.to_string()),
})
}
Err(oj_js::EngineError::Deadline) => Err(format!(
"plugin host timed out after {}s running {hook} (raise OJ_PLUGIN_TIMEOUT for slow plugins)",
self.rpc_wait.as_secs()
)),
Err(oj_js::EngineError::Closed) => Err("plugin host exited".into()),
Err(oj_js::EngineError::MemoryLimit) => {
let msg = format!(
"plugin host exceeded its memory limit ({}MB; raise OJ_PLUGIN_MEMORY_MB) running {hook}",
self.boot.memory_limit_bytes / (1024 * 1024)
);
self.declare_gone(&msg, generation);
Err(msg)
}
Err(oj_js::EngineError::Boot(e)) => Err(e),
Err(oj_js::EngineError::Js(e)) => {
self.mark_initialized();
Err(e)
}
}
}
pub(crate) fn declare_gone(&self, why: &str, generation: u64) {
let mut revive = self.revive.lock().unwrap();
if !first_death_report(&mut revive, generation) {
return;
}
revive.pending_before = oj_js::addons_pending_unsafe_reregistration()
.into_iter()
.collect();
let rss = process_rss_mb()
.map(|m| format!(" (process rss {m}MB)"))
.unwrap_or_default();
eprintln!("oj: {why}; treating the plugin host as gone{rss}");
if let Some(engine) = self.engine.lock().unwrap().take() {
self.retire_engine(engine, &mut revive);
}
drop(revive);
let _ = self.host_gone.send_replace(true);
}
fn retire_engine(&self, engine: std::sync::Arc<oj_js::JsEngine>, revive: &mut ReviveState) {
let addons = oj_js::addons_with_live_registrations();
if addons.is_empty() {
engine.abandon();
return;
}
revive.last = Some(std::time::Instant::now());
let root = self.boot.root.clone();
let registry = self.boot.registry.clone();
let guard = AbandonOnDrop(Some(engine));
tokio::spawn(async move {
let _guard = guard;
if let Err(e) = keep_addons_alive(&root, &addons, registry).await {
eprintln!(
"oj: native-addon keeper unavailable ({e}); a plugin host respawn that would re-register an orphaned addon will be refused"
);
}
});
}
pub fn is_initialized(&self) -> bool {
*self.initialized.borrow()
}
pub fn initialized_updates(&self) -> tokio::sync::watch::Receiver<bool> {
self.initialized.subscribe()
}
pub fn init_deadline_at(&self) -> tokio::time::Instant {
*self.spawned.lock().unwrap() + self.init_wait
}
pub fn host_gone_updates(&self) -> tokio::sync::watch::Receiver<bool> {
self.host_gone.subscribe()
}
pub fn init_failure_updates(&self) -> tokio::sync::watch::Receiver<bool> {
self.init_failed.subscribe()
}
pub fn resync_done_updates(&self) -> tokio::sync::watch::Receiver<u64> {
self.resync_done.subscribe()
}
pub(crate) async fn host_gone_wait(&self) {
let mut rx = self.host_gone.subscribe();
while !*rx.borrow_and_update() {
if rx.changed().await.is_err() {
return;
}
}
}
pub async fn transform(
&self,
code: &str,
id: &str,
resolved: &str,
) -> Result<(String, Vec<String>, Vec<String>, Vec<ChunkEmit>), String> {
let Some(raw) = self.call("transform", &[code, id, resolved]).await? else {
return Ok((code.to_string(), Vec::new(), Vec::new(), Vec::new()));
};
let Ok(v) = serde_json::from_str::<serde_json::Value>(&raw) else {
return Ok((raw, Vec::new(), Vec::new(), Vec::new()));
};
let out = v
.get("code")
.and_then(|c| c.as_str())
.unwrap_or(code)
.to_string();
let str_array = |key: &str| {
v.get(key)
.and_then(|w| w.as_array())
.map(|a| {
a.iter()
.filter_map(|x| x.as_str().map(str::to_string))
.collect()
})
.unwrap_or_default()
};
Ok((
out,
str_array("watchFiles"),
str_array("maps"),
emitted_chunks(&v),
))
}
pub async fn seed_chunk_names(&self, map_json: &str) -> Result<Option<String>, String> {
self.call("seedChunkNames", &[map_json]).await
}
#[inline]
pub async fn has_module_parsed(&self) -> bool {
self.call_flag("hasModuleParsed", false).await
}
#[inline]
pub async fn module_parsed(&self, id: &str) -> Result<(), String> {
self.call_unit("replayModuleParsed", &[id]).await
}
#[inline]
pub async fn resolve_id(&self, source: &str, importer: &str) -> Result<Option<String>, String> {
self.call("resolveId", &[source, importer]).await
}
#[inline]
pub async fn load(&self, id: &str) -> Result<Option<String>, String> {
self.call("load", &[id]).await
}
#[inline]
pub async fn handle_hot_update(
&self,
file: &str,
timestamp: u64,
change_type: &str,
modules_json: &str,
) -> Result<Option<String>, String> {
self.call(
"handleHotUpdate",
&[file, ×tamp.to_string(), change_type, modules_json],
)
.await
}
#[inline]
pub async fn transform_index_html(&self, html: &str, ctx_json: &str) -> Result<String, String> {
Ok(self
.call("transformIndexHtml", &[html, ctx_json])
.await?
.unwrap_or_else(|| html.to_string()))
}
#[inline]
pub async fn build_start(&self) -> Result<Vec<ChunkEmit>, String> {
let Some(raw) = self.call("buildStart", &[]).await? else {
return Ok(Vec::new());
};
Ok(serde_json::from_str::<serde_json::Value>(&raw)
.map(|v| emitted_chunks(&v))
.unwrap_or_default())
}
#[inline]
pub async fn build_end(&self, error: Option<&str>) -> Result<(), String> {
match error {
Some(e) => self.call_unit("buildEnd", &[e]).await,
None => self.call_unit("buildEnd", &[]).await,
}
}
#[inline]
pub async fn render_start(&self) -> Result<(), String> {
self.call_unit("renderStart", &[]).await
}
#[inline]
pub async fn watch_change(&self, file: &str, event: &str) -> Result<(), String> {
self.call_unit("watchChange", &[file, event]).await
}
#[inline]
pub async fn close_bundle(&self) -> Result<(), String> {
self.call_unit("closeBundle", &[]).await
}
#[inline]
pub async fn watch_files(&self) -> Result<Vec<String>, String> {
let Some(json) = self.call("getWatchFiles", &[]).await? else {
return Ok(Vec::new());
};
serde_json::from_str(&json).map_err(|e| e.to_string())
}
#[inline]
pub async fn has_generate_bundle(&self) -> bool {
self.call_flag("hasGenerateBundle", false).await
}
#[inline]
pub async fn has_transform_index_html(&self) -> bool {
!matches!(self.call("hasTransformIndexHtml", &[]).await, Ok(Some(s)) if s == "false")
}
pub async fn prime_hook_plan(self: &std::sync::Arc<Self>) {
use std::sync::atomic::Ordering;
if self.hook_plan_fetched.load(Ordering::Acquire) {
return;
}
if self.hook_plan_prime_started.swap(true, Ordering::AcqRel) {
return;
}
let fetch = self.ensure_hook_plan();
if tokio::time::timeout(std::time::Duration::from_millis(500), fetch)
.await
.is_err()
{
let host = std::sync::Arc::clone(self);
tokio::spawn(async move { host.ensure_hook_plan().await });
}
}
pub async fn build_hook_plan(&self) -> BuildHookPlan {
self.ensure_hook_plan().await;
self.hook_plan.read().unwrap().clone()
}
pub async fn warm_environments(&self, urls: &[String]) -> Result<Option<String>, String> {
let payload = serde_json::to_string(urls).map_err(|e| e.to_string())?;
self.call("warmEnvironments", &[&payload]).await
}
async fn ensure_hook_plan(&self) {
use std::sync::atomic::Ordering;
if self.hook_plan_fetched.load(Ordering::Acquire) {
return;
}
if let Some(v) = self.call_json("getBuildHookPlan").await {
let plan = BuildHookPlan {
transform: HookFilterPlan::from_json(v.get("transform")),
load: HookFilterPlan::from_json(v.get("load")),
resolve_id: HookFilterPlan::from_json(v.get("resolveId")),
};
*self.hook_plan.write().unwrap() = plan;
self.hook_plan_fetched.store(true, Ordering::Release);
}
}
#[inline]
pub fn hook_wants_transform(&self, id: &str, code: &str) -> bool {
self.hook_plan
.read()
.unwrap()
.transform
.wants(id, Some(code))
}
#[inline]
pub fn hook_wants_load(&self, id: &str) -> bool {
self.hook_plan.read().unwrap().load.wants(id, None)
}
#[inline]
pub fn hook_wants_resolve_id(&self, spec: &str) -> bool {
self.hook_plan.read().unwrap().resolve_id.wants(spec, None)
}
#[inline]
pub async fn generate_bundle(
&self,
bundle_json: &str,
is_write: bool,
) -> Result<Option<String>, String> {
self.call("generateBundle", &[bundle_json, bool_arg(is_write)])
.await
}
#[inline]
pub async fn has_render_chunk(&self) -> bool {
self.call_flag("hasRenderChunk", false).await
}
pub async fn render_chunk(
&self,
code: &str,
chunk_json: &str,
) -> Result<Option<String>, String> {
self.call("renderChunk", &[code, chunk_json]).await
}
#[inline]
pub async fn has_write_bundle(&self) -> bool {
self.call_flag("hasWriteBundle", false).await
}
#[inline]
pub async fn write_bundle(&self, bundle_json: &str, is_write: bool) -> Result<(), String> {
self.call_unit("writeBundle", &[bundle_json, bool_arg(is_write)])
.await
}
pub async fn serve_info(&self) -> ServeInfo {
if let Some(info) = *self.serve_info_push.borrow() {
return info;
}
let rpc = self.call("getServeInfo", &[]).await;
if let Some(info) = *self.serve_info_push.borrow() {
return info;
}
rpc.ok()
.flatten()
.and_then(|s| serde_json::from_str::<serde_json::Value>(&s).ok())
.map(|v| ServeInfo::from_json(&v))
.unwrap_or_default()
}
pub fn serve_info_updates(&self) -> tokio::sync::watch::Receiver<Option<ServeInfo>> {
self.serve_info_push.subscribe()
}
pub async fn plugin_count(&self) -> usize {
self.call_ok("getPluginCount")
.await
.and_then(|s| s.parse().ok())
.unwrap_or(1)
}
pub async fn config_defines(&self) -> Vec<(String, String)> {
let Some(v) = self.call_json("getPluginConfig").await else {
return Vec::new();
};
v.get("define")
.and_then(|d| d.as_object())
.map(|d| {
d.iter()
.map(|(k, v)| {
let expr = match v {
serde_json::Value::String(s) => s.clone(),
other => other.to_string(),
};
(k.clone(), expr)
})
.collect()
})
.unwrap_or_default()
}
pub async fn env_delta(&self) -> std::collections::BTreeMap<String, String> {
self.call_ok("getEnvDelta")
.await
.and_then(|s| serde_json::from_str(&s).ok())
.unwrap_or_default()
}
pub async fn has_transform(&self) -> bool {
self.call_flag("getHasTransform", true).await
}
pub async fn has_load(&self) -> bool {
self.call_flag("getHasLoad", false).await
}
pub async fn dep_transform_filters(&self) -> Vec<String> {
self.call_ok("getDepTransformFilters")
.await
.and_then(|s| serde_json::from_str(&s).ok())
.unwrap_or_default()
}
pub async fn dep_load_filters(&self) -> Vec<String> {
self.call_ok("getDepLoadFilters")
.await
.and_then(|s| serde_json::from_str(&s).ok())
.unwrap_or_default()
}
pub async fn resolve_id_filters(&self) -> Vec<String> {
self.call_ok("getResolveIdFilters")
.await
.and_then(|s| serde_json::from_str(&s).ok())
.unwrap_or_default()
}
pub async fn hmr_hooks(&self) -> (bool, bool) {
let Some(v) = self.call_json("getHmrHooks").await else {
return (true, true);
};
let flag = |key: &str| v.get(key).and_then(|b| b.as_bool()).unwrap_or(true);
(flag("watchChange"), flag("handleHotUpdate"))
}
pub fn shutdown(&self) {
self.shut_down
.store(true, std::sync::atomic::Ordering::SeqCst);
let _revive = self.revive.lock().unwrap();
if let Some(engine) = self.engine.lock().unwrap().take() {
engine.abandon();
}
}
pub fn set_server_events_sender(
&self,
tx: tokio::sync::mpsc::UnboundedSender<serde_json::Value>,
) {
let _ = self.server_events.set(tx);
}
pub fn set_import_meta_env(&self, env: Arc<oj_compiler::ImportMetaEnv>) {
let _ = self.boot.import_meta_env.set(env);
}
pub fn set_ws_sender(&self, tx: tokio::sync::broadcast::Sender<String>) {
let _ = self.ws_out.set(tx);
}
#[inline]
pub async fn ws_message(&self, event: &str, data: &str) -> Result<(), String> {
self.call_unit("wsMessage", &[event, data]).await
}
#[inline]
pub async fn ws_connection(&self) -> Result<(), String> {
self.call_unit("wsConnection", &[]).await
}
#[inline]
pub async fn emitted_files(&self) -> Result<Vec<EmittedFile>, String> {
let Some(json) = self.call("getEmittedFiles", &[]).await? else {
return Ok(Vec::new());
};
let arr: Vec<serde_json::Value> = serde_json::from_str(&json).map_err(|e| e.to_string())?;
Ok(arr
.into_iter()
.filter_map(|v| {
Some(EmittedFile {
file_name: v.get("fileName")?.as_str()?.to_string(),
source: v.get("source")?.as_str()?.to_string(),
})
})
.collect())
}
pub async fn get_plugin_css(&self) -> Vec<(String, String)> {
self.call_json("getPluginCss")
.await
.and_then(|v| {
v.as_array().map(|a| {
a.iter()
.filter_map(|e| {
let css = e.get("css")?.as_str()?.to_string();
let id = e
.get("id")
.and_then(|x| x.as_str())
.unwrap_or("")
.to_string();
Some((id, css))
})
.collect()
})
})
.unwrap_or_default()
}
async fn call_unit(&self, hook: &str, args: &[&str]) -> Result<(), String> {
self.call(hook, args).await.map(|_| ())
}
async fn call_ok(&self, hook: &str) -> Option<String> {
self.call(hook, &[]).await.ok().flatten()
}
async fn call_json(&self, hook: &str) -> Option<serde_json::Value> {
serde_json::from_str(&self.call_ok(hook).await?).ok()
}
async fn call_flag(&self, hook: &str, default: bool) -> bool {
match self.call(hook, &[]).await {
Ok(Some(s)) => s == "true",
_ => default,
}
}
}
fn bool_arg(b: bool) -> &'static str {
if b {
"true"
} else {
"false"
}
}
fn emitted_chunks(v: &serde_json::Value) -> Vec<ChunkEmit> {
v.get("emittedChunks")
.and_then(|c| c.as_array())
.map(|a| a.iter().filter_map(ChunkEmit::from_value).collect())
.unwrap_or_default()
}