diff --git a/crates/base/src/deno_runtime.rs b/crates/base/src/deno_runtime.rs index 61c3c9e4b..971732c35 100644 --- a/crates/base/src/deno_runtime.rs +++ b/crates/base/src/deno_runtime.rs @@ -35,6 +35,7 @@ use std::ffi::c_void; use std::fmt; use std::future::Future; use std::marker::PhantomData; +use std::mem::ManuallyDrop; use std::str::FromStr; use std::sync::{Arc, RwLock}; use std::task::Poll; @@ -197,8 +198,8 @@ impl GetRuntimeContext for () { } pub struct DenoRuntime { + pub js_runtime: ManuallyDrop, pub drop_token: CancellationToken, - pub js_runtime: JsRuntime, pub env_vars: HashMap, // TODO: does this need to be pub? pub conf: WorkerRuntimeOpts, @@ -218,14 +219,18 @@ pub struct DenoRuntime { impl Drop for DenoRuntime { fn drop(&mut self) { - self.drop_token.cancel(); - if self.conf.is_user_worker() { self.js_runtime.v8_isolate().remove_gc_prologue_callback( mem_check_gc_prologue_callback_fn, Arc::as_ptr(&self.mem_check) as *mut _, ); } + + unsafe { + ManuallyDrop::drop(&mut self.js_runtime); + } + + self.drop_token.cancel(); } } @@ -250,12 +255,11 @@ where maybe_inspector: Option, ) -> Result { let WorkerContextInitOpts { + mut conf, service_path, no_module_cache, import_map_path, env_vars, - events_rx, - conf, maybe_eszip, maybe_entrypoint, maybe_decorator, @@ -536,7 +540,7 @@ where ..Default::default() }; - let mut js_runtime = JsRuntime::new(runtime_options); + let mut js_runtime = ManuallyDrop::new(JsRuntime::new(runtime_options)); let version: Option<&str> = option_env!("GIT_V_TAG"); { @@ -610,10 +614,10 @@ where let mut env_vars = env_vars.clone(); - if conf.is_events_worker() { - // if worker is an events worker, assert events_rx is to be available - op_state - .put::>(events_rx.unwrap()); + if let Some(opts) = conf.as_events_worker_mut() { + op_state.put::>( + opts.events_msg_rx.take().unwrap(), + ); } if conf.is_main_worker() || conf.is_user_worker() { @@ -1215,7 +1219,6 @@ mod test { static_patterns, maybe_jsx_import_source_config: jsx_import_source_config, - events_rx: None, timing: None, import_map_path: None, @@ -1290,7 +1293,6 @@ mod test { no_module_cache: false, import_map_path: None, env_vars: Default::default(), - events_rx: None, timing: None, maybe_eszip: None, maybe_entrypoint: None, @@ -1337,7 +1339,6 @@ mod test { no_module_cache: false, import_map_path: None, env_vars: Default::default(), - events_rx: None, timing: None, maybe_eszip: Some(EszipPayloadKind::VecKind(eszip_code)), maybe_entrypoint: None, @@ -1401,7 +1402,6 @@ mod test { no_module_cache: false, import_map_path: None, env_vars: Default::default(), - events_rx: None, timing: None, maybe_eszip: Some(EszipPayloadKind::VecKind(eszip_code)), maybe_entrypoint: None, diff --git a/crates/base/src/rt_worker/worker.rs b/crates/base/src/rt_worker/worker.rs index 46f0f33a9..ef5915829 100644 --- a/crates/base/src/rt_worker/worker.rs +++ b/crates/base/src/rt_worker/worker.rs @@ -18,6 +18,7 @@ use sb_workers::context::{UserWorkerMsgs, WorkerContextInitOpts, WorkerExit, Wor use std::any::Any; use std::future::{pending, Future}; use std::pin::Pin; +use std::time::Duration; use tokio::io; use tokio::sync::mpsc::{unbounded_channel, UnboundedReceiver, UnboundedSender}; use tokio::sync::oneshot::{self, Receiver, Sender}; @@ -111,7 +112,6 @@ impl Worker { let method_cloner = self.clone(); let timing = opts.timing.take(); let worker_kind = opts.conf.to_worker_kind(); - let maybe_main_worker_opts = opts.conf.as_main_worker().cloned(); let cancel = self.cancel.clone(); let rt = if worker_kind.is_user_worker() { @@ -130,10 +130,8 @@ impl Worker { let permit = DenoRuntime::acquire().await; let result = match DenoRuntime::new(opts, inspector).await { Ok(new_runtime) => { - let mut runtime = scopeguard::guard(new_runtime, |mut runtime| { - unsafe { - runtime.js_runtime.v8_isolate().enter(); - } + let mut runtime = scopeguard::guard(new_runtime, |mut runtime| unsafe { + runtime.js_runtime.v8_isolate().enter(); }); unsafe { @@ -143,12 +141,10 @@ impl Worker { drop(permit); let metric_src = { - let js_runtime = &mut runtime.js_runtime; - let metric_src = WorkerMetricSource::from_js_runtime(js_runtime); - - if worker_kind.is_main_worker() { - let opts = maybe_main_worker_opts.unwrap(); - let state = js_runtime.op_state(); + let metric_src = + WorkerMetricSource::from_js_runtime(&mut runtime.js_runtime); + if let Some(opts) = runtime.conf.as_main_worker().cloned() { + let state = runtime.js_runtime.op_state(); let mut state_mut = state.borrow_mut(); let metric_src = RuntimeMetricSource::new( metric_src.clone(), @@ -173,7 +169,6 @@ impl Worker { let _cpu_timer; let mut supervise_cancel_token = None; - // TODO: Allow customization of supervisor let termination_fut = if worker_kind.is_user_worker() { // cputimer is returned from supervisor and assigned here to keep it in scope. let Ok((maybe_timer, cancel_token)) = create_supervisor( @@ -207,27 +202,50 @@ impl Worker { ) }; + let maybe_event_worker_ctx = + runtime.conf.as_events_worker().map(|it| { + Duration::from_secs( + it.event_worker_exit_deadline_sec.unwrap_or(10), + ) + }); + base_rt::SUPERVISOR_RT .spawn(async move { token.inbound.cancelled().await; - termination_request_token.cancel(); - - let data_ptr_mut = - Box::into_raw(Box::new(supervisor::IsolateInterruptData { - should_terminate: true, - isolate_memory_usage_tx: None, - })); - - if !thread_safe_handle.request_interrupt( - supervisor::handle_interrupt, - data_ptr_mut as *mut std::ffi::c_void, - ) { - drop(unsafe { Box::from_raw(data_ptr_mut) }); + + let mut already_terminated = false; + if let Some(dur) = maybe_event_worker_ctx { + already_terminated = tokio::time::timeout(dur, async { + while !is_terminated.is_raised() { + waker.wake(); + tokio::task::yield_now().await; + } + }) + .await + .is_ok(); } - while !is_terminated.is_raised() { - waker.wake(); - tokio::task::yield_now().await; + if !already_terminated { + termination_request_token.cancel(); + + let data_ptr_mut = Box::into_raw(Box::new( + supervisor::IsolateInterruptData { + should_terminate: true, + isolate_memory_usage_tx: None, + }, + )); + + if !thread_safe_handle.request_interrupt( + supervisor::handle_interrupt, + data_ptr_mut as *mut std::ffi::c_void, + ) { + drop(unsafe { Box::from_raw(data_ptr_mut) }); + } + + while !is_terminated.is_raised() { + waker.wake(); + tokio::task::yield_now().await; + } } let _ = termination_event_tx.send(WorkerEvents::Shutdown( @@ -251,7 +269,9 @@ impl Worker { let _guard = scopeguard::guard((), |_| { worker_key.and_then(|worker_key_unwrapped| { pool_msg_tx.map(|tx| { - if let Err(err) = tx.send(UserWorkerMsgs::Shutdown(worker_key_unwrapped)) { + if let Err(err) = + tx.send(UserWorkerMsgs::Shutdown(worker_key_unwrapped)) + { error!( "failed to send the shutdown signal to user worker pool: {:?}", err @@ -269,7 +289,6 @@ impl Worker { } }); - let result = method_cloner .handle_creation( &mut runtime, @@ -284,10 +303,10 @@ impl Worker { Ok(WorkerEvents::UncaughtException(ev)) => Some(ev.clone()), Err(err) => Some(UncaughtExceptionEvent { cpu_time_used: 0, - exception: err.to_string() + exception: err.to_string(), }), - _ => None + _ => None, }; if let Some(ev) = maybe_uncaught_exception_event { @@ -342,7 +361,7 @@ impl Worker { }; send_event_if_event_worker_available( - events_msg_tx.clone(), + events_msg_tx.as_ref(), event, event_metadata.clone(), ); diff --git a/crates/base/src/rt_worker/worker_ctx.rs b/crates/base/src/rt_worker/worker_ctx.rs index 685af6fd0..a0b048904 100644 --- a/crates/base/src/rt_worker/worker_ctx.rs +++ b/crates/base/src/rt_worker/worker_ctx.rs @@ -1,5 +1,6 @@ use crate::deno_runtime::DenoRuntime; use crate::inspector_server::Inspector; +use crate::server::ServerFlags; use crate::timeout::{self, CancelOnWriteTimeout, ReadTimeoutStream}; use crate::utils::send_event_if_event_worker_available; @@ -75,7 +76,7 @@ impl TerminationToken { pub fn child_token(&self) -> Self { Self { inbound: self.inbound.child_token(), - outbound: self.outbound.clone(), + outbound: self.outbound.child_token(), } } @@ -648,7 +649,7 @@ pub async fn create_worker>( .as_millis(); send_event_if_event_worker_available( - worker_struct_ref.events_msg_tx.clone(), + worker_struct_ref.events_msg_tx.as_ref(), WorkerEvents::Boot(BootEvent { boot_time: elapsed as usize, }), @@ -749,7 +750,6 @@ pub async fn create_main_worker( service_path, import_map_path, no_module_cache, - events_rx: None, timing: None, maybe_eszip, maybe_entrypoint, @@ -772,9 +772,9 @@ pub async fn create_main_worker( } pub async fn create_events_worker( + flags: &ServerFlags, events_worker_path: PathBuf, import_map_path: Option, - no_module_cache: bool, maybe_entrypoint: Option, maybe_decorator: Option, termination_token: Option, @@ -796,16 +796,18 @@ pub async fn create_events_worker( ( WorkerContextInitOpts { service_path, - no_module_cache, + no_module_cache: flags.no_module_cache, import_map_path, env_vars: std::env::vars().collect(), - events_rx: Some(events_rx), timing: None, maybe_eszip, maybe_entrypoint, maybe_decorator, maybe_module_code: None, - conf: WorkerRuntimeOpts::EventsWorker(EventWorkerRuntimeOpts {}), + conf: WorkerRuntimeOpts::EventsWorker(EventWorkerRuntimeOpts { + events_msg_rx: Some(events_rx), + event_worker_exit_deadline_sec: Some(flags.event_worker_exit_deadline_sec), + }), static_patterns: vec![], maybe_jsx_import_source_config: None, }, @@ -885,7 +887,7 @@ pub async fn create_user_worker_pool( } }, ..worker_options - }, tx, termination_token.as_ref().map(|it| it.child_token())); + }, tx, token.map(TerminationToken::child_token)); } Some(UserWorkerMsgs::Created(key, profile)) => { @@ -916,6 +918,8 @@ pub async fn create_user_worker_pool( } } + worker_pool.worker_event_sender.take(); + Ok(()) } }); diff --git a/crates/base/src/rt_worker/worker_pool.rs b/crates/base/src/rt_worker/worker_pool.rs index 2000fcc1e..a077f44ea 100644 --- a/crates/base/src/rt_worker/worker_pool.rs +++ b/crates/base/src/rt_worker/worker_pool.rs @@ -369,7 +369,6 @@ impl WorkerPool { no_module_cache, import_map_path, env_vars, - events_rx: None, timing: None, conf, maybe_eszip, diff --git a/crates/base/src/server.rs b/crates/base/src/server.rs index 00da7cdb1..f1e6e5c78 100644 --- a/crates/base/src/server.rs +++ b/crates/base/src/server.rs @@ -6,7 +6,6 @@ use crate::rt_worker::worker_pool::WorkerPoolPolicy; use crate::InspectorOption; use anyhow::{anyhow, bail, Context, Error}; use deno_config::JsxImportSourceConfig; -use event_worker::events::WorkerEventWithMetadata; use futures_util::future::{poll_fn, BoxFuture}; use futures_util::{FutureExt, Stream}; use hyper_v014::{server::conn::Http, service::Service, Body, Request, Response}; @@ -105,13 +104,13 @@ impl TerminationTokens { } async fn terminate(&self) { + self.pool.cancel_and_wait().await; + self.main.cancel_and_wait().await; + if let Some(token) = self.event.as_ref() { token.cancel_and_wait().await; } - self.pool.cancel_and_wait().await; - self.main.cancel_and_wait().await; - if let Some(token) = self.input.as_ref() { assert!(token.inbound.is_cancelled()); @@ -246,6 +245,7 @@ pub struct ServerFlags { pub tcp_nodelay: bool, pub graceful_exit_deadline_sec: u64, pub graceful_exit_keepalive_deadline_ms: Option, + pub event_worker_exit_deadline_sec: u64, pub request_wait_timeout_ms: Option, pub request_idle_timeout_ms: Option, pub request_read_timeout_ms: Option, @@ -351,7 +351,8 @@ impl Server { jsx_specifier: Option, jsx_module: Option, ) -> Result { - let mut worker_events_tx: Option> = None; + let mut worker_events_tx = None; + let maybe_events_entrypoint = entrypoints.events; let maybe_main_entrypoint = entrypoints.main; let termination_tokens = @@ -363,9 +364,9 @@ impl Server { let events_path_buf = events_path.to_path_buf(); let (ctx, sender) = create_events_worker( + &flags, events_path_buf, import_map_path.clone(), - flags.no_module_cache, maybe_events_entrypoint, maybe_decorator, Some(termination_tokens.event.clone().unwrap()), @@ -457,7 +458,6 @@ impl Server { let metric_src = self.metric_src.clone(); let termination_tokens = &self.termination_tokens; let input_termination_token = termination_tokens.input.as_ref(); - let flags = self.flags; let mut can_receive_event = false; let mut interrupted = false; @@ -488,7 +488,7 @@ impl Server { mut graceful_exit_deadline_sec, mut graceful_exit_keepalive_deadline_ms, .. - } = flags; + } = self.flags; let request_read_timeout_dur = request_read_timeout_ms.map(Duration::from_millis); let mut terminate_signal_fut = get_termination_signal(); diff --git a/crates/base/src/utils.rs b/crates/base/src/utils.rs index 908cf8813..8eec09a31 100644 --- a/crates/base/src/utils.rs +++ b/crates/base/src/utils.rs @@ -4,7 +4,7 @@ use tokio::sync::mpsc; pub mod units; pub fn send_event_if_event_worker_available( - maybe_event_worker: Option>, + maybe_event_worker: Option<&mpsc::UnboundedSender>, event: WorkerEvents, metadata: EventMetadata, ) { diff --git a/crates/base/src/utils/integration_test_helper.rs b/crates/base/src/utils/integration_test_helper.rs index 423b61848..226c1c05c 100644 --- a/crates/base/src/utils/integration_test_helper.rs +++ b/crates/base/src/utils/integration_test_helper.rs @@ -230,7 +230,6 @@ impl TestBedBuilder { no_module_cache: false, import_map_path: None, env_vars: HashMap::new(), - events_rx: None, timing: None, maybe_eszip: None, maybe_entrypoint: None, diff --git a/crates/base/tests/integration_tests.rs b/crates/base/tests/integration_tests.rs index 99230be01..4427ea457 100644 --- a/crates/base/tests/integration_tests.rs +++ b/crates/base/tests/integration_tests.rs @@ -198,7 +198,6 @@ async fn test_not_trigger_pku_sigsegv_due_to_jit_compilation_non_cli() { no_module_cache: false, import_map_path: None, env_vars: HashMap::new(), - events_rx: None, timing: None, maybe_eszip: None, maybe_entrypoint: None, @@ -357,7 +356,6 @@ async fn test_main_worker_boot_error() { no_module_cache: false, import_map_path: Some("./non-existing-import-map.json".to_string()), env_vars: HashMap::new(), - events_rx: None, timing: None, maybe_eszip: None, maybe_entrypoint: None, @@ -483,7 +481,6 @@ async fn test_main_worker_user_worker_mod_evaluate_exception() { no_module_cache: false, import_map_path: None, env_vars: HashMap::new(), - events_rx: None, timing: None, maybe_eszip: None, maybe_entrypoint: None, @@ -866,7 +863,6 @@ async fn test_worker_boot_invalid_imports() { no_module_cache: false, import_map_path: None, env_vars: HashMap::new(), - events_rx: None, timing: None, maybe_eszip: None, maybe_entrypoint: None, @@ -894,7 +890,6 @@ async fn test_worker_boot_with_0_byte_eszip() { no_module_cache: false, import_map_path: None, env_vars: HashMap::new(), - events_rx: None, timing: None, maybe_eszip: Some(EszipPayloadKind::VecKind(vec![])), maybe_entrypoint: Some("file:///src/index.ts".to_string()), @@ -920,7 +915,6 @@ async fn test_worker_boot_with_invalid_entrypoint() { no_module_cache: false, import_map_path: None, env_vars: HashMap::new(), - events_rx: None, timing: None, maybe_eszip: None, maybe_entrypoint: Some("file:///meow/mmmmeeeow.ts".to_string()), diff --git a/crates/cli/src/flags.rs b/crates/cli/src/flags.rs index 02ba957a8..3a6b37e1c 100644 --- a/crates/cli/src/flags.rs +++ b/crates/cli/src/flags.rs @@ -130,6 +130,15 @@ fn get_start_command() -> Command { .default_value("15") .value_parser(value_parser!(u64).range(..u64::MAX)), ) + .arg( + arg!(--"event-worker-exit-timeout" [SECONDS]) + .help(concat!( + "Maximum time in seconds that can wait for the event worker before terminating ", + "forcibly. (graceful exit)" + )) + .default_value("10") + .value_parser(value_parser!(u64).range(..u64::MAX)) + ) .arg( arg!( --"experimental-graceful-exit-keepalive-deadline-ratio" diff --git a/crates/cli/src/main.rs b/crates/cli/src/main.rs index 6684c4e8d..f0c8c2659 100644 --- a/crates/cli/src/main.rs +++ b/crates/cli/src/main.rs @@ -139,6 +139,11 @@ fn main() -> Result<(), anyhow::Error> { } }); + let event_worker_exit_deadline_sec = sub_matches + .get_one::("event-worker-exit-timeout") + .cloned() + .unwrap_or(0); + let maybe_max_parallelism = sub_matches.get_one::("max-parallelism").cloned(); let maybe_request_wait_timeout = @@ -180,6 +185,7 @@ fn main() -> Result<(), anyhow::Error> { tcp_nodelay, graceful_exit_deadline_sec, graceful_exit_keepalive_deadline_ms, + event_worker_exit_deadline_sec, request_wait_timeout_ms: maybe_request_wait_timeout, request_idle_timeout_ms: maybe_request_idle_timeout, request_read_timeout_ms: maybe_request_read_timeout, diff --git a/crates/event_worker/lib.rs b/crates/event_worker/lib.rs index dcd7af520..2dc0ff771 100644 --- a/crates/event_worker/lib.rs +++ b/crates/event_worker/lib.rs @@ -28,7 +28,10 @@ async fn op_event_accept(state: Rc>) -> Result match data { Some(event) => Ok(RawEvent::Event(Box::new(event))), - None => Ok(RawEvent::Done), + None => { + op_state.waker.wake(); + Ok(RawEvent::Done) + } } } diff --git a/crates/sb_workers/context.rs b/crates/sb_workers/context.rs index 4252494b9..765c8574a 100644 --- a/crates/sb_workers/context.rs +++ b/crates/sb_workers/context.rs @@ -114,10 +114,13 @@ pub struct MainWorkerRuntimeOpts { pub event_worker_metric_src: Option, } -#[derive(Debug, Clone)] -pub struct EventWorkerRuntimeOpts {} +#[derive(Debug)] +pub struct EventWorkerRuntimeOpts { + pub events_msg_rx: Option>, + pub event_worker_exit_deadline_sec: Option, +} -#[derive(Debug, Clone, EnumAsInner)] +#[derive(Debug, EnumAsInner)] pub enum WorkerRuntimeOpts { UserWorker(UserWorkerRuntimeOpts), MainWorker(MainWorkerRuntimeOpts), @@ -190,7 +193,6 @@ pub struct WorkerContextInitOpts { pub no_module_cache: bool, pub import_map_path: Option, pub env_vars: HashMap, - pub events_rx: Option>, pub timing: Option, pub conf: WorkerRuntimeOpts, pub maybe_eszip: Option, diff --git a/crates/sb_workers/lib.rs b/crates/sb_workers/lib.rs index 732e08378..c72d52081 100644 --- a/crates/sb_workers/lib.rs +++ b/crates/sb_workers/lib.rs @@ -142,7 +142,6 @@ pub async fn op_user_worker_create( no_module_cache, import_map_path, env_vars: env_vars_map, - events_rx: None, timing: None, maybe_eszip: maybe_eszip.map(EszipPayloadKind::JsBufferKind), maybe_entrypoint, diff --git a/examples/event-manager/index.ts b/examples/event-manager/index.ts index 73964608b..e0f522336 100644 --- a/examples/event-manager/index.ts +++ b/examples/event-manager/index.ts @@ -17,3 +17,5 @@ for await (const data of eventManager) { } } } + +console.log('event manager exiting'); \ No newline at end of file