From 76fa7b4c96035b3ae3fe49d3b5106e45d90c2325 Mon Sep 17 00:00:00 2001 From: TheHypnoo Date: Sat, 1 Aug 2026 10:19:43 +0200 Subject: [PATCH] fix(cluster): match Node 26 semantics (#7105) Align cluster exports, worker lifecycle, validation, events, networking, and IPC with Node.js 26.5.0. Complete the maintainer audit for EventEmitter semantics, worker identity and lifecycle, advanced IPC framing, scheduling, and parity coverage. Closes #6772 --- CLAUDE.md | 2 +- Cargo.lock | 152 ++--- Cargo.toml | 2 +- .../7105-node-cluster-node26-parity.md | 1 + .../src/lower_call/native_table/node_misc.rs | 40 +- .../src/type_analysis/predicates.rs | 6 +- crates/perry-ext-net/src/lib.rs | 17 +- .../src/lower/expr_call/native_module.rs | 4 + crates/perry-hir/src/lower_decl/block.rs | 66 ++- .../src/child_process/emitter.rs | 33 +- .../perry-runtime/src/child_process/fork.rs | 27 +- .../src/child_process/reactor.rs | 28 +- .../src/child_process/v8_serde.rs | 216 ++++--- .../src/child_process/validate.rs | 68 ++- crates/perry-runtime/src/cluster.rs | 552 ++++++++++++++---- crates/perry-runtime/src/cluster_sched.rs | 219 ++++++- .../src/object/class_registry.rs | 6 +- .../src/object/class_registry/construct.rs | 10 + crates/perry-runtime/src/object/instanceof.rs | 25 +- .../perry-runtime/src/object/native_module.rs | 1 + .../object/native_module/callable_exports.rs | 22 + .../native_module_dispatch/dispatch_a_c.rs | 8 +- .../src/object/native_module_registry.rs | 9 +- crates/perry-runtime/src/process.rs | 15 +- crates/perry-runtime/src/process/ipc.rs | 112 +++- test-parity/node-suite/cluster/README.md | 4 +- test-parity/node-suite/cluster/STATUS.md | 51 +- .../cluster/default-import-shape.ts | 4 + .../cluster/events/listener-validation.ts | 59 ++ .../node-suite/cluster/fork/roles-and-env.ts | 8 +- .../cluster/serialization/advanced-tcp.ts | 27 + .../cluster/serialization/advanced.ts | 2 +- .../cluster/setup/validation-exec-args.ts | 4 +- .../setup/validation-serialization-inspect.ts | 4 +- test-parity/node_suite_baseline.json | 4 +- 35 files changed, 1424 insertions(+), 384 deletions(-) create mode 100644 changelog.d/7105-node-cluster-node26-parity.md create mode 100644 test-parity/node-suite/cluster/events/listener-validation.ts create mode 100644 test-parity/node-suite/cluster/serialization/advanced-tcp.ts diff --git a/CLAUDE.md b/CLAUDE.md index 529f17abbf..ae340d3b58 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -8,7 +8,7 @@ This file provides guidance to Claude Code (claude.ai/code) when working with co Perry is a native TypeScript compiler written in Rust that compiles TypeScript source code directly to native executables. It uses SWC for TypeScript parsing and LLVM for code generation. -**Current Version:** 0.5.1270 +**Current Version:** 0.5.1271 ## TypeScript Parity Status diff --git a/Cargo.lock b/Cargo.lock index eebad50fa5..e159822af1 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -5503,7 +5503,7 @@ checksum = "9b4f627cb1b25917193a259e49bdad08f671f8d9708acfd5fe0a8c1455d87220" [[package]] name = "perry" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "anyhow", "base64", @@ -5563,14 +5563,14 @@ dependencies = [ [[package]] name = "perry-api-manifest" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "serde", ] [[package]] name = "perry-audio-miniaudio" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "cc", "libc", @@ -5578,7 +5578,7 @@ dependencies = [ [[package]] name = "perry-codegen" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "anyhow", "log", @@ -5592,7 +5592,7 @@ dependencies = [ [[package]] name = "perry-codegen-arkts" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "anyhow", "perry-hir", @@ -5600,7 +5600,7 @@ dependencies = [ [[package]] name = "perry-codegen-glance" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "anyhow", "perry-hir", @@ -5608,7 +5608,7 @@ dependencies = [ [[package]] name = "perry-codegen-js" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "anyhow", "perry-dispatch", @@ -5617,7 +5617,7 @@ dependencies = [ [[package]] name = "perry-codegen-swiftui" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "anyhow", "perry-hir", @@ -5625,7 +5625,7 @@ dependencies = [ [[package]] name = "perry-codegen-wasm" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "anyhow", "base64", @@ -5637,7 +5637,7 @@ dependencies = [ [[package]] name = "perry-codegen-wear-tiles" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "anyhow", "perry-hir", @@ -5645,7 +5645,7 @@ dependencies = [ [[package]] name = "perry-container-compose" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "anyhow", "async-trait", @@ -5674,14 +5674,14 @@ dependencies = [ [[package]] name = "perry-container-e2e" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "anyhow", ] [[package]] name = "perry-diagnostics" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "serde", "serde_json", @@ -5689,7 +5689,7 @@ dependencies = [ [[package]] name = "perry-dispatch" -version = "0.5.1270" +version = "0.5.1271" [[package]] name = "perry-doc-fixture-my-bindings" @@ -5700,7 +5700,7 @@ dependencies = [ [[package]] name = "perry-doc-tests" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "anyhow", "clap", @@ -5715,7 +5715,7 @@ dependencies = [ [[package]] name = "perry-ext-ads" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "block2", "objc2", @@ -5725,7 +5725,7 @@ dependencies = [ [[package]] name = "perry-ext-argon2" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "argon2", "perry-ffi", @@ -5733,7 +5733,7 @@ dependencies = [ [[package]] name = "perry-ext-axios" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "perry-ffi", "reqwest", @@ -5742,7 +5742,7 @@ dependencies = [ [[package]] name = "perry-ext-bcrypt" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "bcrypt", "perry-ffi", @@ -5750,7 +5750,7 @@ dependencies = [ [[package]] name = "perry-ext-better-sqlite3" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "perry-ffi", "rusqlite", @@ -5758,7 +5758,7 @@ dependencies = [ [[package]] name = "perry-ext-cheerio" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "perry-ffi", "scraper", @@ -5766,7 +5766,7 @@ dependencies = [ [[package]] name = "perry-ext-commander" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "perry-ffi", "perry-runtime", @@ -5774,7 +5774,7 @@ dependencies = [ [[package]] name = "perry-ext-cron" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "chrono", "cron", @@ -5784,7 +5784,7 @@ dependencies = [ [[package]] name = "perry-ext-dayjs" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "chrono", "perry-ffi", @@ -5792,7 +5792,7 @@ dependencies = [ [[package]] name = "perry-ext-decimal" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "perry-ffi", "rust_decimal", @@ -5800,7 +5800,7 @@ dependencies = [ [[package]] name = "perry-ext-dotenv" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "perry-ffi", "serde_json", @@ -5808,7 +5808,7 @@ dependencies = [ [[package]] name = "perry-ext-ethers" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "perry-ffi", "rand 0.10.1", @@ -5816,7 +5816,7 @@ dependencies = [ [[package]] name = "perry-ext-events" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "perry-ffi", "perry-runtime", @@ -5824,14 +5824,14 @@ dependencies = [ [[package]] name = "perry-ext-exponential-backoff" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "perry-ffi", ] [[package]] name = "perry-ext-fastify" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "bytes", "http-body-util", @@ -5849,7 +5849,7 @@ dependencies = [ [[package]] name = "perry-ext-fetch" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "bytes", "lazy_static", @@ -5862,7 +5862,7 @@ dependencies = [ [[package]] name = "perry-ext-http" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "bytes", "h2", @@ -5886,7 +5886,7 @@ dependencies = [ [[package]] name = "perry-ext-ioredis" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "lazy_static", "perry-ffi", @@ -5896,7 +5896,7 @@ dependencies = [ [[package]] name = "perry-ext-jsonwebtoken" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "base64", "jsonwebtoken", @@ -5907,7 +5907,7 @@ dependencies = [ [[package]] name = "perry-ext-lru-cache" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "lru", "perry-ffi", @@ -5916,7 +5916,7 @@ dependencies = [ [[package]] name = "perry-ext-moment" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "chrono", "perry-ffi", @@ -5924,7 +5924,7 @@ dependencies = [ [[package]] name = "perry-ext-mongodb" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "bson", "futures-util", @@ -5936,7 +5936,7 @@ dependencies = [ [[package]] name = "perry-ext-mysql2" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "chrono", "perry-ffi", @@ -5946,7 +5946,7 @@ dependencies = [ [[package]] name = "perry-ext-nanoid" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "nanoid", "perry-ffi", @@ -5955,7 +5955,7 @@ dependencies = [ [[package]] name = "perry-ext-net" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "bytes", "perry-ffi", @@ -5968,7 +5968,7 @@ dependencies = [ [[package]] name = "perry-ext-node-forge" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "const-oid 0.9.6", "der 0.7.10", @@ -5987,7 +5987,7 @@ dependencies = [ [[package]] name = "perry-ext-nodemailer" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "lettre", "perry-ffi", @@ -5997,7 +5997,7 @@ dependencies = [ [[package]] name = "perry-ext-pdf" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "perry-ffi", "printpdf", @@ -6005,7 +6005,7 @@ dependencies = [ [[package]] name = "perry-ext-pg" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "perry-ffi", "sqlx", @@ -6014,7 +6014,7 @@ dependencies = [ [[package]] name = "perry-ext-ratelimit" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "governor", "perry-ffi", @@ -6022,7 +6022,7 @@ dependencies = [ [[package]] name = "perry-ext-sharp" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "fast_image_resize", "image", @@ -6032,14 +6032,14 @@ dependencies = [ [[package]] name = "perry-ext-slugify" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "perry-ffi", ] [[package]] name = "perry-ext-streams" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "lazy_static", "perry-ffi", @@ -6048,7 +6048,7 @@ dependencies = [ [[package]] name = "perry-ext-undici" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "perry-ffi", "perry-runtime", @@ -6057,7 +6057,7 @@ dependencies = [ [[package]] name = "perry-ext-uuid" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "perry-ffi", "uuid", @@ -6065,7 +6065,7 @@ dependencies = [ [[package]] name = "perry-ext-validator" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "perry-ffi", "regex", @@ -6075,7 +6075,7 @@ dependencies = [ [[package]] name = "perry-ext-ws" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "futures-util", "lazy_static", @@ -6088,7 +6088,7 @@ dependencies = [ [[package]] name = "perry-ext-zlib" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "brotli", "flate2", @@ -6098,7 +6098,7 @@ dependencies = [ [[package]] name = "perry-ffi" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "dashmap", "once_cell", @@ -6107,7 +6107,7 @@ dependencies = [ [[package]] name = "perry-hir" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "anyhow", "perry-api-manifest", @@ -6125,7 +6125,7 @@ dependencies = [ [[package]] name = "perry-parser" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "anyhow", "perry-diagnostics", @@ -6137,7 +6137,7 @@ dependencies = [ [[package]] name = "perry-runtime" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "anyhow", "base64", @@ -6178,14 +6178,14 @@ dependencies = [ [[package]] name = "perry-runtime-static" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "perry-runtime", ] [[package]] name = "perry-stdlib" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "aes 0.8.4", "aes 0.9.1", @@ -6280,14 +6280,14 @@ dependencies = [ [[package]] name = "perry-stdlib-static" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "perry-stdlib", ] [[package]] name = "perry-transform" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "anyhow", "perry-hir", @@ -6296,14 +6296,14 @@ dependencies = [ [[package]] name = "perry-ui" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "perry-ui-model", ] [[package]] name = "perry-ui-android" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "base64", "itoa", @@ -6320,7 +6320,7 @@ dependencies = [ [[package]] name = "perry-ui-geisterhand" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "rand 0.10.1", "serde", @@ -6330,7 +6330,7 @@ dependencies = [ [[package]] name = "perry-ui-gtk4" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "base64", "cairo-rs 0.22.0", @@ -6353,7 +6353,7 @@ dependencies = [ [[package]] name = "perry-ui-ios" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "base64", "block2", @@ -6369,7 +6369,7 @@ dependencies = [ [[package]] name = "perry-ui-macos" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "base64", "block2", @@ -6384,7 +6384,7 @@ dependencies = [ [[package]] name = "perry-ui-model" -version = "0.5.1270" +version = "0.5.1271" [[package]] name = "perry-ui-test" @@ -6395,11 +6395,11 @@ dependencies = [ [[package]] name = "perry-ui-testkit" -version = "0.5.1270" +version = "0.5.1271" [[package]] name = "perry-ui-tvos" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "base64", "block2", @@ -6415,7 +6415,7 @@ dependencies = [ [[package]] name = "perry-ui-visionos" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "base64", "block2", @@ -6431,7 +6431,7 @@ dependencies = [ [[package]] name = "perry-ui-watchos" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "block2", "libc", @@ -6444,7 +6444,7 @@ dependencies = [ [[package]] name = "perry-ui-windows" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "base64", "libc", @@ -6461,14 +6461,14 @@ dependencies = [ [[package]] name = "perry-ui-windows-winui" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "perry-ui-windows", ] [[package]] name = "perry-updater" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "anyhow", "base64", @@ -6484,7 +6484,7 @@ dependencies = [ [[package]] name = "perry-wasm-host" -version = "0.5.1270" +version = "0.5.1271" dependencies = [ "wasmi", ] diff --git a/Cargo.toml b/Cargo.toml index 20a84f7fe3..7534801e3b 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -292,7 +292,7 @@ codegen-units = 16 codegen-units = 16 [workspace.package] -version = "0.5.1270" +version = "0.5.1271" edition = "2021" license = "MIT" repository = "https://github.com/PerryTS/perry" diff --git a/changelog.d/7105-node-cluster-node26-parity.md b/changelog.d/7105-node-cluster-node26-parity.md new file mode 100644 index 0000000000..30df120c86 --- /dev/null +++ b/changelog.d/7105-node-cluster-node26-parity.md @@ -0,0 +1 @@ +Match Node.js 26 behavior across `node:cluster` exports, worker lifecycle, setup validation, networking, and advanced IPC serialization. diff --git a/crates/perry-codegen/src/lower_call/native_table/node_misc.rs b/crates/perry-codegen/src/lower_call/native_table/node_misc.rs index 5e58d13cb1..7a8cfd0fda 100644 --- a/crates/perry-codegen/src/lower_call/native_table/node_misc.rs +++ b/crates/perry-codegen/src/lower_call/native_table/node_misc.rs @@ -110,9 +110,45 @@ pub(super) const NODE_MISC_ROWS: &[NativeModSig] = &[ method: "listenerCount", class_filter: None, runtime: "js_cluster_listener_count", + args: &[NA_F64, NA_F64], + ret: NR_F64, + }, + NativeModSig { + module: "cluster", + has_receiver: false, + method: "listeners", + class_filter: None, + runtime: "js_cluster_listeners", args: &[NA_F64], ret: NR_F64, }, + NativeModSig { + module: "cluster", + has_receiver: false, + method: "rawListeners", + class_filter: None, + runtime: "js_cluster_raw_listeners", + args: &[NA_F64], + ret: NR_F64, + }, + NativeModSig { + module: "cluster", + has_receiver: false, + method: "setMaxListeners", + class_filter: None, + runtime: "js_cluster_set_max_listeners", + args: &[NA_F64], + ret: NR_F64, + }, + NativeModSig { + module: "cluster", + has_receiver: false, + method: "getMaxListeners", + class_filter: None, + runtime: "js_cluster_get_max_listeners", + args: &[], + ret: NR_F64, + }, NativeModSig { module: "cluster", has_receiver: false, @@ -136,8 +172,8 @@ pub(super) const NODE_MISC_ROWS: &[NativeModSig] = &[ has_receiver: false, method: "removeAllListeners", class_filter: None, - runtime: "js_cluster_remove_all_listeners", - args: &[NA_F64], + runtime: "js_cluster_remove_all_listeners_args", + args: &[NA_VARARGS], ret: NR_F64, }, // ========== node:vm ========== diff --git a/crates/perry-codegen/src/type_analysis/predicates.rs b/crates/perry-codegen/src/type_analysis/predicates.rs index 5bde2dd15e..d443481c9f 100644 --- a/crates/perry-codegen/src/type_analysis/predicates.rs +++ b/crates/perry-codegen/src/type_analysis/predicates.rs @@ -711,7 +711,11 @@ pub(crate) fn static_type_of(ctx: &FnCtx<'_>, e: &Expr) -> Option { match (lt, rt) { (Some(a), Some(b)) if a == b => Some(a), (Some(a), Some(b)) => Some(HirType::Union(vec![a, b])), - (Some(t), None) | (None, Some(t)) => Some(t), + // One unknown branch makes the whole conditional unknown. + // Optional chaining lowers to `receiver == null ? undefined : + // receiver.property`; treating that as `Void` when the property + // type is dynamic made `Array.isArray(obj?.value)` constant-fold + // to false without inspecting the runtime value. _ => None, } } diff --git a/crates/perry-ext-net/src/lib.rs b/crates/perry-ext-net/src/lib.rs index 396712853d..9b5187d953 100644 --- a/crates/perry-ext-net/src/lib.rs +++ b/crates/perry-ext-net/src/lib.rs @@ -430,6 +430,13 @@ enum PendingNetEvent { // callback pointer with it. extern "C" { fn js_net_callback_ptr(value: f64) -> i64; + fn js_get_string_pointer_unified(value: f64) -> i64; + fn perry_cluster_worker_listening( + addr_ptr: *const u8, + addr_len: u32, + port: i32, + address_type: i32, + ); } fn push_event(ev: PendingNetEvent) { @@ -700,7 +707,8 @@ pub unsafe extern "C" fn js_net_server_listen(handle: i64, port: f64, arg2: f64, // is left alone.) js_net_validate_listen_port(port); let port_u16 = port as u16; - let host = "0.0.0.0".to_string(); + let host = string_from_header_i64(js_get_string_pointer_unified(arg2)) + .unwrap_or_else(|| "0.0.0.0".to_string()); let (shutdown_tx, mut shutdown_rx) = oneshot::channel::<()>(); @@ -778,6 +786,13 @@ pub unsafe extern "C" fn js_net_server_listen(handle: i64, port: f64, arg2: f64, s.bound_host = local.ip().to_string(); } } + let address = local.ip().to_string(); + perry_cluster_worker_listening( + address.as_ptr(), + address.len() as u32, + local.port() as i32, + if local.is_ipv6() { 6 } else { 4 }, + ); } // bind succeeded — fire `'listening'`. push_event(PendingNetEvent::ServerListening(server_id)); diff --git a/crates/perry-hir/src/lower/expr_call/native_module.rs b/crates/perry-hir/src/lower/expr_call/native_module.rs index 1f3c9d354e..14faff1e0c 100644 --- a/crates/perry-hir/src/lower/expr_call/native_module.rs +++ b/crates/perry-hir/src/lower/expr_call/native_module.rs @@ -31,6 +31,10 @@ fn is_cluster_default_event_emitter_method(method_name: &str) -> bool { | "emit" | "eventNames" | "listenerCount" + | "listeners" + | "rawListeners" + | "getMaxListeners" + | "setMaxListeners" | "removeListener" | "off" | "removeAllListeners" diff --git a/crates/perry-hir/src/lower_decl/block.rs b/crates/perry-hir/src/lower_decl/block.rs index 3727e17af6..52f6122e53 100644 --- a/crates/perry-hir/src/lower_decl/block.rs +++ b/crates/perry-hir/src/lower_decl/block.rs @@ -1679,11 +1679,75 @@ pub fn lower_block_stmt_scoped( let mark = ctx.push_block_scope(); // Via `lower_block_stmt` so this scope's pre-registered forward-captured // lets are re-bound at entry (`rebind_nested_forward_scope_lets`). - let stmts = lower_block_stmt(ctx, block)?; + let stmts = if ctx.current_strict { + lower_strict_block_fn_decls(ctx, block)? + } else { + lower_block_stmt(ctx, block)? + }; ctx.pop_block_scope(mark); Ok(stmts) } +/// Strict-mode block function declarations are lexical bindings initialized +/// when the block is entered. Pre-register their locals before lowering an +/// earlier callback that captures one, then move the declarations' closure +/// initializers ahead of the block's executable statements. +fn lower_strict_block_fn_decls( + ctx: &mut LoweringContext, + block: &ast::BlockStmt, +) -> Result> { + use std::collections::HashSet; + + rebind_nested_forward_scope_lets(ctx, &block.stmts); + + let mut hoisted_ids = HashSet::new(); + for stmt in &block.stmts { + let ast::Stmt::Decl(ast::Decl::Fn(fn_decl)) = stmt else { + continue; + }; + if fn_decl.function.body.is_none() { + continue; + } + let name = fn_decl.ident.sym.to_string(); + let id = ctx + .lookup_local_in_current_scope(&name) + .unwrap_or_else(|| ctx.define_local(name, Type::Any)); + hoisted_ids.insert(id); + } + if hoisted_ids.is_empty() { + return lower_stmts_using_aware(ctx, &block.stmts); + } + + // Lower in source order first: a declaration body may capture lexical + // bindings declared earlier in the block. Only its runtime initializer is + // hoisted after every reference has resolved to the correct LocalId. + let body = lower_stmts_using_aware(ctx, &block.stmts)?; + let mut hoisted = Vec::new(); + let mut other = Vec::new(); + for stmt in body { + let is_hoisted = matches!( + &stmt, + Stmt::Let { id, init: Some(Expr::Closure { .. } | Expr::FuncRef(_)), .. } + if hoisted_ids.contains(id) + ); + if is_hoisted { + hoisted.push(stmt); + } else { + other.push(stmt); + } + } + + let combined: Vec<_> = hoisted.iter().chain(other.iter()).cloned().collect(); + let prealloc = compute_prealloc_for_hoisted_closures(&combined, &hoisted_ids); + let mut result = Vec::new(); + if !prealloc.is_empty() { + result.push(Stmt::PreallocateBoxes(prealloc)); + } + result.extend(hoisted); + result.extend(other); + Ok(result) +} + /// Lower a sequence of body statements, desugaring `using` / `await using` /// declarations into nested try/finally blocks that invoke the bound value's /// `[Symbol.dispose]()` (sync `using`) or `await [Symbol.asyncDispose]()` diff --git a/crates/perry-runtime/src/child_process/emitter.rs b/crates/perry-runtime/src/child_process/emitter.rs index d536c6252c..37ee89e0fe 100644 --- a/crates/perry-runtime/src/child_process/emitter.rs +++ b/crates/perry-runtime/src/child_process/emitter.rs @@ -311,13 +311,29 @@ pub(crate) extern "C" fn cp_method_send( a3: f64, a4: f64, ) -> f64 { + let message_value = JSValue::from_bits(message.to_bits()); + if message_value.is_undefined() { + crate::fs::validate::throw_type_error_with_code( + "The \"message\" argument must be specified", + "ERR_MISSING_ARGS", + ); + } + if unsafe { crate::symbol::js_is_symbol(message) != 0 } + || (message_value.is_pointer() && !crate::fs::extract_closure_ptr(message).is_null()) + { + crate::fs::validate::throw_type_error_with_code( + "The \"message\" argument must be one of type string, object, number, or boolean", + "ERR_INVALID_ARG_TYPE", + ); + } let this = cp_this(closure); // The callback is the last argument when it is a function. dispatch pads // missing slots with `undefined`, so scan slots 4→2 for a closure. - let callback = [a4, a3, a2] - .into_iter() - .find(|v| !crate::fs::extract_closure_ptr(*v).is_null()); + let callback = [a4, a3, a2].into_iter().find(|v| { + JSValue::from_bits(v.to_bits()).is_pointer() + && !crate::fs::extract_closure_ptr(*v).is_null() + }); // A closed IPC channel (after `disconnect()`, or never connected) returns // `false` and reports `ERR_IPC_CHANNEL_CLOSED` to the callback. @@ -394,7 +410,16 @@ pub(crate) extern "C" fn cp_method_disconnect(closure: *const ClosureHeader) -> } cp_set_field(this, b"connected", TAG_FALSE_F64); cp_set_field(this, b"channel", TAG_NULL_F64); - cp_emit(this, "disconnect", &[]); + js_register_closure_arity(cp_disconnect_emit_thunk as *const u8, 0); + let deferred = js_closure_alloc(cp_disconnect_emit_thunk as *const u8, 1); + js_closure_set_capture_ptr(deferred, 0, this.to_bits() as i64); + crate::timer::js_set_immediate_callback(deferred as i64); + cp_undefined() +} + +extern "C" fn cp_disconnect_emit_thunk(closure: *const ClosureHeader) -> f64 { + let child = f64::from_bits(js_closure_get_capture_ptr(closure, 0) as u64); + cp_emit(child, "disconnect", &[]); cp_undefined() } diff --git a/crates/perry-runtime/src/child_process/fork.rs b/crates/perry-runtime/src/child_process/fork.rs index 3b3209888a..bc20139226 100644 --- a/crates/perry-runtime/src/child_process/fork.rs +++ b/crates/perry-runtime/src/child_process/fork.rs @@ -174,12 +174,25 @@ pub extern "C" fn js_child_process_fork(module_ptr: i64, args_ptr: i64, opts_ptr cp_set_field(cp, b"spawnfile", cp_box_string(&exec_path)); // Build the command: [execArgv] [args]. + let native_self = std::env::current_exe() + .ok() + .map(|path| path.to_string_lossy().into_owned()) + .is_some_and(|path| path == exec_path && module == exec_path); let mut command = Command::new(&exec_path); - command.args(&exec_argv); - command.arg(&module); - command.args(&arg_strs); + let native_exec_argv = if native_self { + command.args(&arg_strs); + Some(serde_json::to_string(&exec_argv).unwrap_or_else(|_| "[]".to_string())) + } else { + command.args(&exec_argv); + command.arg(&module); + command.args(&arg_strs); + None + }; cp_apply_argv0(&mut command, opts_val); cp_apply_options(&mut command, opts_val); + if let Some(exec_argv) = native_exec_argv { + command.env("PERRY_PROCESS_EXEC_ARGV", exec_argv); + } cp_apply_detached(&mut command, opts_val); let launch = match cp_apply_live_stdio(&mut command, &stdio_kinds) { Ok(extra_readers) => fork_launch( @@ -280,7 +293,13 @@ fn fork_launch( // The child now holds fd 3; the parent keeps `parent_sock`. drop(child_sock); cp_set_field(cp, b"connected", TAG_TRUE_F64); - let channel = crate::object::js_object_alloc(0, 0); + let channel = cp_build_object( + &[ + ("ref", cp_cast0(cp_method_this0)), + ("unref", cp_cast0(cp_method_this0)), + ], + CP_SHAPE_ID + 0x40, + ); cp_set_field(cp, b"channel", cp_box_ptr(channel as *const u8)); let handle = reactor::cp_register_live_child( cp, diff --git a/crates/perry-runtime/src/child_process/reactor.rs b/crates/perry-runtime/src/child_process/reactor.rs index b3a414e95a..8baeaafb07 100644 --- a/crates/perry-runtime/src/child_process/reactor.rs +++ b/crates/perry-runtime/src/child_process/reactor.rs @@ -607,9 +607,31 @@ pub fn cp_ipc_send_raw_json(handle: u64, json: &str) -> bool { #[cfg(unix)] { use std::io::Write; - let mut frame = Vec::with_capacity(json.len() + 1); - frame.extend_from_slice(json.as_bytes()); - frame.push(b'\n'); + // Match the channel's selected framing. Cluster's query-server reply + // is built as JSON in Rust, but an advanced channel still requires a + // V8-serialized payload just like every user-visible IPC message. + let advanced = { + let guard = cp_live_lock(); + guard + .as_ref() + .and_then(|map| map.get(&handle)) + .map(|lc| lc.ipc_advanced) + .unwrap_or(false) + }; + let frame = if advanced { + let sh = crate::string::js_string_from_bytes(json.as_ptr(), json.len() as u32); + let message = f64::from_bits(unsafe { crate::json::js_json_parse(sh) }.bits()); + let payload = super::v8_serde::v8_serialize(message); + let mut frame = Vec::with_capacity(payload.len() + 4); + frame.extend_from_slice(&(payload.len() as u32).to_be_bytes()); + frame.extend_from_slice(&payload); + frame + } else { + let mut frame = Vec::with_capacity(json.len() + 1); + frame.extend_from_slice(json.as_bytes()); + frame.push(b'\n'); + frame + }; let mut guard = cp_live_lock(); if let Some(map) = guard.as_mut() { if let Some(lc) = map.get_mut(&handle) { diff --git a/crates/perry-runtime/src/child_process/v8_serde.rs b/crates/perry-runtime/src/child_process/v8_serde.rs index f40d2f892b..d57f1c5e4b 100644 --- a/crates/perry-runtime/src/child_process/v8_serde.rs +++ b/crates/perry-runtime/src/child_process/v8_serde.rs @@ -234,95 +234,124 @@ impl Serializer { } if jsval.is_pointer() { let raw = (bits & crate::value::POINTER_MASK) as usize; - if raw >= 0x10000 { - if crate::buffer::is_registered_buffer(raw) { - if self.write_reference_or_register(raw) { - return; - } - self.write_host_buffer(value); - return; - } - if let Some(kind) = crate::typedarray::lookup_typed_array_kind(raw) { - if self.write_reference_or_register(raw) { - return; - } - self.write_host_typed_array(value, kind); - return; - } - if crate::map::is_registered_map(raw) { - if self.depth >= MAX_DEPTH { - self.out.push(TAG_UNDEFINED); - return; - } - if self.write_reference_or_register(raw) { - return; - } - self.write_map(raw as *const crate::map::MapHeader); - return; - } - if crate::set::is_registered_set(raw) { - if self.depth >= MAX_DEPTH { - self.out.push(TAG_UNDEFINED); - return; - } - if self.write_reference_or_register(raw) { - return; - } - self.write_set(raw as *const crate::set::SetHeader); - return; - } - if crate::regex::regex_header_has_magic(raw as *const crate::regex::RegExpHeader) { - if self.write_reference_or_register(raw) { - return; - } - self.write_regexp(raw as *const crate::regex::RegExpHeader); - return; - } - if crate::error::js_error_is_error(value).to_bits() == TAG_TRUE_F64.to_bits() { - if self.depth >= MAX_DEPTH { - self.out.push(TAG_UNDEFINED); - return; - } - if self.write_reference_or_register(raw) { - return; - } - self.write_error(raw as *mut crate::error::ErrorHeader); - return; + if self.write_heap_value(value, raw) { + return; + } + } + // Functions, symbols, unknown — degrade to undefined. + self.out.push(TAG_UNDEFINED); + } + + fn write_heap_value(&mut self, value: f64, raw: usize) -> bool { + if crate::value::addr_class::is_handle_band(raw) { + return false; + } + if crate::buffer::is_registered_buffer(raw) { + if self.write_reference_or_register(raw) { + return true; + } + if crate::buffer::is_array_buffer(raw) { + self.write_array_buffer(raw as *const crate::buffer::BufferHeader); + } else { + self.write_host_buffer(value); + } + return true; + } + if let Some(kind) = crate::typedarray::lookup_typed_array_kind(raw) { + if self.write_reference_or_register(raw) { + return true; + } + self.write_host_typed_array(value, kind); + return true; + } + if crate::map::is_registered_map(raw) { + if self.depth >= MAX_DEPTH { + self.out.push(TAG_UNDEFINED); + return true; + } + if self.write_reference_or_register(raw) { + return true; + } + self.write_map(raw as *const crate::map::MapHeader); + return true; + } + if crate::set::is_registered_set(raw) { + if self.depth >= MAX_DEPTH { + self.out.push(TAG_UNDEFINED); + return true; + } + if self.write_reference_or_register(raw) { + return true; + } + self.write_set(raw as *const crate::set::SetHeader); + return true; + } + if crate::array::js_array_is_array(value).to_bits() == TAG_TRUE_F64.to_bits() { + if self.depth >= MAX_DEPTH { + self.out.push(TAG_UNDEFINED); + return true; + } + if self.write_reference_or_register(raw) { + return true; + } + let array = crate::array::js_array_from_value(value); + self.write_dense_array(array); + return true; + } + if crate::regex::regex_header_has_magic(raw as *const crate::regex::RegExpHeader) { + if self.write_reference_or_register(raw) { + return true; + } + self.write_regexp(raw as *const crate::regex::RegExpHeader); + return true; + } + if crate::error::js_error_is_error(value).to_bits() == TAG_TRUE_F64.to_bits() { + if self.depth >= MAX_DEPTH { + self.out.push(TAG_UNDEFINED); + return true; + } + if self.write_reference_or_register(raw) { + return true; + } + self.write_error(raw as *mut crate::error::ErrorHeader); + return true; + } + if crate::date::is_date_value(value) { + if self.write_reference_or_register(raw) { + return true; + } + self.out.push(TAG_DATE); + self.write_double(crate::date::js_date_get_time(value)); + return true; + } + if let Some(header) = unsafe { crate::value::addr_class::try_read_gc_header(raw) } { + if matches!( + header.obj_type, + crate::gc::GC_TYPE_ARRAY | crate::gc::GC_TYPE_LAZY_ARRAY + ) { + if self.depth >= MAX_DEPTH { + self.out.push(TAG_UNDEFINED); + return true; } - if crate::date::is_date_value(value) { - if self.write_reference_or_register(raw) { - return; - } - self.out.push(TAG_DATE); - self.write_double(crate::date::js_date_get_time(value)); - return; + if self.write_reference_or_register(raw) { + return true; } - if let Some(arr) = cp_array_ptr(value) { - if self.depth >= MAX_DEPTH { - self.out.push(TAG_UNDEFINED); - return; - } - if self.write_reference_or_register(raw) { - return; - } - self.write_dense_array(arr); - return; + self.write_dense_array(raw as *mut crate::array::ArrayHeader); + return true; + } + if header.obj_type == crate::gc::GC_TYPE_OBJECT { + if self.depth >= MAX_DEPTH { + self.out.push(TAG_UNDEFINED); + return true; } - if let Some(obj) = cp_object_ptr(value) { - if self.depth >= MAX_DEPTH { - self.out.push(TAG_UNDEFINED); - return; - } - if self.write_reference_or_register(raw) { - return; - } - self.write_object(obj); - return; + if self.write_reference_or_register(raw) { + return true; } + self.write_object(raw as *const ObjectHeader); + return true; } } - // Functions, symbols, unknown — degrade to undefined. - self.out.push(TAG_UNDEFINED); + false } fn write_number(&mut self, value: f64) { @@ -421,6 +450,17 @@ impl Serializer { self.write_host_view(NODE_BUFFER_VIEW_INDEX, bytes); } + fn write_array_buffer(&mut self, buffer: *const crate::buffer::BufferHeader) { + let len = unsafe { (*buffer).length as usize }; + let data = crate::buffer::buffer_data(buffer); + self.out.push(TAG_ARRAY_BUFFER); + self.write_varint(len as u64); + if !data.is_null() && len != 0 { + self.out + .extend_from_slice(unsafe { std::slice::from_raw_parts(data, len) }); + } + } + fn write_host_typed_array(&mut self, value: f64, kind: u8) { let raw = (value.to_bits() & crate::value::POINTER_MASK) as usize; let ta = raw as *const crate::typedarray::TypedArrayHeader; @@ -941,7 +981,15 @@ impl<'a> Deserializer<'a> { fn read_array_buffer(&mut self) -> Option { let len = self.read_varint()? as usize; let bytes = self.read_raw(len)?; - let v = cp_make_buffer(bytes); + let buffer = crate::buffer::js_array_buffer_new(len as i32); + if buffer.is_null() { + return Some(cp_undefined()); + } + let data = crate::buffer::buffer_data_mut(buffer); + if !data.is_null() && len != 0 { + unsafe { std::ptr::copy_nonoverlapping(bytes.as_ptr(), data, len) }; + } + let v = cp_box_ptr(buffer as *const u8); self.id_table.push(v); Some(v) } diff --git a/crates/perry-runtime/src/child_process/validate.rs b/crates/perry-runtime/src/child_process/validate.rs index 3a6e4cfb69..d1314fb3f5 100644 --- a/crates/perry-runtime/src/child_process/validate.rs +++ b/crates/perry-runtime/src/child_process/validate.rs @@ -176,6 +176,27 @@ fn cp_validate_options(value: f64, sync: bool, allow_null: bool) { crate::fs::validate::throw_type_error_with_code(&message, "ERR_INVALID_ARG_TYPE"); } + // Node normalizes stdio before the remaining spawn options. Preserve that + // precedence when a cumulative cluster.settings snapshot contains more + // than one invalid field. + let stdio = cp_get_field(value, b"stdio"); + if !cp_is_undefined(stdio) { + let mut ipc_count = 0; + if let Some(arr) = cp_array_ptr(stdio) { + for index in 0..unsafe { (*arr).length } { + cp_validate_stdio_entry( + crate::array::js_array_get_f64(arr, index), + sync, + true, + index as usize, + &mut ipc_count, + ); + } + } else { + cp_validate_stdio_entry(stdio, sync, false, 0, &mut ipc_count); + } + } + let required_string = ["cwd", "argv0"]; for name in required_string { let item = cp_get_field(value, name.as_bytes()); @@ -225,6 +246,34 @@ fn cp_validate_options(value: f64, sync: bool, allow_null: bool) { } } + for name in ["uid", "gid"] { + let item = cp_get_field(value, name.as_bytes()); + if !cp_is_undefined(item) && !cp_is_number(item) { + let message = format!( + "The \"options.{name}\" property must be of type number. Received {}", + crate::fs::validate::describe_received(item) + ); + crate::fs::validate::throw_type_error_with_code(&message, "ERR_INVALID_ARG_VALUE"); + } + } + + let inspect_port = cp_get_field(value, b"inspectPort"); + if JSValue::from_bits(inspect_port.to_bits()).is_null() { + crate::fs::validate::throw_range_error_named( + "The value of \"options.inspectPort\" is out of range", + "ERR_SOCKET_BAD_PORT", + ); + } + if !cp_is_undefined(inspect_port) + && !cp_is_number(inspect_port) + && crate::fs::extract_closure_ptr(inspect_port).is_null() + { + crate::fs::validate::throw_type_error_with_code( + "The \"options.inspectPort\" property must be of type number or function", + "ERR_INVALID_ARG_TYPE", + ); + } + let signal = cp_get_field(value, b"killSignal"); if !cp_is_undefined(signal) && !cp_signal_is_valid(signal) { crate::fs::validate::throw_type_error_with_code("Unknown signal", "ERR_UNKNOWN_SIGNAL"); @@ -243,25 +292,6 @@ fn cp_validate_options(value: f64, sync: bool, allow_null: bool) { "ERR_INVALID_ARG_VALUE", ); } - - let stdio = cp_get_field(value, b"stdio"); - if cp_is_undefined(stdio) { - return; - } - let mut ipc_count = 0; - if let Some(arr) = cp_array_ptr(stdio) { - for index in 0..unsafe { (*arr).length } { - cp_validate_stdio_entry( - crate::array::js_array_get_f64(arr, index), - sync, - true, - index as usize, - &mut ipc_count, - ); - } - } else { - cp_validate_stdio_entry(stdio, sync, false, 0, &mut ipc_count); - } } /// `new ChildProcess().spawn(options)` exposes the low-level Node constructor diff --git a/crates/perry-runtime/src/cluster.rs b/crates/perry-runtime/src/cluster.rs index d899e46914..39d43d34a8 100644 --- a/crates/perry-runtime/src/cluster.rs +++ b/crates/perry-runtime/src/cluster.rs @@ -13,7 +13,7 @@ use std::cell::RefCell; use std::collections::HashMap; -use std::sync::Once; +use std::sync::{Once, OnceLock}; use crate::array::ArrayHeader; use crate::closure::{js_closure_get_capture_f64, ClosureHeader}; @@ -54,6 +54,21 @@ thread_local! { } static CLUSTER_INIT: Once = Once::new(); +static CLUSTER_WORKER_ID: OnceLock> = OnceLock::new(); + +fn cluster_worker_id() -> Option<&'static str> { + CLUSTER_WORKER_ID + .get_or_init(|| { + let id = std::env::var("NODE_UNIQUE_ID") + .ok() + .filter(|value| !value.is_empty()); + if id.is_some() { + std::env::remove_var("NODE_UNIQUE_ID"); + } + id + }) + .as_deref() +} fn empty_object_value() -> f64 { let obj = js_object_alloc(0, 0); @@ -63,10 +78,7 @@ fn empty_object_value() -> f64 { /// True when this process was `cluster.fork()`ed — Node's convention is a /// non-empty `NODE_UNIQUE_ID` in the worker's environment. pub fn is_cluster_worker() -> bool { - std::env::var("NODE_UNIQUE_ID") - .ok() - .filter(|s| !s.is_empty()) - .is_some() + cluster_worker_id().is_some() } /// Bind a TCP listener with SO_REUSEPORT (+SO_REUSEADDR) so N cluster @@ -185,10 +197,8 @@ pub fn cluster_property(property: &str) -> Option { // EventEmitter.prototype, not as named module exports). Perry models the // default import as a distinct `cluster.default` native-module namespace whose // EventEmitter method reads resolve here; the namespace import keeps the -// `undefined` shape via `cluster_property`. Real worker-lifecycle events are -// still deferred (closed umbrella #3605) — this is module-level listener -// bookkeeping plus a synchronous `fork` emit so feature-detection and manual -// `emit()` round-trips match Node. +// `undefined` shape via `cluster_property`. The module-level emitter also +// receives the primary-side worker lifecycle events implemented below. // --------------------------------------------------------------------------- #[derive(Clone, Copy)] @@ -197,10 +207,20 @@ struct ClusterListener { once: bool, } -#[derive(Default)] struct ClusterEmitter { events: HashMap>, order: Vec, + max_listeners: f64, +} + +impl Default for ClusterEmitter { + fn default() -> Self { + Self { + events: HashMap::new(), + order: Vec::new(), + max_listeners: 10.0, + } + } } thread_local! { @@ -232,9 +252,7 @@ fn cluster_emitter_event_name(event: f64) -> Option { } fn cluster_register_listener(event: f64, listener: f64, once: bool, prepend: bool) -> f64 { - if !is_closure_value(listener) { - return cluster_default_value(); - } + crate::validators::validate_function(listener, "listener"); if let Some(name) = cluster_emitter_event_name(event) { CLUSTER_EMITTER.with(|emitter| { let mut emitter = emitter.borrow_mut(); @@ -364,17 +382,76 @@ pub extern "C" fn js_cluster_event_names() -> f64 { } #[no_mangle] -pub extern "C" fn js_cluster_listener_count(event: f64) -> f64 { +pub extern "C" fn js_cluster_listeners(event: f64) -> f64 { + ensure_cluster_runtime(); + let Some(name) = cluster_emitter_event_name(event) else { + return alloc_array_value(0); + }; + let callbacks = CLUSTER_EMITTER.with(|emitter| { + emitter + .borrow() + .events + .get(&name) + .map(|listeners| { + listeners + .iter() + .map(|listener| listener.callback_bits) + .collect::>() + }) + .unwrap_or_default() + }); + let mut arr = crate::array::js_array_alloc(callbacks.len() as u32); + for bits in callbacks { + arr = crate::array::js_array_push_f64(arr, f64::from_bits(bits)); + } + box_ptr(arr as *const u8) +} + +#[no_mangle] +pub extern "C" fn js_cluster_raw_listeners(event: f64) -> f64 { + // Perry stores once listeners directly and removes them before dispatch, + // so the raw and unwrapped listener arrays carry the same callable values. + js_cluster_listeners(event) +} + +#[no_mangle] +pub extern "C" fn js_cluster_set_max_listeners(value: f64) -> f64 { + ensure_cluster_runtime(); + let validated = crate::node_stream::validate_max_listeners(value); + CLUSTER_EMITTER.with(|emitter| emitter.borrow_mut().max_listeners = validated); + cluster_default_value() +} + +#[no_mangle] +pub extern "C" fn js_cluster_get_max_listeners() -> f64 { + ensure_cluster_runtime(); + CLUSTER_EMITTER.with(|emitter| emitter.borrow().max_listeners) +} + +#[no_mangle] +pub extern "C" fn js_cluster_listener_count(event: f64, listener: f64) -> f64 { ensure_cluster_runtime(); let Some(name) = cluster_emitter_event_name(event) else { return 0.0; }; + let listener_value = JSValue::from_bits(listener.to_bits()); + let filter_listener = !listener_value.is_undefined() && !listener_value.is_null(); CLUSTER_EMITTER.with(|emitter| { emitter .borrow() .events .get(&name) - .map(|l| l.len() as f64) + .map(|listeners| { + if filter_listener { + let bits = listener.to_bits(); + listeners + .iter() + .filter(|entry| entry.callback_bits == bits) + .count() as f64 + } else { + listeners.len() as f64 + } + }) .unwrap_or(0.0) }) } @@ -382,6 +459,7 @@ pub extern "C" fn js_cluster_listener_count(event: f64) -> f64 { #[no_mangle] pub extern "C" fn js_cluster_remove_listener(event: f64, listener: f64) -> f64 { ensure_cluster_runtime(); + crate::validators::validate_function(listener, "listener"); if let Some(name) = cluster_emitter_event_name(event) { let bits = listener.to_bits(); CLUSTER_EMITTER.with(|emitter| { @@ -401,13 +479,12 @@ pub extern "C" fn js_cluster_remove_listener(event: f64, listener: f64) -> f64 { } #[no_mangle] -pub extern "C" fn js_cluster_remove_all_listeners(event: f64) -> f64 { +pub extern "C" fn js_cluster_remove_all_listeners(event: f64, has_event: i32) -> f64 { ensure_cluster_runtime(); - let jv = JSValue::from_bits(event.to_bits()); - let target = if jv.is_undefined() || jv.is_null() { - None - } else { + let target = if has_event != 0 { cluster_emitter_event_name(event) + } else { + None }; CLUSTER_EMITTER.with(|emitter| { let mut emitter = emitter.borrow_mut(); @@ -425,6 +502,17 @@ pub extern "C" fn js_cluster_remove_all_listeners(event: f64) -> f64 { cluster_default_value() } +#[no_mangle] +pub extern "C" fn js_cluster_remove_all_listeners_args(args: *const ArrayHeader) -> f64 { + let has_event = !args.is_null() && crate::array::js_array_length(args) > 0; + let event = if has_event { + crate::array::js_array_get_f64(args, 0) + } else { + TAG_UNDEFINED_F64 + }; + js_cluster_remove_all_listeners(event, has_event as i32) +} + #[no_mangle] pub extern "C" fn js_cluster_setup_primary(settings: f64) -> f64 { ensure_cluster_runtime(); @@ -448,8 +536,13 @@ pub extern "C" fn js_cluster_fork(env: f64) -> f64 { } let args = get_field(settings, b"args"); + unsafe { + crate::child_process::js_child_process_validate_command(module, b"exec".as_ptr(), 4); + } + crate::child_process::js_child_process_validate_args(args); let args_ptr = array_ptr(args).map(|p| p as i64).unwrap_or(0); let opts = build_fork_options(settings, env); + crate::child_process::js_child_process_validate_options(box_ptr(opts as *const u8), 0, 1); let module_ptr = crate::string::js_string_materialize_to_heap(module) as i64; let worker = crate::child_process::fork::js_child_process_fork(module_ptr, args_ptr, opts as i64); @@ -466,9 +559,7 @@ pub extern "C" fn js_cluster_fork(env: f64) -> f64 { decorate_worker(worker, id); register_worker(id, worker); - // Node fires the cluster-level `fork` event synchronously when the worker - // object is created (`online`/`exit`/etc. remain deferred — #3605). - cluster_emit_event("fork", &[worker]); + defer_cluster_fork_event(worker); worker } @@ -524,11 +615,7 @@ fn cluster_root_scanner(visitor: &mut crate::gc::RuntimeRootVisitor<'_>) { } fn register_cluster_arities() { - let arities: [(*const u8, u32); 7] = [ - (worker_is_connected as *const u8, 0), - (worker_is_dead as *const u8, 0), - (worker_disconnect as *const u8, 0), - (worker_destroy as *const u8, 0), + let arities: [(*const u8, u32); 3] = [ (cluster_internal_online as *const u8, 0), (cluster_internal_disconnect as *const u8, 0), (cluster_internal_exit as *const u8, 2), @@ -536,6 +623,188 @@ fn register_cluster_arities() { for (func, arity) in arities { crate::closure::js_register_closure_arity(func, arity); } + for (func, arity, length) in [ + (cluster_worker_send as *const u8, 4, 0), + (cluster_worker_kill as *const u8, 1, 0), + (cluster_worker_destroy as *const u8, 1, 1), + (cluster_worker_disconnect as *const u8, 0, 0), + (cluster_worker_is_connected as *const u8, 0, 0), + (cluster_worker_is_dead as *const u8, 0, 0), + (cluster_setup_emit_thunk as *const u8, 0, 0), + (cluster_fork_emit_thunk as *const u8, 0, 0), + (cluster_callback_thunk as *const u8, 0, 0), + (cluster_disconnect_emit_thunk as *const u8, 0, 0), + (cluster_exit_emit_thunk as *const u8, 0, 0), + (cluster_internal_message as *const u8, 1, 1), + ] { + crate::closure::js_register_closure_arity(func, arity); + crate::closure::js_register_closure_length(func, length); + } +} + +pub(crate) fn ensure_worker_constructor(value: f64) -> f64 { + let raw = (value.to_bits() & crate::value::POINTER_MASK) as usize; + if raw == 0 || !crate::closure::is_closure_ptr(raw) { + return value; + } + let existing = crate::closure::closure_get_dynamic_prop(raw, "prototype"); + if object_ptr(existing).is_some() { + return value; + } + + ensure_cluster_runtime(); + let proto = js_object_alloc(0, 7); + let proto_value = box_ptr(proto as *const u8); + set_field(proto_value, b"constructor", value); + crate::object::set_builtin_property_attrs( + proto as usize, + "constructor".to_string(), + crate::object::PropertyAttrs::new(true, false, true), + ); + for (name, func, length) in [ + ("send", cluster_worker_send as *const u8, 0), + ("kill", cluster_worker_kill as *const u8, 0), + ("destroy", cluster_worker_destroy as *const u8, 1), + ("disconnect", cluster_worker_disconnect as *const u8, 0), + ("isConnected", cluster_worker_is_connected as *const u8, 0), + ("isDead", cluster_worker_is_dead as *const u8, 0), + ] { + let method = crate::closure::js_closure_alloc(func, 0); + crate::object::set_bound_native_closure_name(method, name); + crate::object::set_builtin_closure_length(method as usize, length); + set_field(proto_value, name.as_bytes(), box_ptr(method as *const u8)); + crate::object::set_builtin_property_attrs( + proto as usize, + name.to_string(), + crate::object::PropertyAttrs::new(true, false, true), + ); + } + + let event_emitter = crate::object::bound_native_callable_export_value("events", "EventEmitter"); + let event_emitter_id = crate::object::function_class_id(event_emitter); + let event_emitter_proto = + crate::object::ensure_function_prototype_object(event_emitter, event_emitter_id); + if !event_emitter_proto.is_null() { + crate::object::prototype_chain::object_set_static_prototype( + proto as usize, + box_ptr(event_emitter_proto as *const u8).to_bits(), + ); + } + + crate::closure::closure_set_dynamic_prop(raw, "prototype", proto_value); + crate::object::set_builtin_property_attrs( + raw, + "prototype".to_string(), + crate::object::PropertyAttrs::new(true, false, false), + ); + value +} + +fn worker_prototype_value() -> f64 { + let constructor = crate::object::bound_native_callable_export_value("cluster", "Worker"); + let constructor = ensure_worker_constructor(constructor); + let raw = (constructor.to_bits() & crate::value::POINTER_MASK) as usize; + crate::closure::closure_get_dynamic_prop(raw, "prototype") +} + +pub extern "C" fn js_cluster_worker_new(options: f64) -> f64 { + ensure_cluster_runtime(); + let worker = alloc_object_value(4); + set_field(worker, b"__clusterWorker", TAG_TRUE_F64); + let id = get_field(options, b"id"); + let state = get_field(options, b"state"); + let process = get_field(options, b"process"); + set_field( + worker, + b"id", + if JSValue::from_bits(id.to_bits()).is_undefined() { + 0.0 + } else { + id + }, + ); + set_field( + worker, + b"state", + if JSValue::from_bits(state.to_bits()).is_undefined() { + box_string("none") + } else { + state + }, + ); + if !JSValue::from_bits(process.to_bits()).is_undefined() { + set_field(worker, b"process", process); + } + if let (Some(obj), Some(proto)) = (object_ptr(worker), object_ptr(worker_prototype_value())) { + crate::object::prototype_chain::object_set_static_prototype( + obj as usize, + box_ptr(proto as *const u8).to_bits(), + ); + } + worker +} + +pub(crate) fn is_worker_instance_value(value: f64) -> bool { + get_field(value, b"__clusterWorker").to_bits() == TAG_TRUE_F64.to_bits() +} + +extern "C" fn cluster_worker_send( + _closure: *const ClosureHeader, + message: f64, + a2: f64, + a3: f64, + a4: f64, +) -> f64 { + if is_self_worker_value(crate::object::js_implicit_this_get()) { + return crate::process::process_ipc_send_call(message, a2, a3, a4); + } + crate::child_process::cp_method_send(std::ptr::null(), message, a2, a3, a4) +} + +extern "C" fn cluster_worker_kill(_closure: *const ClosureHeader, signal: f64) -> f64 { + if is_self_worker_value(crate::object::js_implicit_this_get()) { + let _ = crate::os::js_process_kill(crate::os::js_process_pid(), signal); + return TAG_UNDEFINED_F64; + } + let _ = crate::child_process::cp_method_kill(std::ptr::null(), signal); + TAG_UNDEFINED_F64 +} + +extern "C" fn cluster_worker_destroy(_closure: *const ClosureHeader, signal: f64) -> f64 { + cluster_worker_kill(std::ptr::null(), signal) +} + +extern "C" fn cluster_worker_disconnect(_closure: *const ClosureHeader) -> f64 { + let worker = crate::object::js_implicit_this_get(); + set_field(worker, b"exitedAfterDisconnect", TAG_TRUE_F64); + if is_self_worker_value(worker) { + let _ = crate::process::process_ipc_disconnect_call(); + } else { + let _ = crate::child_process::cp_method_disconnect(std::ptr::null()); + } + worker +} + +extern "C" fn cluster_worker_is_connected(_closure: *const ClosureHeader) -> f64 { + let worker = crate::object::js_implicit_this_get(); + let connected = if is_self_worker_value(worker) { + crate::process::ipc::process_ipc_property("connected").unwrap_or(TAG_FALSE_F64) + } else { + get_field(worker, b"connected") + }; + if connected.to_bits() == TAG_TRUE_F64.to_bits() { + TAG_TRUE_F64 + } else { + TAG_FALSE_F64 + } +} + +extern "C" fn cluster_worker_is_dead(_closure: *const ClosureHeader) -> f64 { + if is_worker_dead(crate::object::js_implicit_this_get()) { + TAG_TRUE_F64 + } else { + TAG_FALSE_F64 + } } fn settings_value() -> f64 { @@ -568,20 +837,45 @@ fn self_worker_value() -> f64 { if bits != 0 { return f64::from_bits(bits); } - let worker = alloc_object_value(3); - let id = std::env::var("NODE_UNIQUE_ID") - .ok() + let worker = alloc_object_value(8); + let id = cluster_worker_id() .and_then(|s| s.parse::().ok()) .unwrap_or(0.0); set_field(worker, b"id", id); - set_field(worker, b"process", TAG_UNDEFINED_F64); + set_field( + worker, + b"process", + crate::object::js_create_native_module_namespace(b"process".as_ptr(), 7), + ); + set_field(worker, b"state", box_string("online")); set_field(worker, b"exitedAfterDisconnect", TAG_FALSE_F64); + set_field(worker, b"__clusterWorker", TAG_TRUE_F64); + set_field(worker, b"__clusterState", box_string("online")); + if let (Some(obj), Some(proto)) = (object_ptr(worker), object_ptr(worker_prototype_value())) + { + crate::object::prototype_chain::object_set_static_prototype( + obj as usize, + box_ptr(proto as *const u8).to_bits(), + ); + } state.borrow_mut().self_worker_bits = worker.to_bits(); worker }) } +fn is_self_worker_value(value: f64) -> bool { + is_cluster_worker() + && CLUSTER_STATE.with(|state| state.borrow().self_worker_bits == value.to_bits()) +} + fn apply_setup_primary(settings_arg: f64) { + let policy = scheduling_policy_value(); + if !matches!(policy, 1 | 2) { + crate::fs::validate::throw_error_with_code( + "Bad cluster scheduling policy: must be one of SCHED_RR or SCHED_NONE", + "ERR_INTERNAL_ASSERTION", + ); + } let previous = settings_value(); let next = alloc_default_settings(); copy_object_fields(previous, next); @@ -592,6 +886,27 @@ fn apply_setup_primary(settings_arg: f64) { state.setup_called = true; state.settings_bits = next.to_bits(); }); + let deferred = crate::closure::js_closure_alloc(cluster_setup_emit_thunk as *const u8, 1); + crate::closure::js_closure_set_capture_f64(deferred, 0, next); + crate::timer::js_set_immediate_callback(deferred as i64); +} + +extern "C" fn cluster_setup_emit_thunk(closure: *const ClosureHeader) -> f64 { + let settings = js_closure_get_capture_f64(closure, 0); + cluster_emit_event("setup", &[settings]); + TAG_UNDEFINED_F64 +} + +fn defer_cluster_fork_event(worker: f64) { + let deferred = crate::closure::js_closure_alloc(cluster_fork_emit_thunk as *const u8, 1); + crate::closure::js_closure_set_capture_f64(deferred, 0, worker); + crate::timer::js_set_immediate_callback(deferred as i64); +} + +extern "C" fn cluster_fork_emit_thunk(closure: *const ClosureHeader) -> f64 { + let worker = js_closure_get_capture_f64(closure, 0); + cluster_emit_event("fork", &[worker]); + TAG_UNDEFINED_F64 } fn alloc_default_settings() -> f64 { @@ -669,6 +984,7 @@ fn build_fork_options(settings: f64, env_arg: f64) -> *mut ObjectHeader { copy_setting_to_option(settings, opts_val, b"cwd"); copy_setting_to_option(settings, opts_val, b"execArgv"); copy_setting_to_option(settings, opts_val, b"execPath"); + copy_setting_to_option(settings, opts_val, b"inspectPort"); copy_setting_to_option(settings, opts_val, b"serialization"); copy_setting_to_option(settings, opts_val, b"silent"); copy_setting_to_option(settings, opts_val, b"stdio"); @@ -709,6 +1025,10 @@ fn build_worker_env(env_arg: f64, worker_id: u32) -> f64 { /// SCHED_RR (2); a user assignment lands in the native-namespace override table /// (CJS exports are mutable), so we read it back from there. SCHED_NONE is `1`. fn scheduling_policy_is_sched_none() -> bool { + scheduling_policy_value() == 1 +} + +fn scheduling_policy_value() -> i32 { for module in ["cluster", "cluster.default"] { if let Some(v) = crate::object::native_namespace_prop_override_get(module, "schedulingPolicy") @@ -719,41 +1039,37 @@ fn scheduling_policy_is_sched_none() -> bool { } else { v as i32 }; - return n == 1; + return n; } } - false + 2 } fn decorate_worker(worker: f64, id: u32) { set_field(worker, b"id", id as f64); set_field(worker, b"process", worker); - set_field(worker, b"exitedAfterDisconnect", TAG_FALSE_F64); + set_field(worker, b"state", box_string("none")); + set_field(worker, b"exitedAfterDisconnect", TAG_UNDEFINED_F64); set_field(worker, b"__clusterWorker", TAG_TRUE_F64); - set_field(worker, b"__clusterState", box_string("online")); - - let original_disconnect = get_field(worker, b"disconnect"); - set_field(worker, b"__clusterDisconnect", original_disconnect); - set_field( - worker, - b"isConnected", - closure_value(worker_is_connected as *const u8, worker), - ); - set_field( - worker, - b"isDead", - closure_value(worker_is_dead as *const u8, worker), - ); - set_field( - worker, - b"disconnect", - closure_value(worker_disconnect as *const u8, worker), - ); - set_field( - worker, - b"destroy", - closure_value(worker_destroy as *const u8, worker), - ); + set_field(worker, b"__clusterState", box_string("none")); + + if let (Some(obj), Some(proto)) = (object_ptr(worker), object_ptr(worker_prototype_value())) { + crate::object::prototype_chain::object_set_static_prototype( + obj as usize, + box_ptr(proto as *const u8).to_bits(), + ); + for name in [ + "send", + "kill", + "destroy", + "disconnect", + "isConnected", + "isDead", + ] { + let key = str_key(name.as_bytes()); + js_object_delete_field(obj, key); + } + } register_listener( worker, @@ -770,6 +1086,11 @@ fn decorate_worker(worker: f64, id: u32) { "exit", closure_value(cluster_internal_exit as *const u8, worker), ); + register_listener( + worker, + "message", + closure_value(cluster_internal_message as *const u8, worker), + ); } pub(crate) fn consume_internal_message(worker: f64, message: f64) -> bool { @@ -876,6 +1197,7 @@ fn mark_worker_online(worker: f64) { } set_field(worker, b"__clusterOnlineEmitted", TAG_TRUE_F64); set_field(worker, b"__clusterState", box_string("online")); + set_field(worker, b"state", box_string("online")); emit(worker, "online", &[]); cluster_emit_event("online", &[worker]); } @@ -929,50 +1251,18 @@ fn drain_disconnect_callbacks_if_idle() { }); for bits in callbacks { - let cb = f64::from_bits(bits); - unsafe { - let _ = crate::closure::js_native_call_value(cb, std::ptr::null(), 0); - } - } -} - -extern "C" fn worker_is_connected(closure: *const ClosureHeader) -> f64 { - let worker = closure_this(closure); - let connected = get_field(worker, b"connected"); - if JSValue::from_bits(connected.to_bits()).is_bool() - && connected.to_bits() == crate::value::TAG_TRUE - { - TAG_TRUE_F64 - } else { - TAG_FALSE_F64 - } -} - -extern "C" fn worker_is_dead(closure: *const ClosureHeader) -> f64 { - if is_worker_dead(closure_this(closure)) { - TAG_TRUE_F64 - } else { - TAG_FALSE_F64 + let deferred = crate::closure::js_closure_alloc(cluster_callback_thunk as *const u8, 1); + crate::closure::js_closure_set_capture_ptr(deferred, 0, bits as i64); + crate::timer::js_set_immediate_callback(deferred as i64); } } -extern "C" fn worker_disconnect(closure: *const ClosureHeader) -> f64 { - let worker = closure_this(closure); - set_field(worker, b"exitedAfterDisconnect", TAG_TRUE_F64); - call_original_disconnect(worker); - worker -} - -extern "C" fn worker_destroy(closure: *const ClosureHeader) -> f64 { - let worker = closure_this(closure); - set_field(worker, b"exitedAfterDisconnect", TAG_TRUE_F64); - let kill = get_field(worker, b"kill"); - if is_closure_value(kill) { - unsafe { - let _ = crate::closure::js_native_call_value(kill, std::ptr::null(), 0); - } +extern "C" fn cluster_callback_thunk(closure: *const ClosureHeader) -> f64 { + let callback = f64::from_bits(crate::closure::js_closure_get_capture_ptr(closure, 0) as u64); + unsafe { + let _ = crate::closure::js_native_call_value(callback, std::ptr::null(), 0); } - worker + TAG_UNDEFINED_F64 } extern "C" fn cluster_internal_online(closure: *const ClosureHeader) -> f64 { @@ -985,6 +1275,11 @@ extern "C" fn cluster_internal_disconnect(closure: *const ClosureHeader) -> f64 TAG_UNDEFINED_F64 } +extern "C" fn cluster_internal_message(closure: *const ClosureHeader, message: f64) -> f64 { + cluster_emit_event("message", &[closure_this(closure), message]); + TAG_UNDEFINED_F64 +} + /// Emit `disconnect` once per worker on both the worker object and the /// cluster — fed by the child reactor's IPC-channel-close event, with an /// exit-time fallback so the Node-mandated disconnect→exit order holds even @@ -995,12 +1290,25 @@ fn mark_worker_disconnected(worker: f64) { return; } set_field(worker, b"__clusterDisconnectEmitted", TAG_TRUE_F64); + set_field(worker, b"state", box_string("disconnected")); + let deferred = crate::closure::js_closure_alloc(cluster_disconnect_emit_thunk as *const u8, 1); + crate::closure::js_closure_set_capture_f64(deferred, 0, worker); + crate::timer::js_set_immediate_callback(deferred as i64); +} + +extern "C" fn cluster_disconnect_emit_thunk(closure: *const ClosureHeader) -> f64 { + let worker = js_closure_get_capture_f64(closure, 0); cluster_emit_event("disconnect", &[worker]); + TAG_UNDEFINED_F64 } extern "C" fn cluster_internal_exit(closure: *const ClosureHeader, code: f64, signal: f64) -> f64 { let worker = closure_this(closure); + if JSValue::from_bits(get_field(worker, b"exitedAfterDisconnect").to_bits()).is_undefined() { + set_field(worker, b"exitedAfterDisconnect", TAG_FALSE_F64); + } set_field(worker, b"__clusterState", box_string("dead")); + set_field(worker, b"state", box_string("dead")); mark_worker_disconnected(worker); // #4962 — drop the dead worker from any SCHED_RR rotation it joined. #[cfg(unix)] @@ -1008,8 +1316,36 @@ extern "C" fn cluster_internal_exit(closure: *const ClosureHeader, code: f64, si crate::cluster_sched::primary_remove_worker(handle); } remove_worker(worker); + let deferred = crate::closure::js_closure_alloc(cluster_exit_emit_thunk as *const u8, 3); + crate::closure::js_closure_set_capture_f64(deferred, 0, worker); + crate::closure::js_closure_set_capture_f64(deferred, 1, code); + crate::closure::js_closure_set_capture_f64(deferred, 2, signal); + crate::timer::js_set_immediate_callback(deferred as i64); + invoke_disconnect_callbacks_if_idle(); + TAG_UNDEFINED_F64 +} + +fn invoke_disconnect_callbacks_if_idle() { + let callbacks = CLUSTER_STATE.with(|state| { + let mut state = state.borrow_mut(); + if !state.worker_bits_by_id.is_empty() { + return Vec::new(); + } + std::mem::take(&mut state.disconnect_callbacks) + }); + for bits in callbacks { + let callback = f64::from_bits(bits); + unsafe { + let _ = crate::closure::js_native_call_value(callback, std::ptr::null(), 0); + } + } +} + +extern "C" fn cluster_exit_emit_thunk(closure: *const ClosureHeader) -> f64 { + let worker = js_closure_get_capture_f64(closure, 0); + let code = js_closure_get_capture_f64(closure, 1); + let signal = js_closure_get_capture_f64(closure, 2); cluster_emit_event("exit", &[worker, code, signal]); - drain_disconnect_callbacks_if_idle(); TAG_UNDEFINED_F64 } @@ -1022,16 +1358,6 @@ fn call_worker_disconnect(worker: f64) { } } -fn call_original_disconnect(worker: f64) { - let original = get_field(worker, b"__clusterDisconnect"); - if !is_closure_value(original) { - return; - } - unsafe { - let _ = crate::closure::js_native_call_value(original, std::ptr::null(), 0); - } -} - fn is_worker_dead(worker: f64) -> bool { let exit_code = get_field(worker, b"exitCode"); let signal_code = get_field(worker, b"signalCode"); @@ -1039,6 +1365,12 @@ fn is_worker_dead(worker: f64) -> bool { || !JSValue::from_bits(signal_code.to_bits()).is_null() } +pub(crate) fn exit_worker_on_unexpected_primary_disconnect(unexpected: bool) { + if unexpected && is_cluster_worker() { + std::process::exit(0); + } +} + fn closure_this(closure: *const ClosureHeader) -> f64 { if closure.is_null() { TAG_UNDEFINED_F64 diff --git a/crates/perry-runtime/src/cluster_sched.rs b/crates/perry-runtime/src/cluster_sched.rs index 274fd664e0..5fbed9cd71 100644 --- a/crates/perry-runtime/src/cluster_sched.rs +++ b/crates/perry-runtime/src/cluster_sched.rs @@ -32,12 +32,13 @@ use std::sync::{Condvar, Mutex}; #[cfg(unix)] use std::time::Duration; -// Binary fd-frame on the IPC byte stream: a leading NUL (never present in the -// newline-delimited JSON the channel otherwise carries) followed by the 4-byte -// big-endian routing key id. The accompanying fd rides in the `SCM_RIGHTS` -// ancillary data of the same `sendmsg`. +// Binary fd-frame on the IPC byte stream: a reserved 0xFE tag followed by the +// 4-byte big-endian routing key id. JSON never starts with this byte, and an +// advanced IPC length prefix cannot start with it for a representable payload +// in this process. The accompanying fd rides in the `SCM_RIGHTS` ancillary +// data of the same `sendmsg`. #[cfg(unix)] -const FD_FRAME_TAG: u8 = 0x00; +const FD_FRAME_TAG: u8 = 0xFE; #[cfg(unix)] const FD_FRAME_LEN: usize = 5; // tag + u32 key id @@ -188,14 +189,16 @@ pub fn worker_recv_fd(key_id: u32) -> RawFd { /// Worker IPC read loop for cluster workers. Replaces the plain /// `BufReader::lines()` reader (which cannot receive `SCM_RIGHTS`) with a -/// `recvmsg` loop that splits the byte stream into JSON frames *and* routes -/// ancillary connection fds. `queryServerReply` frames are consumed here; -/// every other JSON line is handed to `on_message` exactly as before so normal +/// `recvmsg` loop that splits the byte stream into JSON or advanced frames and +/// routes ancillary connection fds. `queryServerReply` frames are consumed +/// here; every other frame is handed to the matching callback so normal /// `process.on('message')` delivery is unaffected. #[cfg(unix)] pub fn worker_recv_loop( stream: std::os::unix::net::UnixStream, + advanced: bool, mut on_message: impl FnMut(String), + mut on_advanced: impl FnMut(Vec), on_closed: impl FnOnce(), ) { let fd = stream.as_raw_fd(); @@ -246,7 +249,11 @@ pub fn worker_recv_loop( } acc.extend_from_slice(&data[..n as usize]); - drain_frames(&mut acc, &mut fd_queue, &mut on_message); + if advanced { + drain_advanced_frames(&mut acc, &mut fd_queue, &mut on_advanced); + } else { + drain_json_frames(&mut acc, &mut fd_queue, &mut on_message); + } } // Channel closed: flag it + drain any already-routed-but-unpulled fds under @@ -276,7 +283,7 @@ pub fn worker_recv_loop( /// queued fd to its injection queue; newline-delimited JSON lines are either /// consumed (queryServerReply) or forwarded to `on_message`. #[cfg(unix)] -fn drain_frames( +fn drain_json_frames( acc: &mut Vec, fd_queue: &mut VecDeque, on_message: &mut impl FnMut(String), @@ -322,6 +329,198 @@ fn drain_frames( } } +/// Split an advanced channel into `[u32 BE length][V8 payload]` frames while +/// still recognizing the reserved fd-passing frame. Query-server replies are +/// decoded without constructing JS heap values on this reader thread because +/// the JS main thread is synchronously waiting for that reply inside `listen()`. +#[cfg(unix)] +fn drain_advanced_frames( + acc: &mut Vec, + fd_queue: &mut VecDeque, + on_message: &mut impl FnMut(Vec), +) { + let mut pos = 0usize; + loop { + if pos >= acc.len() { + break; + } + if acc[pos] == FD_FRAME_TAG { + if acc.len() - pos < FD_FRAME_LEN { + break; + } + let key_id = + u32::from_be_bytes([acc[pos + 1], acc[pos + 2], acc[pos + 3], acc[pos + 4]]); + if let Some(cfd) = fd_queue.pop_front() { + let mut inbox = worker_state().fd_inbox.lock().unwrap(); + inbox.queues.entry(key_id).or_default().push_back(cfd); + drop(inbox); + worker_state().fd_cv.notify_all(); + } + pos += FD_FRAME_LEN; + continue; + } + if acc.len() - pos < 4 { + break; + } + let len = u32::from_be_bytes([acc[pos], acc[pos + 1], acc[pos + 2], acc[pos + 3]]) as usize; + if acc.len() - pos - 4 < len { + break; + } + let start = pos + 4; + let payload = &acc[start..start + len]; + if !try_consume_advanced_query_reply(payload) { + on_message(payload.to_vec()); + } + pos = start + len; + } + if pos > 0 { + acc.drain(..pos); + } +} + +/// Extract the three fields in Perry's flat advanced `queryServerReply` +/// object without constructing JS heap values on the IPC reader thread. +#[cfg(unix)] +fn try_consume_advanced_query_reply(payload: &[u8]) -> bool { + let mut reader = FlatV8Reader::new(payload); + let Some(fields) = reader.read_object() else { + return false; + }; + if fields.act.as_deref() != Some("queryServerReply") { + return false; + } + let Some(req_key) = fields.req_key else { + return false; + }; + let port = fields.port.unwrap_or(-1); + let mut queries = worker_state().queries.lock().unwrap(); + if let Some(slot) = queries.get_mut(&req_key) { + *slot = Some(port); + drop(queries); + worker_state().query_cv.notify_all(); + } + true +} + +#[cfg(unix)] +#[derive(Default)] +struct QueryReplyFields { + act: Option, + req_key: Option, + port: Option, +} + +/// Tiny, allocation-safe V8 reader for a flat object whose keys/values are +/// strings or numbers. This intentionally is not a second structured-clone +/// implementation; it only recognizes the cluster control message above. +#[cfg(unix)] +struct FlatV8Reader<'a> { + bytes: &'a [u8], + pos: usize, +} + +#[cfg(unix)] +impl<'a> FlatV8Reader<'a> { + fn new(bytes: &'a [u8]) -> Self { + Self { bytes, pos: 0 } + } + + fn byte(&mut self) -> Option { + let value = *self.bytes.get(self.pos)?; + self.pos += 1; + Some(value) + } + + fn varint(&mut self) -> Option { + let mut value = 0u64; + let mut shift = 0; + loop { + let byte = self.byte()?; + value |= ((byte & 0x7f) as u64) << shift; + if byte & 0x80 == 0 { + return Some(value); + } + shift += 7; + if shift >= 64 { + return None; + } + } + } + + fn raw(&mut self, len: usize) -> Option<&'a [u8]> { + let end = self.pos.checked_add(len)?; + let bytes = self.bytes.get(self.pos..end)?; + self.pos = end; + Some(bytes) + } + + fn string(&mut self) -> Option { + match self.byte()? { + b'S' | b'"' => { + let len = self.varint()? as usize; + Some(String::from_utf8_lossy(self.raw(len)?).into_owned()) + } + b'c' => { + let len = self.varint()? as usize; + let bytes = self.raw(len)?; + if bytes.len() % 2 != 0 { + return None; + } + let units = bytes + .chunks_exact(2) + .map(|pair| u16::from_le_bytes([pair[0], pair[1]])) + .collect::>(); + String::from_utf16(&units).ok() + } + _ => None, + } + } + + fn number(&mut self) -> Option { + match self.byte()? { + b'I' => { + let value = self.varint()?; + Some(((value >> 1) as i64 ^ -((value & 1) as i64)) as i32) + } + b'U' => Some(self.varint()? as i32), + b'N' => { + let mut bytes = [0u8; 8]; + bytes.copy_from_slice(self.raw(8)?); + Some(f64::from_le_bytes(bytes) as i32) + } + _ => None, + } + } + + fn read_object(&mut self) -> Option { + if self.byte()? != 0xff { + return None; + } + self.varint()?; + if self.byte()? != b'o' { + return None; + } + let mut fields = QueryReplyFields::default(); + loop { + if self.bytes.get(self.pos) == Some(&b'{') { + self.pos += 1; + let _ = self.varint()?; + return Some(fields); + } + let key = self.string()?; + match key.as_str() { + "cmd" => { + let _ = self.string()?; + } + "act" => fields.act = Some(self.string()?), + "reqKey" => fields.req_key = Some(self.string()?), + "port" => fields.port = Some(self.number()?), + _ => return None, + } + } + } +} + /// Resolve a `queryServerReply` line (lightweight scan — runs on the reader /// thread, so no JS-runtime parsing). Returns true if the line was a reply. #[cfg(unix)] diff --git a/crates/perry-runtime/src/object/class_registry.rs b/crates/perry-runtime/src/object/class_registry.rs index 8d5e7a36e4..c2859ae85f 100644 --- a/crates/perry-runtime/src/object/class_registry.rs +++ b/crates/perry-runtime/src/object/class_registry.rs @@ -110,9 +110,9 @@ pub use prototype_methods::{ // ── construct.rs ──────────────────────────────────────────────────────────── pub(crate) use construct::{ extends_target_must_throw, function_would_have_own_prototype, is_callable_function_value, - js_value_is_constructor, lookup_prototype_method, nm_ctor_child_process, nm_ctor_fs, - nm_ctor_readline, nm_ctor_repl, nm_ctor_stream, nm_ctor_tls, nm_ctor_tty, nm_ctor_vm, - nm_ctor_wasi, ordinary_function_prototype_value_for_read, promise_parent_in_chain, + js_value_is_constructor, lookup_prototype_method, nm_ctor_child_process, nm_ctor_cluster, + nm_ctor_fs, nm_ctor_readline, nm_ctor_repl, nm_ctor_stream, nm_ctor_tls, nm_ctor_tty, + nm_ctor_vm, nm_ctor_wasi, ordinary_function_prototype_value_for_read, promise_parent_in_chain, }; pub use construct::{ js_ctor_return_override, js_function_prototype_value_for_read, js_new_function_construct, diff --git a/crates/perry-runtime/src/object/class_registry/construct.rs b/crates/perry-runtime/src/object/class_registry/construct.rs index 8c40a5d40d..69c85162a0 100644 --- a/crates/perry-runtime/src/object/class_registry/construct.rs +++ b/crates/perry-runtime/src/object/class_registry/construct.rs @@ -69,6 +69,16 @@ pub(crate) unsafe fn nm_ctor_child_process( (method == "ChildProcess").then(crate::child_process::cp_build_unstarted_child_process) } +pub(crate) unsafe fn nm_ctor_cluster( + _module: &str, + method: &str, + args_ptr: *const f64, + args_len: usize, +) -> Option { + (method == "Worker") + .then(|| crate::cluster::js_cluster_worker_new(nm_ctor_arg(args_ptr, args_len, 0))) +} + pub(crate) unsafe fn nm_ctor_fs( _module: &str, method: &str, diff --git a/crates/perry-runtime/src/object/instanceof.rs b/crates/perry-runtime/src/object/instanceof.rs index f2917c9287..cff315e217 100644 --- a/crates/perry-runtime/src/object/instanceof.rs +++ b/crates/perry-runtime/src/object/instanceof.rs @@ -301,12 +301,17 @@ pub extern "C" fn js_instanceof_dynamic(value: f64, type_ref: f64) -> f64 { { return f64::from_bits(crate::value::TAG_TRUE); } - if module == "events" - && method == "EventEmitter" - && (is_event_emitter_instance_value(value) - || super::tls_constructor_prototype_is_instance_of(value, method.as_str())) - { - return f64::from_bits(crate::value::TAG_TRUE); + if module == "events" && method == "EventEmitter" { + return f64::from_bits( + if is_event_emitter_instance_value(value) + || super::tls_constructor_prototype_is_instance_of(value, method.as_str()) + || ordinary_has_instance_prototype_walk(value, type_ref) + { + crate::value::TAG_TRUE + } else { + TAG_FALSE + }, + ); } if module == "events" && method == "EventEmitterAsyncResource" @@ -823,6 +828,11 @@ fn dispatch_own_has_instance(cb: f64, value: f64) -> HasInstanceOutcome { } fn is_event_emitter_instance_value(value: f64) -> bool { + if is_native_module_namespace_value(value, "cluster.default") + || crate::cluster::is_worker_instance_value(value) + { + return true; + } if let Some(handle) = small_native_handle_id(value) { if let Some(probe) = crate::object::event_emitter_handle_probe() { return unsafe { probe(handle) }; @@ -835,7 +845,8 @@ fn is_event_emitter_instance_value(value: f64) -> bool { { return true; } - false + let constructor = crate::object::bound_native_callable_export_value("events", "EventEmitter"); + ordinary_has_instance_prototype_walk(value, constructor) } fn is_event_emitter_async_resource_instance_value(value: f64) -> bool { diff --git a/crates/perry-runtime/src/object/native_module.rs b/crates/perry-runtime/src/object/native_module.rs index 8810be7b61..cf21885c54 100644 --- a/crates/perry-runtime/src/object/native_module.rs +++ b/crates/perry-runtime/src/object/native_module.rs @@ -624,6 +624,7 @@ pub(crate) fn canonical_native_callable_property<'a>( ("path" | "path.posix" | "path.win32", "_makeLong") => "toNamespacedPath", ("querystring", "decode") => "parse", ("querystring", "encode") => "stringify", + ("cluster", "setupMaster") => "setupPrimary", _ => property_name, } } diff --git a/crates/perry-runtime/src/object/native_module/callable_exports.rs b/crates/perry-runtime/src/object/native_module/callable_exports.rs index 801dcf11be..f73aaf40f2 100644 --- a/crates/perry-runtime/src/object/native_module/callable_exports.rs +++ b/crates/perry-runtime/src/object/native_module/callable_exports.rs @@ -216,6 +216,10 @@ pub(crate) fn is_cluster_emitter_method(prop: &str) -> bool { | "emit" | "eventNames" | "listenerCount" + | "listeners" + | "rawListeners" + | "getMaxListeners" + | "setMaxListeners" ) } @@ -233,6 +237,7 @@ fn native_callable_export_arity_reference(module: &str, prop: &str) -> Option Some(1), ("cluster", "emit") => Some(1), ("cluster", "eventNames") => Some(0), + ("cluster", "getMaxListeners") => Some(0), ( "cluster", "on" @@ -244,6 +249,7 @@ fn native_callable_export_arity_reference(module: &str, prop: &str) -> Option Some(2), + ("cluster", "listeners" | "rawListeners" | "setMaxListeners") => Some(1), ("cluster", "removeAllListeners") => Some(1), // #6563: node-pty `spawn(file, args, options)`. ("node-pty", "spawn") => Some(3), @@ -1883,6 +1889,18 @@ pub(crate) unsafe fn nm_attach_child_process( } } +pub(crate) unsafe fn nm_attach_cluster( + property_name: &str, + value: f64, + _closure_addr: usize, +) -> f64 { + if property_name == "Worker" { + crate::cluster::ensure_worker_constructor(value) + } else { + value + } +} + #[allow(unused_mut)] pub(crate) unsafe fn nm_attach_sqlite( property_name: &str, @@ -2128,14 +2146,18 @@ static CALLABLE_EXPORT_ARITY_TABLE: &[(&str, &[(&str, u32)])] = &[ ("emit", 1), ("eventNames", 0), ("fork", 1), + ("getMaxListeners", 0), ("listenerCount", 2), + ("listeners", 1), ("off", 2), ("on", 2), ("once", 2), ("prependListener", 2), ("prependOnceListener", 2), + ("rawListeners", 1), ("removeAllListeners", 1), ("removeListener", 2), + ("setMaxListeners", 1), ("setupMaster", 1), ("setupPrimary", 1), ], diff --git a/crates/perry-runtime/src/object/native_module_dispatch/dispatch_a_c.rs b/crates/perry-runtime/src/object/native_module_dispatch/dispatch_a_c.rs index b56233ef73..344f44ae37 100644 --- a/crates/perry-runtime/src/object/native_module_dispatch/dispatch_a_c.rs +++ b/crates/perry-runtime/src/object/native_module_dispatch/dispatch_a_c.rs @@ -527,12 +527,16 @@ pub(crate) unsafe fn nm_dispatch_cluster(ctx: &NmCtx, module_name: &str, method_ } ("cluster", "emit") => crate::cluster::js_cluster_emit(arg(0), pack_args_from(1)), ("cluster", "eventNames") => crate::cluster::js_cluster_event_names(), - ("cluster", "listenerCount") => crate::cluster::js_cluster_listener_count(arg(0)), + ("cluster", "listeners") => crate::cluster::js_cluster_listeners(arg(0)), + ("cluster", "rawListeners") => crate::cluster::js_cluster_raw_listeners(arg(0)), + ("cluster", "setMaxListeners") => crate::cluster::js_cluster_set_max_listeners(arg(0)), + ("cluster", "getMaxListeners") => crate::cluster::js_cluster_get_max_listeners(), + ("cluster", "listenerCount") => crate::cluster::js_cluster_listener_count(arg(0), arg(1)), ("cluster", "removeListener") | ("cluster", "off") => { crate::cluster::js_cluster_remove_listener(arg(0), arg(1)) } ("cluster", "removeAllListeners") => { - crate::cluster::js_cluster_remove_all_listeners(arg(0)) + crate::cluster::js_cluster_remove_all_listeners(arg(0), (args_len > 0) as i32) } // #1577: captured-then-called crypto methods (`const f = diff --git a/crates/perry-runtime/src/object/native_module_registry.rs b/crates/perry-runtime/src/object/native_module_registry.rs index a078a9ef3e..f1eb5642c7 100644 --- a/crates/perry-runtime/src/object/native_module_registry.rs +++ b/crates/perry-runtime/src/object/native_module_registry.rs @@ -252,10 +252,15 @@ pub extern "C" fn js_nm_install_child_process() { } #[no_mangle] pub extern "C" fn js_nm_install_cluster() { + nm_register_attach( + NmBucket::Cluster, + super::native_module::callable_exports::nm_attach_cluster, + ); NM_DISPATCH_REGISTRY[NmBucket::Cluster as usize].store( nm_dispatch_cluster as NmDispatchFn as *mut (), Ordering::Relaxed, ); + nm_register_ctor(NmBucket::Cluster, nm_ctor_cluster); } #[no_mangle] pub extern "C" fn js_nm_install_console() { @@ -610,8 +615,8 @@ pub(crate) fn nm_run_install_all_hook() { // the method-dispatch registry: populated by `js_nm_install_()` (only // the 8 ctor-owning buckets register a fn), looked up by `js_new_function_construct`. use super::class_registry::{ - nm_ctor_child_process, nm_ctor_fs, nm_ctor_readline, nm_ctor_repl, nm_ctor_stream, nm_ctor_tls, - nm_ctor_tty, nm_ctor_vm, nm_ctor_wasi, + nm_ctor_child_process, nm_ctor_cluster, nm_ctor_fs, nm_ctor_readline, nm_ctor_repl, + nm_ctor_stream, nm_ctor_tls, nm_ctor_tty, nm_ctor_vm, nm_ctor_wasi, }; type NmCtorFn = unsafe fn(&str, &str, *const f64, usize) -> Option; diff --git a/crates/perry-runtime/src/process.rs b/crates/perry-runtime/src/process.rs index dd5a5219ca..480975c34f 100644 --- a/crates/perry-runtime/src/process.rs +++ b/crates/perry-runtime/src/process.rs @@ -10,6 +10,7 @@ use crate::string::{js_string_from_bytes, StringHeader}; use crate::value::JSValue; use std::cell::{Cell, RefCell}; use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::OnceLock; mod credentials; mod env_misc; @@ -342,6 +343,17 @@ pub(crate) fn module_array_value(items: &[&str]) -> f64 { f64::from_bits(JSValue::array_ptr(arr).bits()) } +fn process_exec_argv_value() -> f64 { + static EXEC_ARGV: OnceLock> = OnceLock::new(); + let values = EXEC_ARGV.get_or_init(|| { + let raw = std::env::var("PERRY_PROCESS_EXEC_ARGV").unwrap_or_else(|_| "[]".to_string()); + std::env::remove_var("PERRY_PROCESS_EXEC_ARGV"); + serde_json::from_str(&raw).unwrap_or_default() + }); + let refs = values.iter().map(String::as_str).collect::>(); + module_array_value(&refs) +} + pub(crate) fn module_set_value(items: &[&str]) -> f64 { let mut set = crate::set::js_set_alloc(items.len() as u32); for item in items { @@ -690,7 +702,8 @@ pub fn process_metadata_property(property: &str) -> Option { "argv0" | "execPath" => module_string_value(&process_argv0_string()), "config" => report::process_config_value(), "debugPort" => 9229.0, - "execArgv" | "moduleLoadList" => module_array_value(&[]), + "execArgv" => process_exec_argv_value(), + "moduleLoadList" => module_array_value(&[]), "features" => report::process_features_value(), "finalization" => finalization::process_finalization_value(), "permission" => permission::process_permission_value()?, diff --git a/crates/perry-runtime/src/process/ipc.rs b/crates/perry-runtime/src/process/ipc.rs index b090e40a28..551ad7c93b 100644 --- a/crates/perry-runtime/src/process/ipc.rs +++ b/crates/perry-runtime/src/process/ipc.rs @@ -2,7 +2,8 @@ //! //! Node initializes this from `NODE_CHANNEL_FD` during bootstrap. Perry adopts //! the same inherited Unix fd convention used by its `child_process.fork()` -//! parent side and speaks newline-delimited JSON frames for this cut. +//! parent side and speaks either newline-delimited JSON or Node's advanced +//! `[u32 BE length][V8 payload]` frames. #![cfg_attr(not(feature = "proc-ipc"), allow(dead_code))] @@ -16,7 +17,7 @@ use std::collections::VecDeque; use std::sync::{Mutex, MutexGuard, OnceLock}; #[cfg(unix)] -use std::io::{BufRead, Write}; +use std::io::{BufRead, Read, Write}; #[cfg(unix)] use std::os::fd::FromRawFd; #[cfg(unix)] @@ -24,6 +25,7 @@ use std::os::unix::net::UnixStream; enum IpcEvent { Message(String), + MessageAdvanced(Vec), Closed, } @@ -33,6 +35,7 @@ struct ChildIpcState { connected: bool, refed: bool, disconnect_emitted: bool, + advanced: bool, #[cfg(unix)] send: Option, queue: VecDeque, @@ -46,6 +49,7 @@ impl ChildIpcState { connected: false, refed: false, disconnect_emitted: false, + advanced: false, #[cfg(unix)] send: None, queue: VecDeque::new(), @@ -138,9 +142,6 @@ pub(crate) fn process_ipc_ensure_initialized() { #[cfg(unix)] fn initialize_unix_ipc(fd_var: Option, serialization_mode: &str) { - if serialization_mode == "advanced" { - return; - } let Some(fd) = fd_var .and_then(|s| s.parse::().ok()) .filter(|fd| *fd >= 0) @@ -153,31 +154,39 @@ fn initialize_unix_ipc(fd_var: Option, serialization_mode: &str) { Ok(send) => send, Err(_) => return, }; - spawn_ipc_reader(stream); + let advanced = serialization_mode == "advanced"; + spawn_ipc_reader(stream, advanced); let mut state = ipc_lock(); state.available = true; state.connected = true; state.refed = false; state.disconnect_emitted = false; + state.advanced = advanced; state.send = Some(send); } #[cfg(unix)] -fn spawn_ipc_reader(sock: UnixStream) { +fn spawn_ipc_reader(sock: UnixStream, advanced: bool) { // Cluster workers (#4962) need a recvmsg-based reader to receive SCM_RIGHTS - // connection fds for SCHED_RR; it still surfaces ordinary JSON frames - // identically. Plain forks keep the lighter BufReader path. + // connection fds for SCHED_RR in both serialization modes. Plain forks + // keep the lighter framing-specific readers. if crate::cluster::is_cluster_worker() { std::thread::spawn(move || { crate::cluster_sched::worker_recv_loop( sock, + advanced, |line| push_ipc_event(IpcEvent::Message(line)), + |bytes| push_ipc_event(IpcEvent::MessageAdvanced(bytes)), || push_ipc_event(IpcEvent::Closed), ); }); return; } + if advanced { + spawn_ipc_reader_advanced(sock); + return; + } std::thread::spawn(move || { let reader = std::io::BufReader::new(sock); for line in reader.lines() { @@ -191,6 +200,42 @@ fn spawn_ipc_reader(sock: UnixStream) { }); } +#[cfg(unix)] +fn spawn_ipc_reader_advanced(mut sock: UnixStream) { + std::thread::spawn(move || { + let mut acc = Vec::with_capacity(8192); + let mut chunk = [0u8; 8192]; + loop { + let n = match sock.read(&mut chunk) { + Ok(0) => break, + Ok(n) => n, + Err(_) => break, + }; + acc.extend_from_slice(&chunk[..n]); + + let mut consumed = 0; + while acc.len() - consumed >= 4 { + let len = u32::from_be_bytes([ + acc[consumed], + acc[consumed + 1], + acc[consumed + 2], + acc[consumed + 3], + ]) as usize; + if acc.len() - consumed - 4 < len { + break; + } + let start = consumed + 4; + push_ipc_event(IpcEvent::MessageAdvanced(acc[start..start + len].to_vec())); + consumed = start + len; + } + if consumed != 0 { + acc.drain(..consumed); + } + } + push_ipc_event(IpcEvent::Closed); + }); +} + fn push_ipc_event(event: IpcEvent) { ipc_lock().queue.push_back(event); } @@ -302,9 +347,10 @@ extern "C" fn process_ipc_send_fn( } pub(crate) fn process_ipc_send_call(message: f64, a2: f64, a3: f64, a4: f64) -> f64 { - let callback = [a4, a3, a2] - .into_iter() - .find(|v| !crate::fs::extract_closure_ptr(*v).is_null()); + let callback = [a4, a3, a2].into_iter().find(|v| { + JSValue::from_bits(v.to_bits()).is_pointer() + && !crate::fs::extract_closure_ptr(*v).is_null() + }); let ok = process_ipc_send_message(message); if let Some(cb) = callback { defer_send_callback(cb, ok); @@ -331,8 +377,14 @@ fn process_ipc_send_message(message: f64) -> bool { #[cfg(unix)] { - let Some(frame) = json_frame(message) else { - return false; + let advanced = ipc_lock().advanced; + let frame = if advanced { + advanced_frame(message) + } else { + let Some(frame) = json_frame(message) else { + return false; + }; + frame }; let mut emit_disconnect = false; let ok = { @@ -378,9 +430,17 @@ pub(crate) fn process_ipc_send_raw_json(json: &str) -> bool { } #[cfg(unix)] { - let mut frame = Vec::with_capacity(json.len() + 1); - frame.extend_from_slice(json.as_bytes()); - frame.push(b'\n'); + let advanced = ipc_lock().advanced; + let frame = if advanced { + let sh = js_string_from_bytes(json.as_ptr(), json.len() as u32); + let message = f64::from_bits(unsafe { crate::json::js_json_parse(sh) }.bits()); + advanced_frame(message) + } else { + let mut frame = Vec::with_capacity(json.len() + 1); + frame.extend_from_slice(json.as_bytes()); + frame.push(b'\n'); + frame + }; let mut state = ipc_lock(); if !state.available || !state.connected { return false; @@ -392,6 +452,15 @@ pub(crate) fn process_ipc_send_raw_json(json: &str) -> bool { } } +#[cfg(unix)] +fn advanced_frame(message: f64) -> Vec { + let payload = crate::child_process::v8_serialize(message); + let mut frame = Vec::with_capacity(payload.len() + 4); + frame.extend_from_slice(&(payload.len() as u32).to_be_bytes()); + frame.extend_from_slice(&payload); + frame +} + #[cfg(unix)] fn json_frame(message: f64) -> Option> { let sh = unsafe { crate::json::js_json_stringify(message, 0) }; @@ -518,10 +587,17 @@ pub extern "C" fn js_process_ipc_drain() -> i32 { crate::os::emit_process_event("message", &[msg]); count = count.saturating_add(1); } + IpcEvent::MessageAdvanced(bytes) => { + let msg = crate::child_process::v8_deserialize(&bytes); + crate::os::emit_process_event("message", &[msg]); + count = count.saturating_add(1); + } IpcEvent::Closed => { - if mark_closed_from_event() { + let unexpected = mark_closed_from_event(); + if unexpected { crate::os::emit_process_event("disconnect", &[]); } + crate::cluster::exit_worker_on_unexpected_primary_disconnect(unexpected); count = count.saturating_add(1); } } diff --git a/test-parity/node-suite/cluster/README.md b/test-parity/node-suite/cluster/README.md index 1fd806048a..3cae430179 100644 --- a/test-parity/node-suite/cluster/README.md +++ b/test-parity/node-suite/cluster/README.md @@ -10,13 +10,13 @@ terminal path. It does not depend on the historical broad parity files. ## Coverage - module/default export identity, role flags, constants, descriptors, and - EventEmitter/Worker prototype shape; + EventEmitter/Worker prototype shape, listener validation, and bookkeeping; - setup defaults, cumulative snapshots, aliases, setup events, fork option forwarding, and synchronous validation; - fork registry/id/respawn behavior and primary/worker role state; - worker online/message/disconnect/exit ordering, payloads, methods, channel, send callbacks/errors, kill metadata, and cluster-wide disconnect; -- JSON and advanced IPC serialization; +- JSON and advanced IPC serialization, including advanced-mode TCP handoff; - single-worker TCP listening descriptors and request/response behavior on port 0. diff --git a/test-parity/node-suite/cluster/STATUS.md b/test-parity/node-suite/cluster/STATUS.md index e6a4fc0d67..996421d5d9 100644 --- a/test-parity/node-suite/cluster/STATUS.md +++ b/test-parity/node-suite/cluster/STATUS.md @@ -25,31 +25,32 @@ substitute expected-output sources. ## Measured coverage -- 41 granular TypeScript fixtures (40 added over the original one). -- Node v26.5.0: 41/41 complete successfully in three consecutive direct rounds. -- Deno 2.9.2 local comparison: 41/41 complete successfully. -- Bun 1.2.18 local comparison: 40/41 complete successfully; its older local - release fails the `ChildProcess.channel` ref/unref probe. Current Bun source, - rather than this older binary, was used for selection evidence. -- Perry differential result: stable 15/41 (36.6%), with 25 behavioral - differences and one runtime timeout. +- 43 granular TypeScript fixtures (42 added over the original one). +- Node v26.5.x: all 43 complete successfully. +- Deno 2.9.2 local comparison: 41/41 complete successfully before the two + maintainer-audit additions. +- Bun 1.2.18 local comparison: 40/41 complete successfully before the two + maintainer-audit additions; its older local release fails the + `ChildProcess.channel` ref/unref probe. Current Bun source, rather than this + older binary, was used for selection evidence. +- Perry differential result: 43/43 (100%) on the final maintainer-audit run, + with no output differences, compile failures, crashes, or skips. -## Stable diagnostic boundaries +## Repaired diagnostic boundaries -The granular cases intentionally report, rather than repair, mismatches in areas -such as Worker construction/prototypes, setup aliases/events, empty disconnect -timing, fork/event ordering, worker/cluster message forwarding, worker state and -exit payloads, option validation, and TCP listening. +The formerly measured Worker construction/prototype, setup alias/event, empty +disconnect timing, fork/event ordering, worker/cluster message forwarding, +worker state/exit payload, option validation, TCP listening, and advanced IPC +differences now match Node across this suite. ## Stopping exclusions The remaining upstream cases were reviewed and stopped in these categories: -- **Blocked by the single-worker TCP result:** round-robin versus shared - scheduling, multi-worker connection distribution, server restart, backlog, - pipe handles, socket transfer, and shared-handle races. Adding scheduler - assertions before basic listening/request-response parity would obscure the - primary gap. +- **Broader scheduler stress:** multi-worker connection distribution, server + restart, backlog, pipe handles, socket transfer, and shared-handle races. The + deterministic suite now covers single-worker SCHED_RR request/response in both + JSON and advanced serialization modes. - **UDP foundation / duplicate semantics:** dgram sharing, reuse, fd binding, IPv6-only and unshared-UDP disconnect cases. These belong after TCP cluster lifecycle is reliable and otherwise repeat the granular `node:dgram` suite. @@ -71,14 +72,12 @@ access, arbitrary readiness sleep, or scheduler-dependent worker ordering. ## Final verification -- `deno fmt --check test-parity/node-suite/cluster` — 43 files checked. -- Three direct Node v26.5.0 rounds — 41/41 exited successfully each round. -- Two consecutive warmed focused differential rounds — 15/41 each, 25 behavioral - differences and one Perry runtime timeout; no Node failures or compile/link - failures. -- A per-fixture diagnostic run reproduced the same 15 pass, 25 difference, one - runtime-timeout partition and identified the timeout as invalid - serialization/inspect-port validation. +- Three direct Node v26.5.0 rounds — 41/41 exited successfully each round; Node + v26.5.1 also completes both maintainer-audit additions. +- Focused maintainer checks: listener validation/bookkeeping, child-side Worker + shape/send delegation, and advanced-mode TCP handoff all pass. +- Final normal-mode differential run: 43/43 (100%), with no Node failures, + output differences, compile/link failures, crashes, or skips. Only the measured cluster module floor was changed in `node_suite_baseline.json`; its historical full-suite aggregate remains the last diff --git a/test-parity/node-suite/cluster/default-import-shape.ts b/test-parity/node-suite/cluster/default-import-shape.ts index 11eb00fff6..da948f3b1e 100644 --- a/test-parity/node-suite/cluster/default-import-shape.ts +++ b/test-parity/node-suite/cluster/default-import-shape.ts @@ -43,6 +43,10 @@ for ( "emit", "eventNames", "listenerCount", + "listeners", + "rawListeners", + "getMaxListeners", + "setMaxListeners", ] ) { console.log( diff --git a/test-parity/node-suite/cluster/events/listener-validation.ts b/test-parity/node-suite/cluster/events/listener-validation.ts new file mode 100644 index 0000000000..2951231b33 --- /dev/null +++ b/test-parity/node-suite/cluster/events/listener-validation.ts @@ -0,0 +1,59 @@ +// EventEmitter listener arguments are validated synchronously on the cluster +// singleton for both registration and removal methods. +import cluster from "node:cluster"; + +for ( + const method of [ + "on", + "once", + "prependListener", + "prependOnceListener", + "removeListener", + "off", + ] as const +) { + try { + cluster[method]("probe", 42 as never); + } catch (error: any) { + console.log(method, error.name, error.code); + } +} + +const duplicate = () => {}; +cluster.on("count", duplicate); +cluster.on("count", duplicate); +cluster.on("count", () => {}); +console.log( + "filtered count", + cluster.listenerCount("count"), + cluster.listenerCount("count", duplicate), + cluster.listenerCount("count", 42 as never), +); + +cluster.on(null as never, () => {}); +cluster.on(undefined as never, () => {}); +cluster.removeAllListeners(null as never); +console.log( + "explicit event", + cluster.listenerCount("null"), + cluster.listenerCount("undefined"), + cluster.listenerCount("count"), +); +cluster.removeAllListeners(); +console.log("all events", cluster.eventNames().length); + +const regular = () => {}; +cluster.on("surface", regular); +cluster.once("surface", () => {}); +console.log( + "listener arrays", + cluster.listeners("surface").length, + cluster.rawListeners("surface").length, + cluster.listeners("surface")[0] === regular, +); +console.log("max default", cluster.getMaxListeners()); +console.log( + "max set", + cluster.setMaxListeners(17) === cluster, + cluster.getMaxListeners(), +); diff --git a/test-parity/node-suite/cluster/fork/roles-and-env.ts b/test-parity/node-suite/cluster/fork/roles-and-env.ts index f54fca6d01..b02c96ff2b 100644 --- a/test-parity/node-suite/cluster/fork/roles-and-env.ts +++ b/test-parity/node-suite/cluster/fork/roles-and-env.ts @@ -1,5 +1,6 @@ // Worker-side role flags and NODE_UNIQUE_ID consumption from Node's basic test. import cluster from "node:cluster"; +import { EventEmitter } from "node:events"; if (cluster.isPrimary) { const worker = cluster.fork({ CLUSTER_PROBE: "present" }); @@ -11,11 +12,16 @@ if (cluster.isPrimary) { }); worker.once("exit", () => clearTimeout(watchdog)); } else { - process.send!({ + cluster.worker!.send({ isPrimary: cluster.isPrimary, isMaster: cluster.isMaster, isWorker: cluster.isWorker, id: cluster.worker?.id, + state: cluster.worker?.state, + processSame: cluster.worker?.process === process, + connected: cluster.worker?.isConnected(), + workerInstance: cluster.worker instanceof cluster.Worker, + emitterInstance: cluster.worker instanceof EventEmitter, env: process.env.CLUSTER_PROBE, uniqueIdRemoved: process.env.NODE_UNIQUE_ID === undefined, }); diff --git a/test-parity/node-suite/cluster/serialization/advanced-tcp.ts b/test-parity/node-suite/cluster/serialization/advanced-tcp.ts new file mode 100644 index 0000000000..6a7244688c --- /dev/null +++ b/test-parity/node-suite/cluster/serialization/advanced-tcp.ts @@ -0,0 +1,27 @@ +// Advanced IPC framing must coexist with cluster's SCHED_RR socket handoff. +import cluster from "node:cluster"; +import { connect, createServer } from "node:net"; + +if (cluster.isPrimary) cluster.setupPrimary({ serialization: "advanced" }); + +if (cluster.isWorker) { + createServer((socket) => socket.end("advanced-ok")).listen(0, "127.0.0.1"); +} else { + const worker = cluster.fork(); + const watchdog = setTimeout(() => worker.kill(), 5_000); + watchdog.unref(); + worker.once("listening", (address) => { + const socket = connect(address.port, "127.0.0.1"); + let data = ""; + socket.setEncoding("utf8"); + socket.on("data", (chunk) => data += chunk); + socket.once("end", () => { + console.log("response:", data); + worker.disconnect(); + }); + }); + worker.once("exit", (code, signal) => { + clearTimeout(watchdog); + console.log("exit:", code, signal); + }); +} diff --git a/test-parity/node-suite/cluster/serialization/advanced.ts b/test-parity/node-suite/cluster/serialization/advanced.ts index 825ed4a7dd..f5bdd445b3 100644 --- a/test-parity/node-suite/cluster/serialization/advanced.ts +++ b/test-parity/node-suite/cluster/serialization/advanced.ts @@ -2,7 +2,7 @@ // clone values that JSON cannot preserve. import cluster from "node:cluster"; -cluster.setupPrimary({ serialization: "advanced" }); +if (cluster.isPrimary) cluster.setupPrimary({ serialization: "advanced" }); if (cluster.isWorker) { process.once("message", (message: any) => { process.send!({ diff --git a/test-parity/node-suite/cluster/setup/validation-exec-args.ts b/test-parity/node-suite/cluster/setup/validation-exec-args.ts index 0676018f2e..f02fd5a9e7 100644 --- a/test-parity/node-suite/cluster/setup/validation-exec-args.ts +++ b/test-parity/node-suite/cluster/setup/validation-exec-args.ts @@ -6,7 +6,9 @@ for (const [name, value] of [["exec", 1], ["args", 1]] as const) { // setupPrimary merges into cluster.settings, so restore valid defaults for both // fields each iteration and override only the current one with the invalid value; // otherwise the previous iteration's bad field leaks in and masks this case. - cluster.setupPrimary({ exec: process.argv[1], args: [], [name]: value } as any); + cluster.setupPrimary( + { exec: process.argv[1], args: [], [name]: value } as any, + ); try { cluster.fork(); console.log(name, "accepted"); diff --git a/test-parity/node-suite/cluster/setup/validation-serialization-inspect.ts b/test-parity/node-suite/cluster/setup/validation-serialization-inspect.ts index 015d18b589..f0c75cc4a9 100644 --- a/test-parity/node-suite/cluster/setup/validation-serialization-inspect.ts +++ b/test-parity/node-suite/cluster/setup/validation-serialization-inspect.ts @@ -10,7 +10,9 @@ for ( // setupPrimary merges into cluster.settings, so restore valid defaults each // iteration and spread the invalid field last; otherwise the previous case's bad // value (e.g. serialization: "bad") leaks in and masks this case's validation. - cluster.setupPrimary({ serialization: "json", inspectPort: undefined, ...options } as any); + cluster.setupPrimary( + { serialization: "json", inspectPort: undefined, ...options } as any, + ); try { cluster.fork(); console.log(name, "accepted"); diff --git a/test-parity/node_suite_baseline.json b/test-parity/node_suite_baseline.json index 833aae253c..de4f2450aa 100644 --- a/test-parity/node_suite_baseline.json +++ b/test-parity/node_suite_baseline.json @@ -26,8 +26,8 @@ "total": 53 }, "cluster": { - "pass": 15, - "total": 41 + "pass": 43, + "total": 43 }, "console": { "pass": 119,