2020-01-02 15:13:47 -05:00
|
|
|
// Copyright 2018-2020 the Deno authors. All rights reserved. MIT license.
|
2019-08-26 14:50:21 +02:00
|
|
|
use super::dispatch_json::{Deserialize, JsonOp, Value};
|
2019-11-09 21:07:14 +01:00
|
|
|
use crate::deno_error::bad_resource;
|
2019-09-26 00:46:58 +10:00
|
|
|
use crate::deno_error::js_check;
|
2019-08-14 17:03:02 +02:00
|
|
|
use crate::deno_error::DenoError;
|
|
|
|
use crate::deno_error::ErrorKind;
|
2019-10-11 11:41:54 -07:00
|
|
|
use crate::ops::json_op;
|
2019-08-14 17:03:02 +02:00
|
|
|
use crate::startup_data;
|
2020-02-08 20:34:31 +01:00
|
|
|
use crate::state::State;
|
2020-01-21 09:49:47 +01:00
|
|
|
use crate::web_worker::WebWorker;
|
2020-02-08 20:34:31 +01:00
|
|
|
use crate::worker::WorkerChannelsExternal;
|
2020-01-05 11:56:18 -05:00
|
|
|
use deno_core::*;
|
2019-08-14 17:03:02 +02:00
|
|
|
use futures;
|
2019-11-17 01:17:47 +01:00
|
|
|
use futures::future::FutureExt;
|
|
|
|
use futures::future::TryFutureExt;
|
2019-08-14 17:03:02 +02:00
|
|
|
use std;
|
|
|
|
use std::convert::From;
|
2019-10-11 11:41:54 -07:00
|
|
|
use std::sync::atomic::Ordering;
|
|
|
|
|
2020-02-08 20:34:31 +01:00
|
|
|
pub fn init(i: &mut Isolate, s: &State) {
|
2019-10-11 11:41:54 -07:00
|
|
|
i.register_op(
|
|
|
|
"create_worker",
|
|
|
|
s.core_op(json_op(s.stateful_op(op_create_worker))),
|
|
|
|
);
|
2020-01-18 00:43:53 +01:00
|
|
|
i.register_op(
|
|
|
|
"host_close_worker",
|
|
|
|
s.core_op(json_op(s.stateful_op(op_host_close_worker))),
|
|
|
|
);
|
2019-10-11 11:41:54 -07:00
|
|
|
i.register_op(
|
|
|
|
"host_post_message",
|
|
|
|
s.core_op(json_op(s.stateful_op(op_host_post_message))),
|
|
|
|
);
|
|
|
|
i.register_op(
|
|
|
|
"host_get_message",
|
|
|
|
s.core_op(json_op(s.stateful_op(op_host_get_message))),
|
|
|
|
);
|
|
|
|
i.register_op("metrics", s.core_op(json_op(s.stateful_op(op_metrics))));
|
|
|
|
}
|
2019-08-14 17:03:02 +02:00
|
|
|
|
2019-08-26 14:50:21 +02:00
|
|
|
#[derive(Deserialize)]
|
|
|
|
#[serde(rename_all = "camelCase")]
|
|
|
|
struct CreateWorkerArgs {
|
2020-01-29 18:54:23 +01:00
|
|
|
name: Option<String>,
|
2019-08-26 14:50:21 +02:00
|
|
|
specifier: String,
|
|
|
|
has_source_code: bool,
|
|
|
|
source_code: String,
|
2019-08-14 17:03:02 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
/// Create worker as the host
|
2019-10-11 11:41:54 -07:00
|
|
|
fn op_create_worker(
|
2020-02-08 20:34:31 +01:00
|
|
|
state: &State,
|
2019-08-26 14:50:21 +02:00
|
|
|
args: Value,
|
2020-01-24 15:10:49 -05:00
|
|
|
_data: Option<ZeroCopyBuf>,
|
2019-08-26 14:50:21 +02:00
|
|
|
) -> Result<JsonOp, ErrBox> {
|
|
|
|
let args: CreateWorkerArgs = serde_json::from_value(args)?;
|
|
|
|
|
2020-02-03 18:08:44 -05:00
|
|
|
let specifier = args.specifier.clone();
|
2019-08-26 14:50:21 +02:00
|
|
|
let has_source_code = args.has_source_code;
|
2020-02-03 18:08:44 -05:00
|
|
|
let source_code = args.source_code.clone();
|
|
|
|
let args_name = args.name;
|
2019-08-14 17:03:02 +02:00
|
|
|
let parent_state = state.clone();
|
2020-02-08 20:34:31 +01:00
|
|
|
let state = state.borrow();
|
|
|
|
let global_state = state.global_state.clone();
|
|
|
|
let child_permissions = state.permissions.clone();
|
|
|
|
let referrer = state.main_module.to_string();
|
|
|
|
drop(state);
|
|
|
|
|
|
|
|
let (handle_sender, handle_receiver) =
|
|
|
|
std::sync::mpsc::sync_channel::<Result<WorkerChannelsExternal, ErrBox>>(1);
|
|
|
|
|
|
|
|
// TODO(bartlomieju): Isn't this wrong?
|
|
|
|
let result = ModuleSpecifier::resolve_url_or_path(&specifier)?;
|
|
|
|
let module_specifier = if !has_source_code {
|
|
|
|
ModuleSpecifier::resolve_import(&specifier, &referrer)?
|
|
|
|
} else {
|
|
|
|
result
|
|
|
|
};
|
2020-02-03 18:08:44 -05:00
|
|
|
|
|
|
|
std::thread::spawn(move || {
|
2020-02-08 20:34:31 +01:00
|
|
|
let result = State::new_for_worker(
|
|
|
|
global_state,
|
|
|
|
Some(child_permissions), // by default share with parent
|
2020-02-04 20:24:33 +01:00
|
|
|
module_specifier.clone(),
|
2020-02-03 18:08:44 -05:00
|
|
|
);
|
|
|
|
if let Err(err) = result {
|
2020-02-08 20:34:31 +01:00
|
|
|
handle_sender.send(Err(err)).unwrap();
|
2020-02-03 18:08:44 -05:00
|
|
|
return;
|
|
|
|
}
|
|
|
|
let child_state = result.unwrap();
|
|
|
|
let worker_name = args_name.unwrap_or_else(|| {
|
|
|
|
// TODO(bartlomieju): change it to something more descriptive
|
|
|
|
format!("USER-WORKER-{}", specifier)
|
|
|
|
});
|
|
|
|
|
|
|
|
// TODO: add a new option to make child worker not sharing permissions
|
|
|
|
// with parent (aka .clone(), requests from child won't reflect in parent)
|
|
|
|
let mut worker = WebWorker::new(
|
|
|
|
worker_name.to_string(),
|
|
|
|
startup_data::deno_isolate_init(),
|
|
|
|
child_state,
|
|
|
|
);
|
|
|
|
let script = format!("bootstrapWorkerRuntime(\"{}\")", worker_name);
|
|
|
|
js_check(worker.execute(&script));
|
|
|
|
js_check(worker.execute("runWorkerMessageLoop()"));
|
|
|
|
|
2020-02-08 20:34:31 +01:00
|
|
|
handle_sender.send(Ok(worker.thread_safe_handle())).unwrap();
|
2020-02-03 18:08:44 -05:00
|
|
|
|
|
|
|
// Has provided source code, execute immediately.
|
|
|
|
if has_source_code {
|
|
|
|
js_check(worker.execute(&source_code));
|
2020-02-08 20:34:31 +01:00
|
|
|
// FIXME(bartlomieju): runtime is not run in this case
|
2020-02-03 18:08:44 -05:00
|
|
|
return;
|
|
|
|
}
|
2020-01-29 18:54:23 +01:00
|
|
|
|
2020-02-03 18:08:44 -05:00
|
|
|
let fut = async move {
|
2020-02-05 17:16:07 -05:00
|
|
|
let r = worker
|
2020-02-03 18:08:44 -05:00
|
|
|
.execute_mod_async(&module_specifier, None, false)
|
|
|
|
.await;
|
2020-02-05 17:16:07 -05:00
|
|
|
if r.is_ok() {
|
|
|
|
let _ = (&mut *worker).await;
|
|
|
|
}
|
2020-02-03 18:08:44 -05:00
|
|
|
}
|
|
|
|
.boxed_local();
|
2019-08-14 17:03:02 +02:00
|
|
|
|
2020-02-03 18:08:44 -05:00
|
|
|
crate::tokio_util::run_basic(fut);
|
|
|
|
});
|
2020-01-18 00:43:53 +01:00
|
|
|
|
2020-02-08 20:34:31 +01:00
|
|
|
let handle = handle_receiver.recv().unwrap()?;
|
|
|
|
let worker_id = parent_state.add_child_worker(handle);
|
2019-08-26 14:50:21 +02:00
|
|
|
|
2020-02-08 20:34:31 +01:00
|
|
|
Ok(JsonOp::Sync(json!({ "id": worker_id })))
|
2019-11-17 01:17:47 +01:00
|
|
|
}
|
|
|
|
|
2019-08-26 14:50:21 +02:00
|
|
|
#[derive(Deserialize)]
|
2020-01-18 00:43:53 +01:00
|
|
|
struct WorkerArgs {
|
2019-11-09 21:07:14 +01:00
|
|
|
id: i32,
|
2019-08-14 17:03:02 +02:00
|
|
|
}
|
|
|
|
|
2020-01-18 00:43:53 +01:00
|
|
|
fn op_host_close_worker(
|
2020-02-08 20:34:31 +01:00
|
|
|
state: &State,
|
2020-01-18 00:43:53 +01:00
|
|
|
args: Value,
|
2020-01-24 15:10:49 -05:00
|
|
|
_data: Option<ZeroCopyBuf>,
|
2020-01-18 00:43:53 +01:00
|
|
|
) -> Result<JsonOp, ErrBox> {
|
|
|
|
let args: WorkerArgs = serde_json::from_value(args)?;
|
|
|
|
let id = args.id as u32;
|
2020-02-08 20:34:31 +01:00
|
|
|
let mut state = state.borrow_mut();
|
2020-01-18 00:43:53 +01:00
|
|
|
|
2020-02-08 20:34:31 +01:00
|
|
|
let maybe_worker_handle = state.workers.remove(&id);
|
2020-02-03 18:08:44 -05:00
|
|
|
if let Some(worker_handle) = maybe_worker_handle {
|
|
|
|
let mut sender = worker_handle.sender.clone();
|
2020-01-21 17:50:06 +01:00
|
|
|
sender.close_channel();
|
|
|
|
|
2020-02-03 18:08:44 -05:00
|
|
|
let mut receiver =
|
|
|
|
futures::executor::block_on(worker_handle.receiver.lock());
|
2020-01-21 17:50:06 +01:00
|
|
|
receiver.close();
|
2020-01-18 00:43:53 +01:00
|
|
|
};
|
|
|
|
|
|
|
|
Ok(JsonOp::Sync(json!({})))
|
|
|
|
}
|
|
|
|
|
2019-08-26 14:50:21 +02:00
|
|
|
#[derive(Deserialize)]
|
|
|
|
struct HostGetMessageArgs {
|
2019-11-09 21:07:14 +01:00
|
|
|
id: i32,
|
2019-08-14 17:03:02 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
/// Get message from guest worker as host
|
2019-10-11 11:41:54 -07:00
|
|
|
fn op_host_get_message(
|
2020-02-08 20:34:31 +01:00
|
|
|
state: &State,
|
2019-08-26 14:50:21 +02:00
|
|
|
args: Value,
|
2020-01-24 15:10:49 -05:00
|
|
|
_data: Option<ZeroCopyBuf>,
|
2019-08-26 14:50:21 +02:00
|
|
|
) -> Result<JsonOp, ErrBox> {
|
|
|
|
let args: HostGetMessageArgs = serde_json::from_value(args)?;
|
2019-11-09 21:07:14 +01:00
|
|
|
let id = args.id as u32;
|
2020-02-08 20:34:31 +01:00
|
|
|
|
|
|
|
let state = state.borrow();
|
2019-11-09 21:07:14 +01:00
|
|
|
// TODO: don't return bad resource anymore
|
2020-02-08 20:34:31 +01:00
|
|
|
let worker_handle = state.workers.get(&id).ok_or_else(bad_resource)?;
|
2020-02-03 18:08:44 -05:00
|
|
|
let fut = worker_handle.get_message();
|
2020-01-04 15:50:52 +05:30
|
|
|
let op = async move {
|
2020-02-05 17:16:07 -05:00
|
|
|
let maybe_buf = fut.await;
|
2020-01-04 15:50:52 +05:30
|
|
|
Ok(json!({ "data": maybe_buf }))
|
|
|
|
};
|
2020-02-03 18:08:44 -05:00
|
|
|
Ok(JsonOp::Async(op.boxed_local()))
|
2019-08-26 14:50:21 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
#[derive(Deserialize)]
|
|
|
|
struct HostPostMessageArgs {
|
2019-11-09 21:07:14 +01:00
|
|
|
id: i32,
|
2019-08-14 17:03:02 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
/// Post message to guest worker as host
|
2019-10-11 11:41:54 -07:00
|
|
|
fn op_host_post_message(
|
2020-02-08 20:34:31 +01:00
|
|
|
state: &State,
|
2019-08-26 14:50:21 +02:00
|
|
|
args: Value,
|
2020-01-24 15:10:49 -05:00
|
|
|
data: Option<ZeroCopyBuf>,
|
2019-08-26 14:50:21 +02:00
|
|
|
) -> Result<JsonOp, ErrBox> {
|
|
|
|
let args: HostPostMessageArgs = serde_json::from_value(args)?;
|
2019-11-09 21:07:14 +01:00
|
|
|
let id = args.id as u32;
|
|
|
|
let msg = Vec::from(data.unwrap().as_ref()).into_boxed_slice();
|
|
|
|
|
|
|
|
debug!("post message to worker {}", id);
|
2020-02-08 20:34:31 +01:00
|
|
|
let state = state.borrow();
|
2019-11-09 21:07:14 +01:00
|
|
|
// TODO: don't return bad resource anymore
|
2020-02-08 20:34:31 +01:00
|
|
|
let worker_handle = state.workers.get(&id).ok_or_else(bad_resource)?;
|
2020-02-03 18:08:44 -05:00
|
|
|
let fut = worker_handle
|
2019-11-17 14:14:50 +01:00
|
|
|
.post_message(msg)
|
2019-11-20 01:17:05 +01:00
|
|
|
.map_err(|e| DenoError::new(ErrorKind::Other, e.to_string()));
|
|
|
|
futures::executor::block_on(fut)?;
|
2019-08-26 14:50:21 +02:00
|
|
|
Ok(JsonOp::Sync(json!({})))
|
2019-08-14 17:03:02 +02:00
|
|
|
}
|
2019-10-11 11:41:54 -07:00
|
|
|
|
|
|
|
fn op_metrics(
|
2020-02-08 20:34:31 +01:00
|
|
|
state: &State,
|
2019-10-11 11:41:54 -07:00
|
|
|
_args: Value,
|
2020-01-24 15:10:49 -05:00
|
|
|
_zero_copy: Option<ZeroCopyBuf>,
|
2019-10-11 11:41:54 -07:00
|
|
|
) -> Result<JsonOp, ErrBox> {
|
2020-02-08 20:34:31 +01:00
|
|
|
let state = state.borrow();
|
2019-10-11 11:41:54 -07:00
|
|
|
let m = &state.metrics;
|
|
|
|
|
|
|
|
Ok(JsonOp::Sync(json!({
|
|
|
|
"opsDispatched": m.ops_dispatched.load(Ordering::SeqCst) as u64,
|
|
|
|
"opsCompleted": m.ops_completed.load(Ordering::SeqCst) as u64,
|
|
|
|
"bytesSentControl": m.bytes_sent_control.load(Ordering::SeqCst) as u64,
|
|
|
|
"bytesSentData": m.bytes_sent_data.load(Ordering::SeqCst) as u64,
|
|
|
|
"bytesReceived": m.bytes_received.load(Ordering::SeqCst) as u64
|
|
|
|
})))
|
|
|
|
}
|