use std::sync::Arc;
use std::sync::Mutex;
use std::time::Duration;
use deno_core::v8;
use tokio::sync::mpsc;
use tokio::sync::oneshot;
use crate::bridge::EngineHooks;
use crate::host::ModuleHost;
use crate::lock;
use crate::scheduler;
use crate::EngineConfig;
use crate::EngineError;
use crate::EvalInput;
pub(crate) type Reply = oneshot::Sender<Result<serde_json::Value, EngineError>>;
pub(crate) enum Job {
Eval {
input: EvalInput,
deadline: Option<Duration>,
reply: Reply,
},
Gc { reply: std::sync::mpsc::Sender<()> },
Call {
module: String,
export: String,
args: Vec<serde_json::Value>,
deadline: Option<Duration>,
reply: Reply,
},
}
impl Job {
pub(crate) fn refuse_condemned(self) {
match self {
Job::Eval { reply, .. } | Job::Call { reply, .. } => {
let _ = reply.send(Err(EngineError::MemoryLimit));
}
Job::Gc { reply } => {
let _ = reply.send(());
}
}
}
}
struct Running {
tx: mpsc::UnboundedSender<Job>,
thread: std::thread::JoinHandle<()>,
}
pub struct JsEngine {
running: Mutex<Option<Running>>,
isolate: v8::IsolateHandle,
default_deadline: Option<Duration>,
}
impl JsEngine {
pub fn spawn(
config: EngineConfig,
module_host: Option<ModuleHost>,
hooks: Option<EngineHooks>,
) -> Result<JsEngine, EngineError> {
deno_core::JsRuntime::init_platform(None);
let default_deadline = config.default_deadline;
let (tx, rx) = mpsc::unbounded_channel();
if let Some(registry) = &config.registry {
registry.register(&tx);
}
let (ready_tx, ready_rx) = std::sync::mpsc::channel();
let thread = std::thread::Builder::new()
.name("oj-js-engine".into())
.stack_size(8 * 1024 * 1024)
.spawn(move || scheduler::engine_thread(config, module_host, hooks, rx, ready_tx))
.map_err(|e| EngineError::Boot(e.to_string()))?;
match ready_rx.recv() {
Ok(Ok(isolate)) => Ok(JsEngine {
running: Mutex::new(Some(Running { tx, thread })),
isolate,
default_deadline,
}),
Ok(Err(e)) => {
let _ = thread.join();
Err(e)
}
Err(_) => {
let _ = thread.join();
Err(EngineError::Boot(
"engine thread exited during startup".into(),
))
}
}
}
pub fn collect_garbage(&self, wait: Duration) -> bool {
let (ack_tx, ack_rx) = std::sync::mpsc::channel();
let sent = match lock(&self.running).as_ref() {
Some(running) => running.tx.send(Job::Gc { reply: ack_tx }).is_ok(),
None => false,
};
sent && ack_rx.recv_timeout(wait).is_ok()
}
pub fn abandon(&self) {
self.isolate.terminate_execution();
drop(lock(&self.running).take());
}
pub async fn eval(
&self,
input: EvalInput,
deadline: Option<Duration>,
) -> Result<serde_json::Value, EngineError> {
let deadline = deadline.or(self.default_deadline);
self.request(|reply| Job::Eval {
input,
deadline,
reply,
})
.await
}
pub async fn call(
&self,
module: impl Into<String>,
export: &str,
args: Vec<serde_json::Value>,
deadline: Option<Duration>,
) -> Result<serde_json::Value, EngineError> {
let module = module.into();
let export = export.to_string();
let deadline = deadline.or(self.default_deadline);
self.request(|reply| Job::Call {
module,
export,
args,
deadline,
reply,
})
.await
}
async fn request(
&self,
make_job: impl FnOnce(Reply) -> Job,
) -> Result<serde_json::Value, EngineError> {
let (reply_tx, reply_rx) = oneshot::channel();
lock(&self.running)
.as_ref()
.ok_or(EngineError::Closed)?
.tx
.send(make_job(reply_tx))
.map_err(|_| EngineError::Closed)?;
reply_rx.await.map_err(|_| EngineError::Closed)?
}
}
impl Drop for JsEngine {
fn drop(&mut self) {
if let Some(Running { tx, thread }) = lock(&self.running).take() {
drop(tx);
let _ = thread.join();
}
}
}
#[derive(Clone, Default)]
pub struct EngineRegistry {
engines: Arc<Mutex<Vec<mpsc::WeakUnboundedSender<Job>>>>,
}
impl std::fmt::Debug for EngineRegistry {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "EngineRegistry({} slots)", lock(&self.engines).len())
}
}
impl EngineRegistry {
pub fn new() -> EngineRegistry {
EngineRegistry::default()
}
fn register(&self, tx: &mpsc::UnboundedSender<Job>) {
let mut engines = lock(&self.engines);
engines.retain(|weak| weak.upgrade().is_some_and(|tx| !tx.is_closed()));
engines.push(tx.downgrade());
}
pub fn collect_garbage(&self, wait: Duration) -> usize {
let (ack_tx, ack_rx) = std::sync::mpsc::channel();
let sent = lock(&self.engines)
.iter()
.filter_map(mpsc::WeakUnboundedSender::upgrade)
.filter(|tx| {
tx.send(Job::Gc {
reply: ack_tx.clone(),
})
.is_ok()
})
.count();
drop(ack_tx);
let deadline = std::time::Instant::now() + wait;
let mut collected = 0;
while collected < sent {
let left = deadline.saturating_duration_since(std::time::Instant::now());
if left.is_zero() || ack_rx.recv_timeout(left).is_err() {
break;
}
collected += 1;
}
collected
}
}