Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions changelog.d/8825-async-hooks-lifecycle.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
### Fixed

- Completed Node `async_hooks` lifecycle support across async resources, event emitters, HTTP, sockets, workers, zlib, DNS, and WebCrypto. Provider scopes now restore execution and `AsyncLocalStorage` state when hooks or callbacks throw, deferred destroy hooks run at the correct lifecycle boundary, and allocation-sensitive values remain rooted across moving garbage collections.
25 changes: 17 additions & 8 deletions crates/perry-codegen/src/expr/this_super_call.rs
Original file line number Diff line number Diff line change
Expand Up @@ -261,23 +261,32 @@ pub(crate) fn lower(ctx: &mut FnCtx<'_>, expr: &Expr) -> Result<String> {
let async_parent = ctx
.classes
.get(&current_class_name)
.and_then(|class| class.extends_name.clone());
.filter(|class| class.extends_expr.is_none() && !class.heritage_lexically_shadowed)
.and_then(|class| class.extends_name.clone())
.filter(|parent| !ctx.classes.contains_key(parent.as_str()));
if matches!(
async_parent.as_deref(),
Some("EventEmitterAsyncResource" | "AsyncLocalStorage" | "AsyncResource")
) {
let undef = double_literal(f64::from_bits(crate::nanbox::TAG_UNDEFINED));
let zero_idx = "0".to_string();
let one_idx = "1".to_string();
let first =
ctx.block()
.call(DOUBLE, "js_array_get_f64", &[(I64, &arr), (I32, &zero_idx)]);
let second =
ctx.block()
.call(DOUBLE, "js_array_get_f64", &[(I64, &arr), (I32, &one_idx)]);
rooting::with_rooted_group(ctx, 3, |ctx, group| {
rooting::with_rooted_group(ctx, 4, |ctx, group| {
let this_root = group.adopt_emitted(ctx, Repr::Boxed, &this_box, true);
let arr_root = group.adopt_emitted(ctx, Repr::Ptr, &arr, true);
let arr = group.reread_emitted(ctx, arr_root);
let first = ctx.block().call(
DOUBLE,
"js_array_get_f64",
&[(I64, &arr), (I32, &zero_idx)],
);
let first_root = group.adopt_emitted(ctx, Repr::Boxed, &first, true);
let arr = group.reread_emitted(ctx, arr_root);
let second = ctx.block().call(
DOUBLE,
"js_array_get_f64",
&[(I64, &arr), (I32, &one_idx)],
);
let second_root = group.adopt_emitted(ctx, Repr::Boxed, &second, true);
let this_box = group.reread_emitted(ctx, this_root);
match async_parent.as_deref() {
Expand Down
11 changes: 5 additions & 6 deletions crates/perry-codegen/src/lower_call/builtin.rs
Original file line number Diff line number Diff line change
Expand Up @@ -153,6 +153,7 @@ pub(super) fn lower_builtin_new<'a>(
Some(index) => group.reread(ctx, index)?,
None => double_literal(f64::from_bits(crate::nanbox::TAG_UNDEFINED)),
};
let options = group.adopt_emitted(ctx, crate::rooting::Repr::Boxed, &options, true);
let runtime = if import_src.is_some_and(|source| {
source.strip_prefix("node:").unwrap_or(source) == "dns/promises"
}) {
Expand All @@ -163,12 +164,10 @@ pub(super) fn lower_builtin_new<'a>(
ctx.pending_declares
.push((runtime.to_string(), DOUBLE, vec![I64]));
let zero = "0".to_string();
let args_array = ctx.block().call(I64, "js_array_alloc", &[(I32, &zero)]);
let args_array = ctx.block().call(
I64,
"js_array_push_f64",
&[(I64, &args_array), (DOUBLE, &options)],
);
let args_array = group.begin_array(ctx, &zero);
let options = group.reread_emitted(ctx, options);
group.push_array(ctx, args_array, &options);
let args_array = group.read_array(ctx, args_array);
Ok(Some(ctx.block().call(
DOUBLE,
runtime,
Expand Down
32 changes: 32 additions & 0 deletions crates/perry-ext-events/src/emit_scope.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
use super::*;

pub(super) struct EventEmitterEmitCall {
pub(super) handle: Handle,
pub(super) event_value: TransientRootedNanbox,
pub(super) args_ptr: TransientRootedAddr,
}

pub(super) unsafe extern "C" fn event_emitter_emit_thunk(data: *mut c_void) -> f64 {
let call = &mut *(data as *mut EventEmitterEmitCall);
let Some(event_name) = event_name_from_bits(call.event_value.get().to_bits() as i64) else {
return f64::from_bits(0x7FFC_0000_0000_0003);
};
js_event_emitter_emit_impl(
call.handle,
&event_name,
call.args_ptr.get() as *mut ArrayHeader,
)
}

pub(super) struct EventEmitterEmit0Call {
pub(super) handle: Handle,
pub(super) event_value: TransientRootedNanbox,
}

pub(super) unsafe extern "C" fn event_emitter_emit0_thunk(data: *mut c_void) -> f64 {
let call = &mut *(data as *mut EventEmitterEmit0Call);
let Some(event_name) = event_name_from_bits(call.event_value.get().to_bits() as i64) else {
return f64::from_bits(0x7FFC_0000_0000_0003);
};
js_event_emitter_emit0_impl(call.handle, &event_name)
}
98 changes: 49 additions & 49 deletions crates/perry-ext-events/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,14 +23,20 @@ use perry_ffi::{
error_value_with_code, js_array_alloc, js_array_get, js_array_length, js_array_push,
js_array_set, js_object_alloc_with_shape, js_object_set_field, nanbox_string_bits, read_string,
throw_with_code, ArrayHeader, ErrorKind, Handle, JsPromise, JsString, JsValue, ObjectHeader,
Promise, RawClosureHeader, StringHeader, TransientRootScope,
Promise, RawClosureHeader, StringHeader, TransientRootScope, TransientRootedAddr,
TransientRootedNanbox,
};
use std::collections::{HashMap, HashSet};
use std::ffi::c_void;
use std::sync::{Mutex, MutexGuard, Once, OnceLock};

mod error_monitor;
use error_monitor::dispatch_error_monitor;
mod emit_scope;
use emit_scope::{
event_emitter_emit0_thunk, event_emitter_emit_thunk, EventEmitterEmit0Call,
EventEmitterEmitCall,
};
mod max_listeners;
mod messages;
mod module_helpers;
Expand Down Expand Up @@ -686,16 +692,29 @@ unsafe fn string_from_header(ptr: *const StringHeader) -> Option<String> {
read_string(handle).map(String::from)
}

fn is_raw_string_header_bits(raw: u64) -> bool {
(0x10000..MAX_HEAP_POINTER).contains(&raw) && (raw & TAG_MASK) == 0
}

unsafe fn event_name_from_bits(event_bits: i64) -> Option<String> {
let raw = event_bits as u64;
if (0x10000..MAX_HEAP_POINTER).contains(&raw) && (raw & TAG_MASK) == 0 {
if is_raw_string_header_bits(raw) {
return string_from_header(raw as *const StringHeader);
}

let rendered = js_jsvalue_to_string(f64::from_bits(raw));
string_from_header(rendered as *const StringHeader)
}

fn event_value_from_bits(event_bits: i64) -> f64 {
let raw = event_bits as u64;
if is_raw_string_header_bits(raw) {
f64::from_bits(nanbox_string_bits(raw as *mut StringHeader))
} else {
f64::from_bits(raw)
}
}

fn event_bits_from_string_ptr(ptr: *const StringHeader) -> i64 {
f64::from_bits(nanbox_string_bits(ptr as *mut StringHeader)).to_bits() as i64
}
Expand Down Expand Up @@ -1296,16 +1315,19 @@ pub unsafe extern "C" fn js_event_emitter_emit(
event_bits: i64,
args_ptr: *mut ArrayHeader,
) -> f64 {
if event_name_from_bits(event_bits).is_none() {
return f64::from_bits(0x7FFC_0000_0000_0003);
}
let roots = TransientRootScope::enter();
let event_value = roots.root_nanbox(event_value_from_bits(event_bits));
let args_ptr = roots.root_addr(args_ptr as i64);
let async_id = event_emitter_async_id(handle);
if async_id == 0 {
return js_event_emitter_emit_impl(handle, event_bits, args_ptr);
let Some(event_name) = event_name_from_bits(event_value.get().to_bits() as i64) else {
return f64::from_bits(0x7FFC_0000_0000_0003);
};
return js_event_emitter_emit_impl(handle, &event_name, args_ptr.get() as *mut ArrayHeader);
}
let mut call = EventEmitterEmitCall {
handle,
event_bits,
event_value,
args_ptr,
};
js_async_hooks_provider_run_catching(
Expand All @@ -1315,42 +1337,28 @@ pub unsafe extern "C" fn js_event_emitter_emit(
)
}

struct EventEmitterEmitCall {
handle: Handle,
event_bits: i64,
args_ptr: *mut ArrayHeader,
}

unsafe extern "C" fn event_emitter_emit_thunk(data: *mut std::ffi::c_void) -> f64 {
let call = &mut *(data as *mut EventEmitterEmitCall);
js_event_emitter_emit_impl(call.handle, call.event_bits, call.args_ptr)
}

unsafe fn js_event_emitter_emit_impl(
handle: Handle,
event_bits: i64,
event_name: &str,
args_ptr: *mut ArrayHeader,
) -> f64 {
const TAG_FALSE_F64: f64 = f64::from_bits(0x7FFC_0000_0000_0003);
const TAG_TRUE_F64: f64 = f64::from_bits(0x7FFC_0000_0000_0004);
let Some(event_name) = event_name_from_bits(event_bits) else {
return TAG_FALSE_F64;
};
let mut had_listeners = false;
let mut domain_error: Option<(Handle, f64)> = None;
let mut throw_error: Option<f64> = None;
if let Some(emitter) = get_event_emitter_mut(handle) {
let snapshot: Vec<Listener> = match emitter.events.get(&event_name) {
let snapshot: Vec<Listener> = match emitter.events.get(event_name) {
Some(v) if !v.is_empty() => v.clone(),
_ => Vec::new(),
};
if !snapshot.is_empty() {
had_listeners = true;
if snapshot.iter().any(|l| l.once) {
if let Some(v) = emitter.events.get_mut(&event_name) {
if let Some(v) = emitter.events.get_mut(event_name) {
v.retain(|l| !l.once);
}
emitter.prune_event_if_empty(&event_name);
emitter.prune_event_if_empty(event_name);
}
}

Expand All @@ -1374,7 +1382,7 @@ unsafe fn js_event_emitter_emit_impl(
}

if domain_error.is_none() && throw_error.is_none() {
drain_pending_once_promises(emitter, &event_name, args_ptr);
drain_pending_once_promises(emitter, event_name, args_ptr);

let capture_rejections = emitter.capture_rejections && event_name != "error";
for l in snapshot {
Expand Down Expand Up @@ -1412,52 +1420,44 @@ unsafe fn js_event_emitter_emit_impl(
/// `event_name_ptr` must be null or a Perry-runtime `StringHeader`.
#[no_mangle]
pub unsafe extern "C" fn js_event_emitter_emit0(handle: Handle, event_bits: i64) -> f64 {
if event_name_from_bits(event_bits).is_none() {
return f64::from_bits(0x7FFC_0000_0000_0003);
}
let roots = TransientRootScope::enter();
let event_value = roots.root_nanbox(event_value_from_bits(event_bits));
let async_id = event_emitter_async_id(handle);
if async_id == 0 {
return js_event_emitter_emit0_impl(handle, event_bits);
let Some(event_name) = event_name_from_bits(event_value.get().to_bits() as i64) else {
return f64::from_bits(0x7FFC_0000_0000_0003);
};
return js_event_emitter_emit0_impl(handle, &event_name);
}
let mut call = EventEmitterEmit0Call { handle, event_bits };
let mut call = EventEmitterEmit0Call {
handle,
event_value,
};
js_async_hooks_provider_run_catching(
async_id,
event_emitter_emit0_thunk,
&mut call as *mut EventEmitterEmit0Call as *mut std::ffi::c_void,
)
}

struct EventEmitterEmit0Call {
handle: Handle,
event_bits: i64,
}

unsafe extern "C" fn event_emitter_emit0_thunk(data: *mut std::ffi::c_void) -> f64 {
let call = &mut *(data as *mut EventEmitterEmit0Call);
js_event_emitter_emit0_impl(call.handle, call.event_bits)
}

unsafe fn js_event_emitter_emit0_impl(handle: Handle, event_bits: i64) -> f64 {
unsafe fn js_event_emitter_emit0_impl(handle: Handle, event_name: &str) -> f64 {
const TAG_FALSE_F64: f64 = f64::from_bits(0x7FFC_0000_0000_0003);
const TAG_TRUE_F64: f64 = f64::from_bits(0x7FFC_0000_0000_0004);
let Some(event_name) = event_name_from_bits(event_bits) else {
return TAG_FALSE_F64;
};
let mut had_listeners = false;
let mut domain_error: Option<(Handle, f64)> = None;
let mut throw_error: Option<f64> = None;
if let Some(emitter) = get_event_emitter_mut(handle) {
let snapshot: Vec<Listener> = match emitter.events.get(&event_name) {
let snapshot: Vec<Listener> = match emitter.events.get(event_name) {
Some(v) if !v.is_empty() => v.clone(),
_ => Vec::new(),
};
if !snapshot.is_empty() {
had_listeners = true;
if snapshot.iter().any(|l| l.once) {
if let Some(v) = emitter.events.get_mut(&event_name) {
if let Some(v) = emitter.events.get_mut(event_name) {
v.retain(|l| !l.once);
}
emitter.prune_event_if_empty(&event_name);
emitter.prune_event_if_empty(event_name);
}
}

Expand All @@ -1480,7 +1480,7 @@ unsafe fn js_event_emitter_emit0_impl(handle: Handle, event_bits: i64) -> f64 {
}
}
if domain_error.is_none() && throw_error.is_none() {
drain_pending_once_promises(emitter, &event_name, empty_args);
drain_pending_once_promises(emitter, event_name, empty_args);

let capture_rejections = emitter.capture_rejections && event_name != "error";
for l in snapshot {
Expand Down
6 changes: 4 additions & 2 deletions crates/perry-ext-http/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1802,16 +1802,18 @@ pub unsafe extern "C" fn js_http_once(
if callback == 0 {
return handle;
}
let roots = perry_ffi::TransientRootScope::enter();
let callback = roots.root_addr(callback);
let wrapper =
client_request_surface::create_client_once_wrapper(handle, &event, callback, false);
client_request_surface::create_client_once_wrapper(handle, &event, callback.get(), false);
let mut matched = false;
with_handle_mut::<ClientRequestHandle, _, _>(handle, |request| {
request
.listeners
.entry(event.clone())
.or_default()
.push(ClientEventListener {
callback,
callback: callback.get(),
raw_wrapper: wrapper,
once: true,
});
Expand Down
10 changes: 0 additions & 10 deletions crates/perry-ext-http/src/server/handle_dispatch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -126,8 +126,6 @@ extern "C" {
fn js_node_http_im_resume(handle: i64);
fn js_node_http_im_destroy(handle: i64);
fn js_node_http_im_on(handle: i64, event_name_ptr: *const StringHeader, callback: i64) -> f64;
fn js_node_http_im_once(handle: i64, event_name_ptr: *const StringHeader, callback: i64)
-> f64;
fn js_node_http_im_set_encoding(handle: i64, encoding_ptr: *const StringHeader) -> i64;
fn js_node_http_im_set_timeout(handle: i64, msecs: f64, callback: i64) -> i64;
fn js_node_http_im_read(handle: i64) -> f64;
Expand Down Expand Up @@ -374,14 +372,6 @@ pub unsafe extern "C" fn js_ext_http_server_dispatch_method(
}
self_ref
}
"once" if args.len() >= 2 => {
let event_ptr = string_arg(args[0]);
if event_ptr.is_null() {
return self_ref;
}
js_node_http_im_once(handle, event_ptr, closure_arg(Some(args[1])));
self_ref
}
"once" if args.len() >= 2 => {
let event =
read_string_header(string_arg(args[0]) as *mut StringHeader).unwrap_or_default();
Expand Down
1 change: 1 addition & 0 deletions crates/perry-ext-net/src/adopt.rs
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,7 @@ pub fn adopt_upgraded_tcp_stream(stream: tokio::net::TcpStream) -> i64 {
destroyed: false,
bytes_read: 0,
bytes_written: 0,
bytes_queued: 0,
timeout: None,
type_of_service: 0,
server_id: None,
Expand Down
2 changes: 2 additions & 0 deletions crates/perry-ext-net/src/ipc.rs
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,7 @@ fn allocate_socket() -> (i64, mpsc::UnboundedReceiver<SocketCommand>) {
destroyed: false,
bytes_read: 0,
bytes_written: 0,
bytes_queued: 0,
timeout: None,
type_of_service: 0,
server_id: None,
Expand Down Expand Up @@ -114,6 +115,7 @@ pub(crate) fn register_accepted_transport(
destroyed: false,
bytes_read: 0,
bytes_written: 0,
bytes_queued: 0,
timeout: None,
type_of_service: 0,
server_id: Some(server_id),
Expand Down
Loading
Loading