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,