diff --git a/crates/ion/src/env.rs b/crates/ion/src/env.rs index b74778b..735cf90 100644 --- a/crates/ion/src/env.rs +++ b/crates/ion/src/env.rs @@ -1,6 +1,4 @@ -use std::cell::RefCell; use std::future::Future; -use std::rc::Rc; use std::sync::Arc; use flume::Sender; @@ -27,7 +25,6 @@ pub struct Env { pub(crate) context: sys::GlobalContext, pub(crate) background_task_manager: Arc, pub(crate) global_refs: RefCounter, - pub(crate) shutdown_requested: Rc>, pub(crate) tx: Sender, pub(crate) finalizer_registry: FinalizerRegistery, pub(crate) global_this: sys::GlobalThis, @@ -40,7 +37,6 @@ impl Env { context: sys::GlobalContext, background_task_manager: Arc, global_refs: RefCounter, - shutdown_requested: Rc>, tx: Sender, finalizer_registry: FinalizerRegistery, global_this: sys::GlobalThis, @@ -53,7 +49,6 @@ impl Env { background_task_manager, inner: std::ptr::null_mut(), global_refs, - shutdown_requested, finalizer_registry, tx, }); @@ -85,19 +80,6 @@ impl Env { pub fn dec_ref(&self) { self.global_refs.dec(); - let shutdown_requested = { - let shutdown_requested = self.shutdown_requested.borrow(); - *shutdown_requested - }; - - if self.global_refs.count() == 0 && shutdown_requested { - self.tx - .try_send(JsWorkerEvent::RequestContextShutdown { - id: self.realm_id, - resolve: None, - }) - .unwrap(); - } } pub fn ref_count(&self) -> usize { diff --git a/crates/ion/src/error.rs b/crates/ion/src/error.rs index ecab1cd..db565f7 100644 --- a/crates/ion/src/error.rs +++ b/crates/ion/src/error.rs @@ -13,6 +13,8 @@ pub enum Error { PlatformCommunicationError, PlatformInitializeError, PlatformDisposeError, + WorkerAlreadyShutdown, + ContextAlreadyShutdown, IsolateNotInitializedError, EventLoopNotInitializedError, WorkerInitializeError, diff --git a/crates/ion/src/js_context.rs b/crates/ion/src/js_context.rs index 5947dc4..d5f6c8d 100644 --- a/crates/ion/src/js_context.rs +++ b/crates/ion/src/js_context.rs @@ -1,3 +1,5 @@ +use std::sync::Arc; + use flume::Sender; use flume::bounded; @@ -5,13 +7,16 @@ use crate::Env; use crate::Error; use crate::JsUnknown; use crate::platform::worker::JsWorkerEvent; -use crate::utils::channel::oneshot; +use crate::platform::worker_handle_state::WorkerHandleState; +use crate::utils::complete_signal::CompleteSignal; /// This is a handle to a v8::Context #[derive(Debug, Clone)] pub struct JsContext { + pub(crate) worker_handle_state: Arc, pub(crate) id: usize, pub(crate) tx: Sender, + pub(crate) context_shutdown_sig: CompleteSignal, } impl JsContext { @@ -83,25 +88,64 @@ impl JsContext { let specifier = specifier.as_ref().to_string(); self.exec_blocking(move |env| env.import(specifier)) } -} -impl Drop for JsContext { - fn drop(&mut self) { - let (tx, rx) = oneshot(); + /// Wait for the context to complete all activity + pub fn join(self) -> crate::Result<()> { + if !self.worker_handle_state.worker_handle_active() { + return Err(crate::Error::WorkerAlreadyShutdown); + } + + self.worker_handle_state + .context_handle_set_status(&self.id, false); if self .tx - .send(JsWorkerEvent::RequestContextShutdown { - id: self.id, - resolve: Some(tx), + .send(JsWorkerEvent::ContextHandleDeactivated { + id: self.id.clone(), }) .is_err() { - panic!("Cannot drop JsContext 1") - }; + return Err(crate::Error::ContextAlreadyShutdown); + } + + self.context_shutdown_sig.wait(); + + Ok(()) + } + + /// Wait for the context to complete all activity + pub async fn join_async(self) -> crate::Result<()> { + if !self.worker_handle_state.worker_handle_active() { + return Err(crate::Error::WorkerAlreadyShutdown); + } - if rx.recv().is_err() { - panic!("Cannot drop JsContext 2") + self.worker_handle_state + .context_handle_set_status(&self.id, false); + + if self + .tx + .send(JsWorkerEvent::ContextHandleDeactivated { + id: self.id.clone(), + }) + .is_err() + { + return Err(crate::Error::ContextAlreadyShutdown); } + + self.context_shutdown_sig.wait_async().await; + + Ok(()) + } +} + +impl Drop for JsContext { + fn drop(&mut self) { + self.worker_handle_state + .context_handle_set_status(&self.id, false); + + drop( + self.tx + .try_send(JsWorkerEvent::ContextHandleDropped { id: self.id }), + ); } } diff --git a/crates/ion/src/js_runtime.rs b/crates/ion/src/js_runtime.rs index 25b1b64..503e600 100644 --- a/crates/ion/src/js_runtime.rs +++ b/crates/ion/src/js_runtime.rs @@ -14,6 +14,8 @@ use crate::JsWorker; use crate::JsWorkerOptions; use crate::platform::platform::HAS_INIT; use crate::platform::platform::PLATFORM; +use crate::platform::worker_handle_state::WorkerHandleState; +use crate::utils::complete_signal::CompleteSignal; static JS_RUNTIME: OnceLock>> = OnceLock::new(); @@ -173,12 +175,17 @@ impl JsRuntime { pub fn spawn_worker( &self, options: JsWorkerOptions, - ) -> crate::Result> { + ) -> crate::Result { + let worker_handle_state = Arc::new(WorkerHandleState::default()); + let worker_shutdown_sig = CompleteSignal::default(); + let (tx, rx) = bounded(1); if self .tx .send(PlatformEvent::SpawnWorker { + worker_shutdown_sig: worker_shutdown_sig.clone(), + worker_handle_state: Arc::clone(&worker_handle_state), extensions: options.extensions, transformers: options.transformers, resolvers: options.resolvers, @@ -193,7 +200,12 @@ impl JsRuntime { return Err(Error::WorkerInitializeError); }; - Ok(Arc::new(JsWorker::new(tx, handle))) + Ok(JsWorker::new( + worker_handle_state, + tx, + Arc::new(handle), + worker_shutdown_sig, + )) } } diff --git a/crates/ion/src/js_worker.rs b/crates/ion/src/js_worker.rs index 7272b76..26fde70 100644 --- a/crates/ion/src/js_worker.rs +++ b/crates/ion/src/js_worker.rs @@ -11,7 +11,8 @@ use crate::JsExtension; use crate::JsResolver; use crate::JsTransformer; use crate::platform::worker::JsWorkerEvent; -use crate::utils::channel::oneshot; +use crate::platform::worker_handle_state::WorkerHandleState; +use crate::utils::complete_signal::CompleteSignal; #[derive(Default)] pub struct JsWorkerOptions { @@ -30,25 +31,39 @@ pub struct JsWorkerOptions { /// to be used to execute JavaScript #[derive(Debug)] pub struct JsWorker { + worker_handle_state: Arc, tx: Sender, - handle: Mutex>>>, + handle: Arc>>>>, + worker_shutdown_sig: CompleteSignal, } impl JsWorker { pub(crate) fn new( + worker_handle_state: Arc, tx: Sender, - handle: Mutex>>>, + handle: Arc>>>>, + worker_shutdown_sig: CompleteSignal, ) -> Self { - JsWorker { tx, handle } + JsWorker { + tx, + handle, + worker_handle_state, + worker_shutdown_sig, + } } /// Create a handle to a v8::Context associated with this v8::Isolate - pub fn create_context(&self) -> crate::Result> { + pub fn create_context(&self) -> crate::Result { + let context_shutdown_sig = CompleteSignal::default(); + let (tx, rx) = bounded(1); if self .tx - .send(JsWorkerEvent::CreateContext { resolve: tx }) + .send(JsWorkerEvent::CreateContext { + resolve: tx, + context_shutdown_sig: context_shutdown_sig.clone(), + }) .is_err() { return Err(Error::WorkerInitializeError); @@ -58,7 +73,15 @@ impl JsWorker { return Err(Error::WorkerInitializeError); }; - Ok(Arc::new(JsContext { id, tx })) + self.worker_handle_state + .context_handle_set_status(&id, true); + + Ok(JsContext { + id, + tx, + worker_handle_state: Arc::clone(&self.worker_handle_state), + context_shutdown_sig, + }) } pub fn run_garbage_collection_for_testing(&self) -> crate::Result<()> { @@ -74,24 +97,36 @@ impl JsWorker { Ok(rx.recv()?) } -} -impl Drop for JsWorker { - fn drop(&mut self) { - let (tx, rx) = oneshot(); + /// Wait for all of the contexts within the worker to complete all activity + pub fn join(self) -> crate::Result<()> { + self.worker_handle_state.worker_handle_deactivate(); + self.tx + .send(JsWorkerEvent::WorkerHandleDeactivated) + .unwrap(); - if self - .tx - .send(JsWorkerEvent::RequestShutdown { resolve: tx }) - .is_err() - { - panic!("Cannot drop JsWorker 1"); + self.worker_shutdown_sig.wait(); + + let Ok(mut handle) = self.handle.lock() else { + panic!("Cannot drop JsWorker 3"); }; - if rx.recv().is_err() { - panic!("Cannot drop JsWorker 2"); + if let Some(handle) = handle.take() { + drop(handle.join().unwrap()); } + Ok(()) + } + + /// Wait for all of the contexts within the worker to complete all activity + pub async fn join_async(self) -> crate::Result<()> { + self.worker_handle_state.worker_handle_deactivate(); + self.tx + .send(JsWorkerEvent::WorkerHandleDeactivated) + .unwrap(); + + self.worker_shutdown_sig.wait_async().await; + let Ok(mut handle) = self.handle.lock() else { panic!("Cannot drop JsWorker 3"); }; @@ -99,5 +134,14 @@ impl Drop for JsWorker { if let Some(handle) = handle.take() { drop(handle.join().unwrap()); } + + Ok(()) + } +} + +impl Drop for JsWorker { + fn drop(&mut self) { + self.worker_handle_state.worker_handle_deactivate(); + drop(self.tx.try_send(JsWorkerEvent::WorkerHandleDropped)); } } diff --git a/crates/ion/src/platform/mod.rs b/crates/ion/src/platform/mod.rs index e0088a9..8d1dc62 100644 --- a/crates/ion/src/platform/mod.rs +++ b/crates/ion/src/platform/mod.rs @@ -9,5 +9,6 @@ mod realm; pub mod resolve; pub(crate) mod sys; pub(crate) mod worker; +pub(crate) mod worker_handle_state; pub(crate) use realm::*; diff --git a/crates/ion/src/platform/platform.rs b/crates/ion/src/platform/platform.rs index 347e598..cdacb63 100644 --- a/crates/ion/src/platform/platform.rs +++ b/crates/ion/src/platform/platform.rs @@ -16,6 +16,8 @@ use crate::JsTransformer; use crate::platform::background_worker::BackgroundTaskManager; use crate::platform::worker::JsWorkerEvent; use crate::platform::worker::start_js_worker_thread; +use crate::platform::worker_handle_state::WorkerHandleState; +use crate::utils::complete_signal::CompleteSignal; pub(crate) enum PlatformEvent { Init { @@ -26,6 +28,8 @@ pub(crate) enum PlatformEvent { transformers: Vec, }, SpawnWorker { + worker_shutdown_sig: CompleteSignal, + worker_handle_state: Arc, extensions: Vec, resolvers: Vec, transformers: Vec, @@ -96,6 +100,8 @@ pub(crate) static PLATFORM: LazyLock> = LazyLock::new(|| { } } PlatformEvent::SpawnWorker { + worker_shutdown_sig, + worker_handle_state, resolve, extensions: init_extensions, resolvers: init_resolvers, @@ -117,6 +123,8 @@ pub(crate) static PLATFORM: LazyLock> = LazyLock::new(|| { } let (tx, handle) = start_js_worker_thread( + worker_shutdown_sig, + worker_handle_state, background_task_manager.clone(), worker_extensions, worker_resolvers, diff --git a/crates/ion/src/platform/realm.rs b/crates/ion/src/platform/realm.rs index 8b51201..f638985 100644 --- a/crates/ion/src/platform/realm.rs +++ b/crates/ion/src/platform/realm.rs @@ -1,6 +1,4 @@ -use std::cell::RefCell; use std::collections::HashMap; -use std::rc::Rc; use std::sync::Arc; use flume::Sender; @@ -18,6 +16,7 @@ use crate::platform::sys; use crate::platform::worker::JsWorkerEvent; use crate::utils::RefCounter; use crate::utils::channel::oneshot; +use crate::utils::complete_signal::CompleteSignal; // Container that constructs a V8 context and preserves the internals until dropped pub struct JsRealm { @@ -31,10 +30,10 @@ pub struct JsRealm { /// Used to tell the Worker if there are any long-lived async tasks /// that should prevent the context from being shutdown pub(crate) global_refs: RefCounter, - pub(crate) shutdown_requested: Rc>, pub(crate) modules: ModuleMap, - pub(crate) tx: Sender, pub(crate) global_this: sys::GlobalThis, + pub(crate) context_shutdown_sig: CompleteSignal, + pub(crate) tx: Sender, } impl JsRealm { @@ -45,12 +44,12 @@ impl JsRealm { transformers: HashMap>, background_task_manager: Arc, tx: Sender, + context_shutdown_sig: CompleteSignal, ) -> Box { let context = sys::GlobalContext::new(unsafe { &mut *isolate }); let global_this = sys::GlobalThis::new(&context); let global_refs = RefCounter::new(0); - let shutdown_requested = Rc::new(RefCell::new(false)); let finalizer_registry = FinalizerRegistery::new(isolate); // TODO make these RefCells @@ -61,7 +60,6 @@ impl JsRealm { context.clone(), Arc::clone(&background_task_manager), global_refs.clone(), - Rc::clone(&shutdown_requested), tx.clone(), finalizer_registry.clone(), global_this.clone(), @@ -76,9 +74,9 @@ impl JsRealm { resolvers, transformers, global_refs, - shutdown_requested, finalizer_registry, global_this, + context_shutdown_sig, tx, }); @@ -133,7 +131,7 @@ impl JsRealm { ) -> crate::Result { let (tx, rx) = oneshot(); self.background_task_manager.spawn(async move { - tx.try_send(fut.await).unwrap(); + tx.try_send(fut.await).expect("Unable to resolve"); Ok(()) })?; rx.recv()? diff --git a/crates/ion/src/platform/worker.rs b/crates/ion/src/platform/worker.rs index 5f8a67a..d79b586 100644 --- a/crates/ion/src/platform/worker.rs +++ b/crates/ion/src/platform/worker.rs @@ -18,20 +18,16 @@ use crate::JsResolver; use crate::JsTransformer; use crate::fs::FileSystem; use crate::platform::background_worker::BackgroundTaskManager; +use crate::platform::worker_handle_state::WorkerHandleState; use crate::utils::HashMapExt; use crate::utils::PathExt; +use crate::utils::complete_signal::CompleteSignal; pub(crate) enum JsWorkerEvent { CreateContext { + context_shutdown_sig: CompleteSignal, resolve: Sender<(usize, Sender)>, }, - BackgroundTaskComplete { - id: usize, - }, - RequestContextShutdown { - resolve: Option>, - id: usize, - }, Exec { id: usize, #[allow(clippy::type_complexity)] @@ -42,17 +38,27 @@ pub(crate) enum JsWorkerEvent { id: usize, specifier: String, }, - RequestShutdown { - resolve: Sender<()>, - }, RunGarbageCollectionForTesting { resolve: Sender<()>, }, + WorkerHandleDropped, + WorkerHandleDeactivated, + ContextHandleDropped { + id: usize, + }, + ContextHandleDeactivated { + id: usize, + }, + BackgroundTaskComplete { + id: usize, + }, } // Create a dedicated thread to host the isolate #[allow(clippy::type_complexity)] pub(crate) fn start_js_worker_thread( + worker_shutdown_sig: CompleteSignal, + worker_handle_state: Arc, background_task_manager: Arc, extensions: Vec>, resolvers: Vec, @@ -68,6 +74,8 @@ pub(crate) fn start_js_worker_thread( let tx: Sender = tx.clone(); move || { worker_thread( + worker_shutdown_sig.clone(), + worker_handle_state, tx, rx, background_task_manager, @@ -82,6 +90,8 @@ pub(crate) fn start_js_worker_thread( } fn worker_thread( + worker_shutdown_sig: CompleteSignal, + worker_handle_state: Arc, tx: Sender, rx: Receiver, background_task_manager: Arc, @@ -98,15 +108,14 @@ fn worker_thread( // Maintain a store of Global to help with cleanup on shutdown. let mut realms = HashMap::>::new(); - // Cleanup hooks - let mut shutdown_context_senders = HashMap::>>::new(); - let mut shutdown_senders = Vec::>::new(); - let mut shutdown_requested = false; - while let Ok(event) = rx.recv() { - // println!("{:?} {:?}", active_context, event); + // println!(" {:?}", event); + match event { - JsWorkerEvent::CreateContext { resolve } => { + JsWorkerEvent::CreateContext { + resolve, + context_shutdown_sig, + } => { let realm = JsRealm::new( isolate_ptr, fs.clone(), @@ -114,6 +123,7 @@ fn worker_thread( transformers.clone(), background_task_manager.clone(), tx.clone(), + context_shutdown_sig, ); let realm_id = realm.id(); @@ -122,26 +132,25 @@ fn worker_thread( realms.insert(realm_id, realm); resolve.try_send((realm_id, tx.clone()))?; } - JsWorkerEvent::RequestContextShutdown { id, resolve } => { - // Store shutdown resolvers for when the context is closed - if let Some(resolve) = resolve { - shutdown_context_senders - .entry(id) - .or_default() - .push(resolve); - } + JsWorkerEvent::Exec { id, callback, span } => { + let realm = realms.try_get(&id)?; - // If there are async tasks pending then wait for them to complete - { - let realm = realms.try_get_mut(&id)?; - let mut realm_shutdown_requested = realm.shutdown_requested.borrow_mut(); - (*realm_shutdown_requested) = true; - if realm.global_refs.count() != 0 { - continue; - } + let _span_guard = span.enter(); + if let Err(err) = callback(realm.env()) { + // TODO global error handler + panic!("Callback errored {:?}", err) }; - // If there are no async tasks then shutdown the context + if worker_handle_state.worker_handle_active() + && worker_handle_state.context_handle_active(&id) + { + continue; + } + + if realm.global_refs.count() != 0 { + continue; + } + let Some(realm) = realms.remove(&id) else { continue; }; @@ -150,32 +159,30 @@ fn worker_thread( finalizer_registry.clear(); drop(finalizer_registry); - for resolver in shutdown_context_senders.remove(&id).unwrap_or_default() { - let _ = resolver.try_send(()); + realm.context_shutdown_sig.done(); + } + JsWorkerEvent::BackgroundTaskComplete { id } => { + let realm = realms.try_get(&id)?; + + if worker_handle_state.worker_handle_active() + && worker_handle_state.context_handle_active(&id) + { + continue; } - if shutdown_requested && realms.is_empty() { - for sender in shutdown_senders { - let _ = sender.try_send(()); - } - break; + if realm.global_refs.count() != 0 { + continue; } - } - JsWorkerEvent::Exec { id, callback, span } => { - let realm = realms.try_get(&id)?; - let _span_guard = span.enter(); - if let Err(err) = callback(realm.env()) { - // TODO global error handler - panic!("Callback errored {:?}", err) + let Some(realm) = realms.remove(&id) else { + continue; }; - } - JsWorkerEvent::BackgroundTaskComplete { id } => { - let realm = realms.try_get(&id)?; - let realm_shutdown_requested = realm.shutdown_requested.borrow(); - if *realm_shutdown_requested && realm.global_refs.count() == 0 { - tx.try_send(JsWorkerEvent::RequestContextShutdown { id, resolve: None })?; - } + + let finalizer_registry = realm.finalizer_registry; + finalizer_registry.clear(); + drop(finalizer_registry); + + realm.context_shutdown_sig.done(); } JsWorkerEvent::Import { id, specifier } => { Module::v8_initialize( @@ -185,45 +192,102 @@ fn worker_thread( std::env::current_dir()?.try_to_string()?, )?; } - JsWorkerEvent::RequestShutdown { resolve } => { - shutdown_senders.push(resolve); - shutdown_requested = true; - if !realms.is_empty() { - continue; - } + JsWorkerEvent::RunGarbageCollectionForTesting { resolve } => { + isolate.request_garbage_collection_for_testing(v8::GarbageCollectionType::Full); + resolve.try_send(())?; + } + JsWorkerEvent::WorkerHandleDropped => { + for id in realms.keys().cloned().collect::>() { + let Some(realm) = realms.remove(&id) else { + continue; + }; - for sender in shutdown_senders { - let _ = sender.try_send(()); + let finalizer_registry = realm.finalizer_registry; + finalizer_registry.clear(); + drop(finalizer_registry); + + realm.context_shutdown_sig.done(); } break; } - JsWorkerEvent::RunGarbageCollectionForTesting { resolve } => { - isolate.request_garbage_collection_for_testing(v8::GarbageCollectionType::Full); - resolve.try_send(())?; + JsWorkerEvent::WorkerHandleDeactivated => { + let mut to_drop = vec![]; + + for (id, realm) in realms.iter() { + if realm.global_refs.count() != 0 { + continue; + } + to_drop.push(id.clone()); + } + + for id in to_drop { + let Some(realm) = realms.remove(&id) else { + continue; + }; + + let finalizer_registry = realm.finalizer_registry; + finalizer_registry.clear(); + drop(finalizer_registry); + + realm.context_shutdown_sig.done(); + } + } + JsWorkerEvent::ContextHandleDeactivated { id } => { + if realms.try_get_mut(&id)?.global_refs.count() != 0 { + continue; + } + + let Some(realm) = realms.remove(&id) else { + continue; + }; + + let finalizer_registry = realm.finalizer_registry; + finalizer_registry.clear(); + drop(finalizer_registry); + + realm.context_shutdown_sig.done(); } + JsWorkerEvent::ContextHandleDropped { id } => { + let Some(realm) = realms.remove(&id) else { + continue; + }; + + let finalizer_registry = realm.finalizer_registry; + finalizer_registry.clear(); + drop(finalizer_registry); + + realm.context_shutdown_sig.done(); + } + } + + if !worker_handle_state.worker_handle_active() && realms.len() == 0 { + break; } } + worker_shutdown_sig.done(); + Ok(()) } #[allow(unused)] +#[rustfmt::skip] impl std::fmt::Debug for JsWorkerEvent { fn fmt( &self, f: &mut std::fmt::Formatter<'_>, ) -> std::fmt::Result { match self { - Self::CreateContext { resolve } => write!(f, "CreateContext"), - Self::BackgroundTaskComplete { id } => write!(f, "BackgroundTaskComplete"), - Self::RequestContextShutdown { id, resolve } => write!(f, "RequestContextShutdown"), - Self::Exec { id, callback, span } => write!(f, "Exec"), - Self::Import { id, specifier } => write!(f, "Import"), - Self::RequestShutdown { resolve } => write!(f, "RequestShutdown"), - Self::RunGarbageCollectionForTesting { resolve } => { - write!(f, "RunGarbageCollectionForTesting") - } + Self::CreateContext { resolve, context_shutdown_sig } => write!(f, "CreateContext"), + Self::Exec { id, callback, span } => write!(f, "Exec [id={}]", id), + Self::BackgroundTaskComplete { id } => write!(f, "BackgroundTaskComplete [id={}]", id), + Self::Import { id, specifier } => write!(f, "Import"), + Self::WorkerHandleDropped => write!(f, "WorkerHandleDropped"), + Self::WorkerHandleDeactivated => write!(f, "WorkerHandleDeactivated"), + Self::ContextHandleDropped { id } => write!(f, "ContextHandleDropped"), + Self::ContextHandleDeactivated { id } => write!(f, "ContextHandleDeactivated [id={}]", id), + Self::RunGarbageCollectionForTesting { resolve } => write!(f, "RunGarbageCollectionForTesting"), } } } diff --git a/crates/ion/src/platform/worker_handle_state.rs b/crates/ion/src/platform/worker_handle_state.rs new file mode 100644 index 0000000..caa7eca --- /dev/null +++ b/crates/ion/src/platform/worker_handle_state.rs @@ -0,0 +1,100 @@ +use std::collections::HashMap; +use std::sync::atomic::AtomicBool; +use std::sync::atomic::Ordering; + +use parking_lot::Mutex; +use parking_lot::RwLock; + +pub type WorkerShutdownCallback = Box; +pub type ContextShutdownCallback = Box; + +pub struct WorkerHandleState { + pub(crate) worker_handle_active: AtomicBool, + pub(crate) context_handle_active: RwLock>, + pub(crate) worker_shutdown: Mutex>, + pub(crate) context_shutdown: Mutex>>, +} + +impl Default for WorkerHandleState { + fn default() -> Self { + Self { + worker_handle_active: AtomicBool::new(true), + context_handle_active: Default::default(), + worker_shutdown: Default::default(), + context_shutdown: Default::default(), + } + } +} + +impl WorkerHandleState { + pub(crate) fn worker_handle_active(&self) -> bool { + self.worker_handle_active.load(Ordering::Relaxed) + } + + pub(crate) fn worker_handle_deactivate(&self) { + self.worker_handle_active.swap(false, Ordering::Relaxed); + } + + pub(crate) fn context_handle_active( + &self, + id: &usize, + ) -> bool { + self.context_handle_active + .read() + .get(id) + .unwrap_or(&false) + .clone() + } + + pub(crate) fn context_handle_set_status( + &self, + id: &usize, + status: bool, + ) { + self.context_handle_active + .write() + .insert(id.clone(), status); + } + + // pub(crate) fn add_worker_shutdown_callback( + // &self, + // callback: impl 'static + Send + Sync + FnOnce(), + // ) { + // self.worker_shutdown.lock().push(Box::new(callback)); + // } + + // pub(crate) fn add_context_shutdown_callback( + // &self, + // id: usize, + // callback: impl 'static + Send + Sync + FnOnce(), + // ) { + // self.context_shutdown + // .lock() + // .entry(id) + // .or_default() + // .push(Box::new(callback)); + // } + + // pub(crate) fn take_worker_shutdown_callbacks(&self) -> Vec { + // std::mem::take(&mut *self.worker_shutdown.lock()) + // } + + // pub(crate) fn take_context_shutdown_callbacks( + // &self, + // id: usize, + // ) -> Vec { + // std::mem::take(&mut *self.context_shutdown.lock().entry(id).or_default()) + // } +} + +impl std::fmt::Debug for WorkerHandleState { + fn fmt( + &self, + f: &mut std::fmt::Formatter<'_>, + ) -> std::fmt::Result { + f.debug_struct("CallbackRegistry") + .field("worker_shutdown", &self.worker_shutdown.lock().len()) + .field("context_shutdown", &self.context_shutdown.lock().len()) + .finish() + } +} diff --git a/crates/ion/src/utils/complete_signal.rs b/crates/ion/src/utils/complete_signal.rs new file mode 100644 index 0000000..d172743 --- /dev/null +++ b/crates/ion/src/utils/complete_signal.rs @@ -0,0 +1,127 @@ +use std::sync::Arc; +use std::sync::Condvar; +use std::sync::Mutex; + +use tokio::sync::Notify; + +#[derive(Clone, Default)] +pub struct CompleteSignal { + inner: Arc, +} + +impl std::fmt::Debug for CompleteSignal { + fn fmt( + &self, + f: &mut std::fmt::Formatter<'_>, + ) -> std::fmt::Result { + f.debug_struct("CompleteSignal").finish() + } +} + +struct Inner { + completed: Mutex, + condvar: Condvar, + notify: Notify, +} + +impl Default for Inner { + fn default() -> Self { + Self { + completed: Mutex::new(false), + condvar: Condvar::new(), + notify: Notify::new(), + } + } +} + +impl CompleteSignal { + pub fn new() -> Self { + Self::default() + } + + pub fn done(&self) { + let mut completed = self.inner.completed.lock().unwrap(); + if !*completed { + *completed = true; + self.inner.condvar.notify_all(); + self.inner.notify.notify_waiters(); + } + } + + pub fn wait(&self) { + let mut completed = self.inner.completed.lock().unwrap(); + while !*completed { + completed = self.inner.condvar.wait(completed).unwrap(); + } + } + + pub async fn wait_async(&self) { + { + let completed = self.inner.completed.lock().unwrap(); + if *completed { + return; + } + } + + let notified = self.inner.notify.notified(); + + { + let completed = self.inner.completed.lock().unwrap(); + if *completed { + return; + } + } + + notified.await; + } + + pub fn is_done(&self) -> bool { + *self.inner.completed.lock().unwrap() + } +} + +#[cfg(test)] +mod tests { + use std::thread; + use std::time::Duration; + + use super::*; + + #[tokio::test] + async fn test_complete_signal() { + let sig = CompleteSignal::default(); + + thread::spawn({ + let sig = sig.clone(); + move || { + println!("** Controller WAITING"); + thread::sleep(Duration::from_millis(100)); + println!("** Controller DONE"); + sig.done(); + } + }); + + thread::spawn({ + let sig = sig.clone(); + move || { + println!("Sync WAITING"); + sig.wait(); + println!("Sync DONE 1"); + sig.wait(); + sig.wait(); + println!("Sync DONE 2"); + } + }); + + println!("Async WAITING"); + sig.wait_async().await; + println!("Async DONE 1"); + + // Subsequent calls after the signal is complete will complete immediately + sig.wait_async().await; + sig.wait_async().await; + println!("Async DONE 2"); + + assert!(sig.is_done()); + } +} diff --git a/crates/ion/src/utils/mod.rs b/crates/ion/src/utils/mod.rs index 54b4dcd..6e4bd30 100644 --- a/crates/ion/src/utils/mod.rs +++ b/crates/ion/src/utils/mod.rs @@ -1,4 +1,5 @@ pub mod channel; +pub mod complete_signal; pub mod debug; pub mod hash; pub mod hash_map_ext; diff --git a/examples/src/_utils/mod.rs b/examples/src/_utils/mod.rs index b42be1c..ed514d2 100644 --- a/examples/src/_utils/mod.rs +++ b/examples/src/_utils/mod.rs @@ -1,2 +1,3 @@ pub mod memory_blob; pub mod memory_usage; +pub mod thread_id; diff --git a/examples/src/_utils/thread_id.rs b/examples/src/_utils/thread_id.rs new file mode 100644 index 0000000..9aca507 --- /dev/null +++ b/examples/src/_utils/thread_id.rs @@ -0,0 +1,12 @@ +use std::sync::atomic::AtomicUsize; +use std::sync::atomic::Ordering; + +static NEXT_THREAD_ID: AtomicUsize = AtomicUsize::new(1); + +thread_local! { + static THREAD_ID: usize = NEXT_THREAD_ID.fetch_add(1, Ordering::Relaxed); +} + +pub fn thread_id() -> usize { + THREAD_ID.with(|&id| id) +} diff --git a/examples/src/background_tasks/mod.rs b/examples/src/background_tasks/mod.rs index a83c75d..b40b7e3 100644 --- a/examples/src/background_tasks/mod.rs +++ b/examples/src/background_tasks/mod.rs @@ -54,5 +54,7 @@ pub fn main() -> anyhow::Result<()> { Ok(()) })?; + ctx.join()?; + Ok(()) } diff --git a/examples/src/basic/mod.rs b/examples/src/basic/mod.rs index 1a07431..4547388 100644 --- a/examples/src/basic/mod.rs +++ b/examples/src/basic/mod.rs @@ -22,5 +22,7 @@ pub fn main() -> anyhow::Result<()> { Ok(()) })?; + ctx.join()?; + Ok(()) } diff --git a/examples/src/basic_join/basic_join.test.ts b/examples/src/basic_join/basic_join.test.ts new file mode 100644 index 0000000..bde7120 --- /dev/null +++ b/examples/src/basic_join/basic_join.test.ts @@ -0,0 +1,537 @@ +import { executeExample } from "../../test-utils/run_test.ts"; +import { assert, assertEquals, assertObjectMatch } from "jsr:@std/assert@^1"; + +type Result = { + thread: number; + message: string; + js_context?: number; + event_loop?: boolean; +}; + +async function executeBasicJoin(caseName: string): Promise> { + const result = await executeExample("basic_join", [caseName]); + return result.split("\n").map((record) => JSON.parse(record)); +} + +type ProcessedRecords = { + main: Array; + jsContexts: Record>; + eventLoop: Record>; +}; + +function processResults(input: Array): ProcessedRecords { + const processed: ProcessedRecords = { + main: [], + jsContexts: {}, + eventLoop: {}, + }; + + for (const result of input) { + // No js_context -> goes to main + if (result.js_context === undefined) { + processed.main.push(result); + continue; + } + + // Has event_loop flag -> goes to eventLoop + if (result.event_loop) { + if (!processed.eventLoop[result.js_context]) { + processed.eventLoop[result.js_context] = []; + } + processed.eventLoop[result.js_context].push(result); + continue; + } + + // Has js_context but no event_loop -> goes to jsContexts + if (!processed.jsContexts[result.js_context]) { + processed.jsContexts[result.js_context] = []; + } + processed.jsContexts[result.js_context].push(result); + } + + return processed; +} + +async function run(caseName: string): Promise { + return processResults(await executeBasicJoin(caseName)); +} + +function assertArraysMatch, Y extends Array>( + arr1: T, + arr2: Y, + msg?: string +): void { + return assertObjectMatch({ arr: arr1 }, { arr: arr2 }, msg); +} + +Deno.test("should_cancel_when_dropped", async () => { + const example = "should_cancel_when_dropped"; + const results = await run(example); + + // The code on the main thread will always run + assertArraysMatch(results.main, [ + { thread: 1, message: "start" }, + { thread: 1, message: "end" }, + ]); + + // The code on the JavaScript thread may or may not run + assert( + (results.jsContexts[0] || []).length === 0 || + (results.jsContexts[0] || []).length === 2 + ); + + // The code on the Event Loop should not run + assertObjectMatch(results.eventLoop, {}); +}); + +Deno.test("should_cancel_when_dropped_multiple", async () => { + const example = "should_cancel_when_dropped_multiple"; + const results = await run(example); + + // The code on the main thread will always run + assertArraysMatch(results.main, [ + { thread: 1, message: "start" }, + { thread: 1, message: "end" }, + ]); + + // The code on the JavaScript thread may or may not run + assert( + (results.jsContexts[0] || []).length === 0 || + (results.jsContexts[0] || []).length === 4 + ); + + // The code on the Event Loop should not run + assertObjectMatch(results.eventLoop, {}); +}); + +Deno.test("should_cancel_blocking_when_dropped", async () => { + const example = "should_cancel_blocking_when_dropped"; + const results = await run(example); + + assertArraysMatch(results.main, [ + { thread: 1, message: "start" }, + { thread: 1, message: "end" }, + ]); + + assertObjectMatch(results.jsContexts, { + "0": [ + { thread: 2, js_context: 0, message: "start" }, + { thread: 2, js_context: 0, message: "end" }, + ], + }); + + // The code on the Event Loop may or may not run, but not progress when the task sleeps + assertObjectMatch(results.eventLoop, {}); +}); + +Deno.test("should_cancel_blocking_when_dropped_multiple", async () => { + const example = "should_cancel_blocking_when_dropped_multiple"; + const results = await run(example); + + assertArraysMatch(results.main, [ + { thread: 1, message: "start" }, + { thread: 1, message: "end" }, + ]); + + assertObjectMatch(results.jsContexts, { + "0": [ + { thread: 2, js_context: 0, message: "start" }, + { thread: 2, js_context: 0, message: "end" }, + { thread: 2, js_context: 0, message: "start" }, + { thread: 2, js_context: 0, message: "end" }, + ], + }); + + // The code on the Event Loop may or may not run, but not progress when the task sleeps + assertObjectMatch(results.eventLoop, {}); +}); + +Deno.test("should_wait_for_code_to_finish", async () => { + const example = "should_wait_for_code_to_finish"; + const results = await run(example); + assertObjectMatch(results, { + main: [ + { thread: 1, message: "start" }, + { thread: 1, message: "end" }, + ], + jsContexts: { + "0": [ + { thread: 2, js_context: 0, message: "start" }, + { thread: 2, js_context: 0, message: "end" }, + { thread: 2, js_context: 0, message: "resolved" }, + ], + }, + eventLoop: { + "0": [ + { + thread: 1000, + js_context: 0, + event_loop: true, + message: "start", + }, + { + thread: 1000, + js_context: 0, + event_loop: true, + message: "end", + }, + ], + }, + }); +}); + +Deno.test("should_wait_for_code_to_finish_multiple", async () => { + const example = "should_wait_for_code_to_finish_multiple"; + const results = await run(example); + assertObjectMatch(results, { + main: [ + { thread: 1, message: "start" }, + { thread: 1, message: "end" }, + ], + jsContexts: { + "0": [ + { thread: 2, js_context: 0, message: "start" }, + { thread: 2, js_context: 0, message: "end" }, + { thread: 2, js_context: 0, message: "start" }, + { thread: 2, js_context: 0, message: "end" }, + { thread: 2, js_context: 0, message: "resolved" }, + { thread: 2, js_context: 0, message: "resolved" }, + ], + }, + eventLoop: { + "0": [ + { + thread: 1000, + js_context: 0, + event_loop: true, + message: "start", + }, + { + thread: 1000, + js_context: 0, + event_loop: true, + message: "start", + }, + { + thread: 1000, + js_context: 0, + event_loop: true, + message: "end", + }, + { + thread: 1000, + js_context: 0, + event_loop: true, + message: "end", + }, + ], + }, + }); +}); + +Deno.test("should_wait_for_code_to_finish_blocking", async () => { + const example = "should_wait_for_code_to_finish_blocking"; + const results = await run(example); + assertObjectMatch(results, { + main: [ + { thread: 1, message: "start" }, + { thread: 1, message: "end" }, + ], + jsContexts: { + "0": [ + { thread: 2, js_context: 0, message: "start" }, + { thread: 2, js_context: 0, message: "end" }, + { thread: 2, js_context: 0, message: "resolved" }, + ], + }, + eventLoop: { + "0": [ + { + thread: 1000, + js_context: 0, + event_loop: true, + message: "start", + }, + { + thread: 1000, + js_context: 0, + event_loop: true, + message: "end", + }, + ], + }, + }); +}); + +Deno.test("should_wait_for_code_to_finish_worker", async () => { + const example = "should_wait_for_code_to_finish_worker"; + const results = await run(example); + assertObjectMatch(results, { + main: [ + { thread: 1, message: "start" }, + { thread: 1, message: "end" }, + ], + jsContexts: { + "0": [ + { thread: 2, js_context: 0, message: "start" }, + { thread: 2, js_context: 0, message: "end" }, + { thread: 2, js_context: 0, message: "resolved" }, + ], + }, + eventLoop: { + "0": [ + { + thread: 1000, + js_context: 0, + event_loop: true, + message: "start", + }, + { + thread: 1000, + js_context: 0, + event_loop: true, + message: "end", + }, + ], + }, + }); +}); + +Deno.test("should_wait_for_code_to_finish_worker_blocking", async () => { + const example = "should_wait_for_code_to_finish_worker_blocking"; + const results = await run(example); + assertObjectMatch(results, { + main: [ + { thread: 1, message: "start" }, + { thread: 1, message: "end" }, + ], + jsContexts: { + "0": [ + { thread: 2, js_context: 0, message: "start" }, + { thread: 2, js_context: 0, message: "end" }, + { thread: 2, js_context: 0, message: "resolved" }, + ], + }, + eventLoop: { + "0": [ + { + thread: 1000, + js_context: 0, + event_loop: true, + message: "start", + }, + { + thread: 1000, + js_context: 0, + event_loop: true, + message: "end", + }, + ], + }, + }); +}); + +Deno.test("should_wait_for_code_to_finish_context", async () => { + const example = "should_wait_for_code_to_finish_context"; + const results = await run(example); + assertObjectMatch(results, { + main: [ + { thread: 1, message: "start" }, + { thread: 1, message: "end" }, + ], + jsContexts: { + "0": [ + { thread: 2, js_context: 0, message: "start" }, + { thread: 2, js_context: 0, message: "end" }, + { thread: 2, js_context: 0, message: "resolved" }, + ], + }, + eventLoop: { + "0": [ + { + thread: 1000, + js_context: 0, + event_loop: true, + message: "start", + }, + { + thread: 1000, + js_context: 0, + event_loop: true, + message: "end", + }, + ], + }, + }); +}); + +Deno.test("should_wait_for_code_to_finish_context_blocking", async () => { + const example = "should_wait_for_code_to_finish_context_blocking"; + const results = await run(example); + assertObjectMatch(results, { + main: [ + { thread: 1, message: "start" }, + { thread: 1, message: "end" }, + ], + jsContexts: { + "0": [ + { thread: 2, js_context: 0, message: "start" }, + { thread: 2, js_context: 0, message: "end" }, + { thread: 2, js_context: 0, message: "resolved" }, + ], + }, + eventLoop: { + "0": [ + { + thread: 1000, + js_context: 0, + event_loop: true, + message: "start", + }, + { + thread: 1000, + js_context: 0, + event_loop: true, + message: "end", + }, + ], + }, + }); +}); + +Deno.test("should_wait_for_code_to_finish_contexts_blocking", async () => { + const example = "should_wait_for_code_to_finish_contexts_blocking"; + const results = processResults(await executeBasicJoin(example)); + assertObjectMatch(results, { + main: [ + { thread: 1, message: "start" }, + { thread: 1, message: "end" }, + ], + jsContexts: { + "0": [ + { thread: 2, js_context: 0, message: "start" }, + { thread: 2, js_context: 0, message: "end" }, + { thread: 2, js_context: 0, message: "resolved" }, + ], + "1": [ + { thread: 2, js_context: 1, message: "start" }, + { thread: 2, js_context: 1, message: "end" }, + { thread: 2, js_context: 1, message: "resolved" }, + ], + }, + eventLoop: { + "0": [ + { + thread: 1000, + js_context: 0, + event_loop: true, + message: "start", + }, + { + thread: 1000, + js_context: 0, + event_loop: true, + message: "end", + }, + ], + "1": [ + { + thread: 1000, + js_context: 1, + event_loop: true, + message: "start", + }, + { + thread: 1000, + js_context: 1, + event_loop: true, + message: "end", + }, + ], + }, + }); +}); + +Deno.test("should_wait_for_code_to_finish_contexts", async () => { + const example = "should_wait_for_code_to_finish_contexts"; + const results = await run(example); + assertObjectMatch(results, { + main: [ + { thread: 1, message: "start" }, + { thread: 1, message: "end" }, + ], + jsContexts: { + "0": [ + { thread: 2, js_context: 0, message: "start" }, + { thread: 2, js_context: 0, message: "end" }, + { thread: 2, js_context: 0, message: "start" }, + { thread: 2, js_context: 0, message: "end" }, + { thread: 2, js_context: 0, message: "resolved" }, + { thread: 2, js_context: 0, message: "resolved" }, + ], + }, + eventLoop: { + "0": [ + { + thread: 1000, + js_context: 0, + event_loop: true, + message: "start", + }, + { + thread: 1000, + js_context: 0, + event_loop: true, + message: "start", + }, + { + thread: 1000, + js_context: 0, + event_loop: true, + message: "end", + }, + { + thread: 1000, + js_context: 0, + event_loop: true, + message: "end", + }, + ], + }, + }); +}); + +Deno.test("should_not_run_code_after_joining", async () => { + const example = "should_not_run_code_after_joining"; + const results = await run(example); + assertObjectMatch(results, { + main: [ + { thread: 1, message: "start" }, + { thread: 1, message: "did_not_run" }, + { thread: 1, message: "end" }, + ], + jsContexts: { + "0": [ + { thread: 2, js_context: 0, message: "start" }, + { thread: 2, js_context: 0, message: "end" }, + { thread: 2, js_context: 0, message: "resolved" }, + ], + }, + eventLoop: { + "0": [ + { + thread: 1000, + js_context: 0, + event_loop: true, + message: "start", + }, + { + thread: 1000, + js_context: 0, + event_loop: true, + message: "end", + }, + ], + }, + }); +}); diff --git a/examples/src/basic_join/mod.rs b/examples/src/basic_join/mod.rs new file mode 100644 index 0000000..32bd888 --- /dev/null +++ b/examples/src/basic_join/mod.rs @@ -0,0 +1,264 @@ +use std::sync::Arc; +use std::time::Duration; + +use ion::*; +use serde::Serialize; + +use crate::_utils::thread_id; + +#[derive(Serialize)] +struct Report { + thread: usize, + #[serde(skip_serializing_if = "Option::is_none")] + js_context: Option, + #[serde(skip_serializing_if = "Option::is_none")] + event_loop: Option, + message: String, +} + +impl Report { + fn print( + js_context: Option, + event_loop: Option, + message: &str, + ) { + let thread = if event_loop.is_none() { + thread_id::thread_id() + } else { + 1000 + }; + println!( + "{}", + serde_json::to_string(&Report { + thread, + js_context: js_context, + event_loop: event_loop, + message: message.to_string(), + }) + .unwrap() + ) + } +} + +pub fn main() -> anyhow::Result<()> { + let case = std::env::args() + .collect::>() + .get(2) + .cloned() + .expect("No code provided"); + + let runtime = JsRuntime::initialize_once(JsRuntimeOptions::default())?; + + Report::print(None, None, "start"); + #[rustfmt::skip] + match case.as_str() { + "should_cancel_when_dropped" => should_cancel_when_dropped(runtime), + "should_cancel_when_dropped_multiple" => should_cancel_when_dropped_multiple(runtime), + "should_cancel_blocking_when_dropped" => should_cancel_blocking_when_dropped(runtime), + "should_cancel_blocking_when_dropped_multiple" => should_cancel_blocking_when_dropped_multiple(runtime), + "should_wait_for_code_to_finish" => should_wait_for_code_to_finish(runtime), + "should_wait_for_code_to_finish_multiple" => should_wait_for_code_to_finish_multiple(runtime), + "should_wait_for_code_to_finish_blocking" => should_wait_for_code_to_finish_blocking(runtime), + "should_wait_for_code_to_finish_worker" => should_wait_for_code_to_finish_worker(runtime), + "should_wait_for_code_to_finish_worker_blocking" => should_wait_for_code_to_finish_worker_blocking(runtime), + "should_wait_for_code_to_finish_context" => should_wait_for_code_to_finish_context(runtime), + "should_wait_for_code_to_finish_context_blocking" => should_wait_for_code_to_finish_context_blocking(runtime), + "should_wait_for_code_to_finish_contexts" => should_wait_for_code_to_finish_contexts(runtime), + "should_wait_for_code_to_finish_contexts_blocking" => should_wait_for_code_to_finish_contexts_blocking(runtime), + "should_not_run_code_after_joining" => should_not_run_code_after_joining(runtime), + _ => panic!("No Case Selected"), + }?; + + Report::print(None, None, "end"); + Ok(()) +} + +fn non_blocking_exec(context: usize) -> Box ion::Result<()>> { + return Box::new(move |env| { + env.inc_ref(); + Report::print(Some(context), None, "start"); + + env.spawn_background({ + let env = env.as_async(); + + async move { + Report::print(Some(context), Some(true), "start"); + tokio::time::sleep(Duration::from_millis(100)).await; + Report::print(Some(context), Some(true), "end"); + + env.exec_async(move |env| { + env.dec_ref(); + Report::print(Some(context), None, "resolved"); + Ok(()) + }) + .await + } + })?; + + Report::print(Some(context), None, "end"); + Ok(()) + }); +} + +fn should_cancel_when_dropped(runtime: Arc) -> anyhow::Result<()> { + let w0 = runtime.spawn_worker(JsWorkerOptions::default())?; + let c0 = w0.create_context()?; + + c0.exec(non_blocking_exec(0))?; + Ok(()) +} + +fn should_cancel_when_dropped_multiple(runtime: Arc) -> anyhow::Result<()> { + let w0 = runtime.spawn_worker(JsWorkerOptions::default())?; + let c0 = w0.create_context()?; + + c0.exec(non_blocking_exec(0))?; + c0.exec(non_blocking_exec(0))?; + + Ok(()) +} + +fn should_cancel_blocking_when_dropped(runtime: Arc) -> anyhow::Result<()> { + let w0 = runtime.spawn_worker(JsWorkerOptions::default())?; + let c0 = w0.create_context()?; + + c0.exec_blocking(non_blocking_exec(0))?; + + Ok(()) +} + +fn should_cancel_blocking_when_dropped_multiple(runtime: Arc) -> anyhow::Result<()> { + let w0 = runtime.spawn_worker(JsWorkerOptions::default())?; + let c0 = w0.create_context()?; + + c0.exec_blocking(non_blocking_exec(0))?; + c0.exec_blocking(non_blocking_exec(0))?; + + Ok(()) +} + +fn should_wait_for_code_to_finish(runtime: Arc) -> anyhow::Result<()> { + let w0 = runtime.spawn_worker(JsWorkerOptions::default())?; + let c0 = w0.create_context()?; + + c0.exec(non_blocking_exec(0))?; + + c0.join()?; + w0.join()?; + + Ok(()) +} + +fn should_wait_for_code_to_finish_multiple(runtime: Arc) -> anyhow::Result<()> { + let w0 = runtime.spawn_worker(JsWorkerOptions::default())?; + let c0 = w0.create_context()?; + + c0.exec(non_blocking_exec(0))?; + c0.exec(non_blocking_exec(0))?; + + c0.join()?; + w0.join()?; + + Ok(()) +} + +fn should_wait_for_code_to_finish_blocking(runtime: Arc) -> anyhow::Result<()> { + let w0 = runtime.spawn_worker(JsWorkerOptions::default())?; + let c0 = w0.create_context()?; + + c0.exec_blocking(non_blocking_exec(0))?; + + c0.join()?; + w0.join()?; + + Ok(()) +} + +fn should_wait_for_code_to_finish_worker(runtime: Arc) -> anyhow::Result<()> { + let w0 = runtime.spawn_worker(JsWorkerOptions::default())?; + let c0 = w0.create_context()?; + + c0.exec(non_blocking_exec(0))?; + + w0.join()?; + + Ok(()) +} + +fn should_wait_for_code_to_finish_worker_blocking(runtime: Arc) -> anyhow::Result<()> { + let w0 = runtime.spawn_worker(JsWorkerOptions::default())?; + let c0 = w0.create_context()?; + + c0.exec_blocking(non_blocking_exec(0))?; + + w0.join()?; + + Ok(()) +} + +fn should_wait_for_code_to_finish_context(runtime: Arc) -> anyhow::Result<()> { + let w0 = runtime.spawn_worker(JsWorkerOptions::default())?; + let c0 = w0.create_context()?; + + c0.exec(non_blocking_exec(0))?; + + c0.join()?; + + Ok(()) +} + +fn should_wait_for_code_to_finish_context_blocking(runtime: Arc) -> anyhow::Result<()> { + let w0 = runtime.spawn_worker(JsWorkerOptions::default())?; + let c0 = w0.create_context()?; + + c0.exec_blocking(non_blocking_exec(0))?; + + c0.join()?; + + Ok(()) +} + +fn should_wait_for_code_to_finish_contexts(runtime: Arc) -> anyhow::Result<()> { + let w0 = runtime.spawn_worker(JsWorkerOptions::default())?; + + let c0 = w0.create_context()?; + let c1 = w0.create_context()?; + + c0.exec(non_blocking_exec(0))?; + c1.exec(non_blocking_exec(0))?; + + c0.join()?; + c1.join()?; + + Ok(()) +} + +fn should_wait_for_code_to_finish_contexts_blocking(runtime: Arc) -> anyhow::Result<()> { + let w0 = runtime.spawn_worker(JsWorkerOptions::default())?; + + let c0 = w0.create_context()?; + let c1 = w0.create_context()?; + + c0.exec_blocking(non_blocking_exec(0))?; + c1.exec_blocking(non_blocking_exec(1))?; + + c0.join()?; + c1.join()?; + + Ok(()) +} + +fn should_not_run_code_after_joining(runtime: Arc) -> anyhow::Result<()> { + let w0 = runtime.spawn_worker(JsWorkerOptions::default())?; + let c0 = w0.create_context()?; + + c0.exec_blocking(non_blocking_exec(0))?; + + w0.join()?; + + if c0.exec_blocking(non_blocking_exec(0)).is_err() { + Report::print(None, None, "did_not_run"); + }; + + Ok(()) +} diff --git a/examples/src/custom_extension/mod.rs b/examples/src/custom_extension/mod.rs index 1c4a9d7..edbaaae 100644 --- a/examples/src/custom_extension/mod.rs +++ b/examples/src/custom_extension/mod.rs @@ -24,6 +24,8 @@ pub fn main() -> anyhow::Result<()> { Ok(()) })?; + ctx.join()?; + Ok(()) } diff --git a/examples/src/custom_resolver/mod.rs b/examples/src/custom_resolver/mod.rs index a21cc5d..c4c9e8c 100644 --- a/examples/src/custom_resolver/mod.rs +++ b/examples/src/custom_resolver/mod.rs @@ -23,6 +23,7 @@ pub fn main() -> anyhow::Result<()> { ctx.import(&entry_point)?; + ctx.join()?; Ok(()) } diff --git a/examples/src/deferred/mod.rs b/examples/src/deferred/mod.rs index 26b2f3f..4c69977 100644 --- a/examples/src/deferred/mod.rs +++ b/examples/src/deferred/mod.rs @@ -63,5 +63,6 @@ pub fn main() -> anyhow::Result<()> { Ok(()) })?; + ctx.join()?; Ok(()) } diff --git a/examples/src/eval/mod.rs b/examples/src/eval/mod.rs index 323e226..19f31b9 100644 --- a/examples/src/eval/mod.rs +++ b/examples/src/eval/mod.rs @@ -28,5 +28,7 @@ pub fn main() -> anyhow::Result<()> { let ctx = worker.create_context()?; ctx.eval(code)?; + ctx.join()?; + Ok(()) } diff --git a/examples/src/external_value/mod.rs b/examples/src/external_value/mod.rs index bc25342..d0ebc57 100644 --- a/examples/src/external_value/mod.rs +++ b/examples/src/external_value/mod.rs @@ -40,5 +40,6 @@ pub fn main() -> anyhow::Result<()> { Ok(()) })?; + ctx.join()?; Ok(()) } diff --git a/examples/src/main.rs b/examples/src/main.rs index 38d894d..45237aa 100644 --- a/examples/src/main.rs +++ b/examples/src/main.rs @@ -1,14 +1,15 @@ -#![deny(unused_crate_dependencies)] +// #![deny(unused_crate_dependencies)] mod _utils; mod background_tasks; mod basic; +mod basic_join; mod context_multiplexing; mod custom_extension; mod custom_resolver; mod deferred; mod eval; mod external_value; -mod http_server; +// mod http_server; mod memory_usage_context; mod memory_usage_external_value; mod memory_usage_module; @@ -34,11 +35,12 @@ fn main() -> anyhow::Result<()> { match example.as_str() { "basic" => basic::main(), + "basic_join" => basic_join::main(), "custom_extension" => custom_extension::main(), "custom_resolver" => custom_resolver::main(), "deferred" => deferred::main(), "eval" => eval::main(), - "http_server" => http_server::main(), + // "http_server" => http_server::main(), "promise" => promise::main(), "run" => run::main(), "set_interval" => set_interval::main(), diff --git a/examples/src/memory_usage_context/mod.rs b/examples/src/memory_usage_context/mod.rs index e02c187..316e253 100644 --- a/examples/src/memory_usage_context/mod.rs +++ b/examples/src/memory_usage_context/mod.rs @@ -34,7 +34,7 @@ pub fn main() -> anyhow::Result<()> { let ctx1 = worker.create_context()?; ctx0.eval("globalThis.value = []")?; - for i in 0..100 { + for i in 0..1 { ctx0.eval(format!("globalThis.value.push({})", i))?; } @@ -43,12 +43,12 @@ pub fn main() -> anyhow::Result<()> { ctx1.eval(format!("globalThis.value.push({})", i))?; } - drop(ctx0); - drop(ctx1); + ctx0.join()?; + ctx1.join()?; }; worker.run_garbage_collection_for_testing()?; - drop(worker); + worker.join()?; println!("{}", memu.megabytes().json()); } diff --git a/examples/src/multiple_workers/mod.rs b/examples/src/multiple_workers/mod.rs index 9ae7b7a..65b396b 100644 --- a/examples/src/multiple_workers/mod.rs +++ b/examples/src/multiple_workers/mod.rs @@ -30,5 +30,9 @@ pub fn main() -> anyhow::Result<()> { wrk2ctx1.eval("console.log('wrk2ctx1')")?; wrk3ctx1.eval("console.log('wrk3ctx1')")?; + wrk1ctx1.join()?; + wrk2ctx1.join()?; + wrk3ctx1.join()?; + Ok(()) } diff --git a/examples/src/promise/mod.rs b/examples/src/promise/mod.rs index 5ab3b8a..b51a9c0 100644 --- a/examples/src/promise/mod.rs +++ b/examples/src/promise/mod.rs @@ -55,5 +55,6 @@ pub fn main() -> anyhow::Result<()> { Ok(()) })?; + ctx.join()?; Ok(()) } diff --git a/examples/src/run/mod.rs b/examples/src/run/mod.rs index dc5e6d1..de58f77 100644 --- a/examples/src/run/mod.rs +++ b/examples/src/run/mod.rs @@ -43,6 +43,7 @@ pub fn main() -> anyhow::Result<()> { let ctx = worker.create_context()?; ctx.import(file_path.try_to_string()?)?; + ctx.join()?; Ok(()) } diff --git a/examples/src/set_interval/mod.rs b/examples/src/set_interval/mod.rs index 4eb79a9..2359f52 100644 --- a/examples/src/set_interval/mod.rs +++ b/examples/src/set_interval/mod.rs @@ -40,5 +40,7 @@ pub fn main() -> anyhow::Result<()> { Ok(()) })?; + + ctx.join()?; Ok(()) } diff --git a/examples/src/set_timeout/mod.rs b/examples/src/set_timeout/mod.rs index 2275acf..562b7e5 100644 --- a/examples/src/set_timeout/mod.rs +++ b/examples/src/set_timeout/mod.rs @@ -42,5 +42,7 @@ pub fn main() -> anyhow::Result<()> { Ok(()) })?; + + ctx.join()?; Ok(()) } diff --git a/examples/src/thread_safe_function/mod.rs b/examples/src/thread_safe_function/mod.rs index b8c594e..86e6c17 100644 --- a/examples/src/thread_safe_function/mod.rs +++ b/examples/src/thread_safe_function/mod.rs @@ -78,5 +78,6 @@ pub fn main() -> anyhow::Result<()> { Ok(()) })?; + ctx.join()?; Ok(()) } diff --git a/examples/src/thread_safe_promise/mod.rs b/examples/src/thread_safe_promise/mod.rs index f5eeba1..2550f0a 100644 --- a/examples/src/thread_safe_promise/mod.rs +++ b/examples/src/thread_safe_promise/mod.rs @@ -47,5 +47,6 @@ pub fn main() -> anyhow::Result<()> { println!("[Rust] Got {}", result); + ctx.join()?; Ok(()) } diff --git a/examples/src/transformers/mod.rs b/examples/src/transformers/mod.rs index cc9fdbb..83e1f2e 100644 --- a/examples/src/transformers/mod.rs +++ b/examples/src/transformers/mod.rs @@ -34,5 +34,6 @@ pub fn main() -> anyhow::Result<()> { ctx.exec_blocking(move |env| env.import(entry_point.try_to_string()?))?; + ctx.join()?; Ok(()) } diff --git a/examples/src/typescript/mod.rs b/examples/src/typescript/mod.rs index 3e8ae40..3e68663 100644 --- a/examples/src/typescript/mod.rs +++ b/examples/src/typescript/mod.rs @@ -34,5 +34,6 @@ pub fn main() -> anyhow::Result<()> { ctx.exec_blocking(move |env| env.import(entry_point.try_to_string()?))?; + ctx.join()?; Ok(()) } diff --git a/examples/test-utils/run_test.ts b/examples/test-utils/run_test.ts index a435f5f..d8591c3 100644 --- a/examples/test-utils/run_test.ts +++ b/examples/test-utils/run_test.ts @@ -6,7 +6,7 @@ export async function executeExample(testName: string, args: string[] = [], env: const command = new Deno.Command(Paths["~/"]("target", "debug", binName), { args: [testName, ...args], stdout: "piped", - stderr: "piped", + stderr: "inherit", cwd: Paths["~"], env: { ...Deno.env.toObject(), @@ -14,12 +14,11 @@ export async function executeExample(testName: string, args: string[] = [], env: } }); - const { code, stdout, stderr } = await command.output(); + const { code, stdout } = await command.output(); if (code !== 0) { - const errorText = new TextDecoder().decode(stderr); throw new Error( - `Test '${testName}' failed with exit code ${code}:\n${errorText}` + `Test '${testName}' failed with exit code ${code}:\n` ); }