From 48723266d954d40c752422e6358e451886adf76f Mon Sep 17 00:00:00 2001 From: Jarred Sumner Date: Fri, 21 Aug 2026 03:06:34 +0000 Subject: [PATCH 1/4] node:net: keep the process alive while a socket that was unref'd or paused early is still connecting `socket.unref()` (or `pause()`) called before `connect()` or while the connection was still pending released the handle's hold on the event loop immediately, so a process with nothing else pending exited before 'connect' ever fired (breaks testcontainers' Ryuk client). Node keeps the loop alive for a pending connect regardless of the handle's ref state and only lets go once it is established. The native socket now always holds the loop for the connect attempt and records ref()/unref() issued before it is established, applying that preference in on_open. An autoSelectFamily retry handle also inherits a pause() issued mid-connect, not just unref(). Fixes #37086 --- src/js/node/net.ts | 4 +- src/runtime/node/node_net_binding.rs | 2 +- src/runtime/socket/Listener.rs | 24 +++++------ src/runtime/socket/socket_body.rs | 41 +++++++----------- test/js/node/net/node-net.test.ts | 63 ++++++++++++++++++++++++---- 5 files changed, 84 insertions(+), 50 deletions(-) diff --git a/src/js/node/net.ts b/src/js/node/net.ts index 1e9d83b4dda9..11122852ca96 100644 --- a/src/js/node/net.ts +++ b/src/js/node/net.ts @@ -4272,8 +4272,8 @@ function initSocketHandle(self) { const handle = self._handle; if (handle) { handle[owner_symbol] = self; - // A fresh handle (e.g. an autoSelectFamily retry) inherits a prior unref(). - if (self[kUserUnrefed]) handle.unref?.(); + // A fresh handle (e.g. an autoSelectFamily retry) inherits a prior unref()/pause(). + if (self[kUserUnrefed] || self[kPausedUnref]) handle.unref?.(); } } diff --git a/src/runtime/node/node_net_binding.rs b/src/runtime/node/node_net_binding.rs index 23c77f09f614..d7e6f32468ce 100644 --- a/src/runtime/node/node_net_binding.rs +++ b/src/runtime/node/node_net_binding.rs @@ -147,7 +147,7 @@ pub(crate) fn new_detached_socket(global: &JSGlobalObject, frame: &CallFrame) -> flags: Cell::new(SocketFlags::default() | SocketFlags::DEFERS_SERVER_IDENTITY), this_value: JsCell::new(jsc::JsRef::empty()), poll_ref: JsCell::new(KeepAlive::init()), - ref_pollref_on_connect: Cell::new(true), + user_wants_ref: Cell::new(true), connection: JsCell::new(None), server_name: JsCell::new(None), buffered_data_for_node_net: Default::default(), diff --git a/src/runtime/socket/Listener.rs b/src/runtime/socket/Listener.rs index d36d771db941..d39d1dea807c 100644 --- a/src/runtime/socket/Listener.rs +++ b/src/runtime/socket/Listener.rs @@ -626,7 +626,7 @@ impl Listener { owned_ssl_ctx: Cell::new(None), this_value: JsCell::new(jsc::JsRef::empty()), poll_ref: JsCell::new(KeepAlive::init()), - ref_pollref_on_connect: Cell::new(true), + user_wants_ref: Cell::new(true), connection: JsCell::new(None), local_binding: JsCell::new(None), server_name: JsCell::new(None), @@ -672,7 +672,7 @@ impl Listener { owned_ssl_ctx: Cell::new(None), this_value: JsCell::new(jsc::JsRef::empty()), poll_ref: JsCell::new(KeepAlive::init()), - ref_pollref_on_connect: Cell::new(true), + user_wants_ref: Cell::new(true), connection: JsCell::new(None), local_binding: JsCell::new(None), server_name: JsCell::new(None), @@ -1256,7 +1256,7 @@ impl Listener { flags: Cell::new(SocketFlags::default()), this_value: JsCell::new(jsc::JsRef::empty()), poll_ref: JsCell::new(KeepAlive::init()), - ref_pollref_on_connect: Cell::new(true), + user_wants_ref: Cell::new(true), buffered_data_for_node_net: Default::default(), bytes_written: Cell::new(0), native_callback: JsCell::new(crate::socket::NativeCallbacks::None), @@ -1342,7 +1342,7 @@ impl Listener { flags: Cell::new(SocketFlags::default()), this_value: JsCell::new(jsc::JsRef::empty()), poll_ref: JsCell::new(KeepAlive::init()), - ref_pollref_on_connect: Cell::new(true), + user_wants_ref: Cell::new(true), buffered_data_for_node_net: Default::default(), bytes_written: Cell::new(0), native_callback: JsCell::new(crate::socket::NativeCallbacks::None), @@ -1584,7 +1584,7 @@ fn connect_finish( flags: Cell::new(SocketFlags::default()), this_value: JsCell::new(jsc::JsRef::empty()), poll_ref: JsCell::new(KeepAlive::init()), - ref_pollref_on_connect: Cell::new(true), + user_wants_ref: Cell::new(true), buffered_data_for_node_net: Default::default(), bytes_written: Cell::new(0), native_callback: JsCell::new(crate::socket::NativeCallbacks::None), @@ -1619,6 +1619,12 @@ fn connect_finish( f.set(SocketFlags::ALLOW_HALF_OPEN, allow_half_open); socket_ref.flags.set(f); } + // The connect attempt holds the loop whatever `.ref()`/`.unref()` said so + // far (node:net hands out the handle before connecting); `on_open` applies + // `user_wants_ref`, and every failure path releases `poll_ref` itself. + socket_ref + .poll_ref + .with_mut(|p| p.ref_(bun_io::js_vm_ctx())); // Note: `do_connect` reads `self.connection` directly so no second // borrow is needed here. // An already-open fd socket runs `on_open` synchronously; what settling @@ -1687,14 +1693,6 @@ fn connect_finish( } }; - // if this is from node:net there's surface where the user can .ref() and .deref() - // before the connection starts. make sure we honor that here. - if socket_ref.ref_pollref_on_connect.get() && !socket_ref.socket.get().is_closed() { - socket_ref - .poll_ref - .with_mut(|p| p.ref_(bun_io::js_vm_ctx())); - } - // What settling the connect promise in `on_open` left pending (allocation // failure, a terminating VM). if let Some(err) = opened_err { diff --git a/src/runtime/socket/socket_body.rs b/src/runtime/socket/socket_body.rs index 40d5befa04f1..da7b4200eed9 100644 --- a/src/runtime/socket/socket_body.rs +++ b/src/runtime/socket/socket_body.rs @@ -308,7 +308,7 @@ pub struct NewSocket { /// downgraded to weak once the socket is closed/inactive so GC can reclaim it. pub this_value: JsCell, pub poll_ref: JsCell, - pub(crate) ref_pollref_on_connect: Cell, + pub(crate) user_wants_ref: Cell, pub(crate) connection: JsCell>, /// `localAddress`/`localPort` from the connect options: the socket is /// bound to this address before connecting. Always a literal IP. @@ -1423,6 +1423,9 @@ impl NewSocket { // update the internal socket instance to the one that was just connected // This socket must be replaced because the previous one is a connecting socket not a uSockets socket this.socket.set(socket); + if !this.user_wants_ref.get() { + this.poll_ref.with_mut(|p| p.unref(js_loop_ctx())); + } jsc::mark_binding!(); // Add SNI support for TLS (mongodb and others requires this) @@ -3220,41 +3223,29 @@ impl NewSocket { #[bun_jsc::host_fn(method)] pub(crate) fn js_ref( this: &Self, - global: &JSGlobalObject, + _global: &JSGlobalObject, _frame: &CallFrame, ) -> JsResult { jsc::mark_binding!(); - if this.socket.get().is_detached() { - this.ref_pollref_on_connect.set(true); - } - if this.socket.get().is_detached() { - return Ok(JSValue::UNDEFINED); + this.user_wants_ref.set(true); + // Not yet established: `connect_finish` holds the loop and `on_open` applies this. + if this.socket.get().is_established() { + this.poll_ref.with_mut(|p| p.ref_(js_loop_ctx())); } - let _ = global; - this.poll_ref.with_mut(|p| { - p.ref_(bun_io::posix_event_loop::get_vm_ctx( - bun_io::AllocatorType::Js, - )) - }); Ok(JSValue::UNDEFINED) } #[bun_jsc::host_fn(method)] pub(crate) fn js_unref( this: &Self, - global: &JSGlobalObject, + _global: &JSGlobalObject, _frame: &CallFrame, ) -> JsResult { jsc::mark_binding!(); - if this.socket.get().is_detached() { - this.ref_pollref_on_connect.set(false); + this.user_wants_ref.set(false); + if this.socket.get().is_established() { + this.poll_ref.with_mut(|p| p.unref(js_loop_ctx())); } - let _ = global; - this.poll_ref.with_mut(|p| { - p.unref(bun_io::posix_event_loop::get_vm_ctx( - bun_io::AllocatorType::Js, - )) - }); Ok(JSValue::UNDEFINED) } @@ -3595,7 +3586,7 @@ impl NewSocket { flags: Cell::new(initial_flags), this_value: JsCell::new(JsRef::empty()), poll_ref: JsCell::new(KeepAlive::init()), - ref_pollref_on_connect: Cell::new(true), + user_wants_ref: Cell::new(true), buffered_data_for_node_net: JsCell::new(Vec::new()), bytes_written: Cell::new(0), native_callback: JsCell::new(NativeCallbacks::None), @@ -3702,7 +3693,7 @@ impl NewSocket { flags: Cell::new(Flags::BYPASS_TLS | Flags::IS_ACTIVE | Flags::OWNED_PROTOS), this_value: JsCell::new(JsRef::empty()), poll_ref: JsCell::new(KeepAlive::init()), - ref_pollref_on_connect: Cell::new(true), + user_wants_ref: Cell::new(true), buffered_data_for_node_net: JsCell::new(Vec::new()), bytes_written: Cell::new(0), native_callback: JsCell::new(NativeCallbacks::None), @@ -4822,7 +4813,7 @@ pub fn js_upgrade_duplex_to_tls( flags: Cell::new(initial_flags), this_value: JsCell::new(JsRef::empty()), poll_ref: JsCell::new(KeepAlive::init()), - ref_pollref_on_connect: Cell::new(true), + user_wants_ref: Cell::new(true), buffered_data_for_node_net: JsCell::new(Vec::new()), bytes_written: Cell::new(0), native_callback: JsCell::new(NativeCallbacks::None), diff --git a/test/js/node/net/node-net.test.ts b/test/js/node/net/node-net.test.ts index eadca761d75b..9d2973e142db 100644 --- a/test/js/node/net/node-net.test.ts +++ b/test/js/node/net/node-net.test.ts @@ -699,9 +699,9 @@ it("unref should exit when no more work pending", async () => { expect(await process.exited).toBe(0); }); -// An unref() applied while lookup is pending must survive the autoSelectFamily handle reinit. -it("unref survives an autoSelectFamily retry", async () => { - // IPv4-only server + injected lookup listing ::1 first forces a refused attempt then a retry; unref() runs mid-lookup. +// IPv4-only server + injected lookup listing ::1 first forces a refused attempt then a retry; the call runs mid-lookup +// and must carry over to the retry handle. The pending connect holds the loop by itself; once connected it lets go. +it.concurrent.each(["s.unref()", "s.pause()"])("%s survives an autoSelectFamily retry", async call => { const server = createServer(() => {}); await new Promise(resolve => server.listen(0, "127.0.0.1", resolve)); try { @@ -714,27 +714,72 @@ it("unref survives an autoSelectFamily retry", async () => { const lookup = (host, opts, cb) => setTimeout(() => cb(null, [{ address: "::1", family: 6 }, { address: "127.0.0.1", family: 4 }]), 10); const s = net.connect({ host: "localhost", port: ${server.address().port}, autoSelectFamily: true, lookup }); - s.on("data", () => {}); s.on("error", e => process.stdout.write("error " + e.code + "\\n")); s.on("connect", () => process.stdout.write("connected " + s.remoteAddress + "\\n")); - s.unref(); - // Sentinel keeping the loop alive across the refuse + retry. - setTimeout(() => process.stdout.write("timer\\n"), 500); + ${call}; `, ], env: bunEnv, stdout: "pipe", stderr: "inherit", }); - // After the sentinel timer only the unref'd socket remains, so the process must exit. const [stdout, exitCode] = await Promise.all([proc.stdout.text(), proc.exited]); - expect(stdout.trim().split("\n").sort()).toEqual(["connected 127.0.0.1", "timer"]); + expect(stdout).toBe("connected 127.0.0.1\n"); expect(exitCode).toBe(0); } finally { server.close(); } }); +// https://github.com/oven-sh/bun/issues/37086 — node's pending uv_connect_t keeps the loop alive even on an +// unref'd/non-reading handle, so unref()/pause() issued before or while connecting only take effect once connected. +describe.concurrent("unref()/pause() around connect()", () => { + async function run(client: string, onConnection = "c => c.unref()") { + await using proc = Bun.spawn({ + cmd: [ + bunExe(), + "-e", + ` + const net = require("net"); + const server = net.createServer(${onConnection}); + server.listen(0, "127.0.0.1", () => { + const port = server.address().port; + const s = new net.Socket(); + s.on("connect", () => process.stdout.write("connected\\n")); + s.on("close", () => { process.stdout.write("closed\\n"); server.close(); }); + ${client} + }); + server.unref(); + `, + ], + env: bunEnv, + stdout: "pipe", + stderr: "inherit", + }); + const [stdout, exitCode] = await Promise.all([proc.stdout.text(), proc.exited]); + return { stdout, exitCode }; + } + + it.each([ + ["unref() before connect()", `s.unref(); s.connect(port, "127.0.0.1");`], + ["unref() while connecting", `s.connect(port, "127.0.0.1"); s.unref();`], + ["pause() while connecting", `s.connect(port, "127.0.0.1"); s.pause();`], + ])("%s waits for the connection, then lets the process exit", async (_, client) => { + const { stdout, exitCode } = await run(client); + expect(stdout).toBe("connected\n"); + expect(exitCode).toBe(0); + }); + + it("ref() after unref() while connecting keeps holding the loop", async () => { + const { stdout, exitCode } = await run( + `s.connect(port, "127.0.0.1"); s.unref(); s.ref(); s.resume();`, + "c => { c.unref(); c.end(); }", + ); + expect(stdout).toBe("connected\nclosed\n"); + expect(exitCode).toBe(0); + }); +}); + it("socket should keep process alive if unref is not called", async () => { const process = Bun.spawn({ cmd: [bunExe(), join(import.meta.dir, "node-ref-default-fixture.js")], From 2bcc4412a15cdd25a20e72756dca528fb72daeee Mon Sep 17 00:00:00 2001 From: Jarred Sumner Date: Fri, 21 Aug 2026 05:29:00 +0000 Subject: [PATCH 2/4] Address review: hold the loop from connect start, TLS upgrade inherits the ref preference, cover pause()/in-flight cases - connect_finish takes the loop hold before do_connect so a synchronously opened fd socket has on_open apply the preference after it, instead of reconstructing what happened afterwards. - upgradeTLS carries user_wants_ref onto the TLS socket. - node:net: initSocketHandle now runs for every connect() (including a reconnect through a live handle) and re-applies unref()/pause() or ref() for the new connection; the pre-connect nextTick read(0) is left to afterConnect while still connecting so a pause() issued mid-connect is not undone by a read queued before it. - Tests: pipe and assert stderr, add pause()-before-connect and unref()/pause()/ref() while the native connect is actually in flight, and a pause() variant of the autoSelectFamily retry test. --- src/js/node/net.ts | 18 ++++++++-------- src/runtime/socket/Listener.rs | 4 +--- src/runtime/socket/socket_body.rs | 2 +- test/js/node/net/node-net.test.ts | 35 +++++++++++++++++++------------ 4 files changed, 33 insertions(+), 26 deletions(-) diff --git a/src/js/node/net.ts b/src/js/node/net.ts index 11122852ca96..dbb7cbc6f85b 100644 --- a/src/js/node/net.ts +++ b/src/js/node/net.ts @@ -1926,14 +1926,12 @@ Socket.prototype.connect = function connect(...args) { this.pause(); } else { process.nextTick(() => { - // Honor pause()/resume() calls made while connecting — only start - // reading if the user hasn't explicitly paused the stream. Matches - // Node's afterConnect, which calls socket.read(0) only when not paused: + // An already-open handle (fd, wrapped duplex) starts reading here unless + // the user paused; read(0) does that without switching to flowing mode. + // A pending connect gets this from afterConnect instead, so a pause() + // that lands before then is still honored: // https://github.com/nodejs/node/blob/843dc5f0d5ad/lib/net.js#L1649 - // read(0) starts the handle reading without switching the stream into - // flowing mode, so data that arrives before a 'data' listener is - // attached stays buffered instead of being emitted to nobody. - if (!this.isPaused()) this.read(0); + if (!this.connecting && !this.isPaused()) this.read(0); }); if (fd == null) this.connecting = true; } @@ -2134,8 +2132,8 @@ Socket.prototype.connect = function connect(...args) { if (!this._handle) { this._handle = newDetachedSocket(typeof this[bunTlsSymbol] === "function"); - initSocketHandle(this); } + initSocketHandle(this); if (!pipe) { lookupAndConnect(this, options); @@ -4272,8 +4270,10 @@ function initSocketHandle(self) { const handle = self._handle; if (handle) { handle[owner_symbol] = self; - // A fresh handle (e.g. an autoSelectFamily retry) inherits a prior unref()/pause(). + // The new connection (fresh handle, autoSelectFamily retry, or reconnect through a + // live handle) inherits a prior unref()/pause(), not the previous connection's hold. if (self[kUserUnrefed] || self[kPausedUnref]) handle.unref?.(); + else handle.ref?.(); } } diff --git a/src/runtime/socket/Listener.rs b/src/runtime/socket/Listener.rs index d39d1dea807c..3f2d4f466e46 100644 --- a/src/runtime/socket/Listener.rs +++ b/src/runtime/socket/Listener.rs @@ -1619,9 +1619,7 @@ fn connect_finish( f.set(SocketFlags::ALLOW_HALF_OPEN, allow_half_open); socket_ref.flags.set(f); } - // The connect attempt holds the loop whatever `.ref()`/`.unref()` said so - // far (node:net hands out the handle before connecting); `on_open` applies - // `user_wants_ref`, and every failure path releases `poll_ref` itself. + // Held for the connect attempt regardless of `user_wants_ref`; `on_open` applies that. socket_ref .poll_ref .with_mut(|p| p.ref_(bun_io::js_vm_ctx())); diff --git a/src/runtime/socket/socket_body.rs b/src/runtime/socket/socket_body.rs index da7b4200eed9..c9bcea862fd8 100644 --- a/src/runtime/socket/socket_body.rs +++ b/src/runtime/socket/socket_body.rs @@ -3586,7 +3586,7 @@ impl NewSocket { flags: Cell::new(initial_flags), this_value: JsCell::new(JsRef::empty()), poll_ref: JsCell::new(KeepAlive::init()), - user_wants_ref: Cell::new(true), + user_wants_ref: Cell::new(this.user_wants_ref.get()), buffered_data_for_node_net: JsCell::new(Vec::new()), bytes_written: Cell::new(0), native_callback: JsCell::new(NativeCallbacks::None), diff --git a/test/js/node/net/node-net.test.ts b/test/js/node/net/node-net.test.ts index 9d2973e142db..6089a25e9a12 100644 --- a/test/js/node/net/node-net.test.ts +++ b/test/js/node/net/node-net.test.ts @@ -721,10 +721,11 @@ it.concurrent.each(["s.unref()", "s.pause()"])("%s survives an autoSelectFamily ], env: bunEnv, stdout: "pipe", - stderr: "inherit", + stderr: "pipe", }); - const [stdout, exitCode] = await Promise.all([proc.stdout.text(), proc.exited]); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); expect(stdout).toBe("connected 127.0.0.1\n"); + expect(stderr).toBe(""); expect(exitCode).toBe(0); } finally { server.close(); @@ -754,28 +755,36 @@ describe.concurrent("unref()/pause() around connect()", () => { ], env: bunEnv, stdout: "pipe", - stderr: "inherit", + stderr: "pipe", }); - const [stdout, exitCode] = await Promise.all([proc.stdout.text(), proc.exited]); - return { stdout, exitCode }; + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + return { stdout, stderr, exitCode }; } + // net.ts hands the connect to the native socket on the next tick, so two ticks in the + // attempt is in flight (the handle exists but is not yet established). + const inFlight = (code: string) => `process.nextTick(() => process.nextTick(() => { ${code} }));`; it.each([ ["unref() before connect()", `s.unref(); s.connect(port, "127.0.0.1");`], - ["unref() while connecting", `s.connect(port, "127.0.0.1"); s.unref();`], - ["pause() while connecting", `s.connect(port, "127.0.0.1"); s.pause();`], + ["unref() right after connect()", `s.connect(port, "127.0.0.1"); s.unref();`], + ["unref() while the connect is in flight", `s.connect(port, "127.0.0.1"); ${inFlight("s.unref();")}`], + ["pause() before connect()", `s.pause(); s.connect(port, "127.0.0.1");`], + ["pause() right after connect()", `s.connect(port, "127.0.0.1"); s.pause();`], + ["pause() while the connect is in flight", `s.connect(port, "127.0.0.1"); ${inFlight("s.pause();")}`], ])("%s waits for the connection, then lets the process exit", async (_, client) => { - const { stdout, exitCode } = await run(client); + const { stdout, stderr, exitCode } = await run(client); expect(stdout).toBe("connected\n"); + expect(stderr).toBe(""); expect(exitCode).toBe(0); }); - it("ref() after unref() while connecting keeps holding the loop", async () => { - const { stdout, exitCode } = await run( - `s.connect(port, "127.0.0.1"); s.unref(); s.ref(); s.resume();`, - "c => { c.unref(); c.end(); }", - ); + it.each([ + ["right after connect()", `s.connect(port, "127.0.0.1"); s.unref(); s.ref(); s.resume();`], + ["while the connect is in flight", `s.connect(port, "127.0.0.1"); ${inFlight("s.unref(); s.ref(); s.resume();")}`], + ])("ref() after unref() %s keeps holding the loop", async (_, client) => { + const { stdout, stderr, exitCode } = await run(client, "c => { c.unref(); c.end(); }"); expect(stdout).toBe("connected\nclosed\n"); + expect(stderr).toBe(""); expect(exitCode).toBe(0); }); }); From 2b7d543ddba12b3409ecb1461031667137f2a88f Mon Sep 17 00:00:00 2001 From: Jarred Sumner Date: Fri, 21 Aug 2026 06:20:36 +0000 Subject: [PATCH 3/4] Simplify: latch only pre-open ref()/unref(), don't pause a half-connected socket natively - js_ref/js_unref act on poll_ref directly once established and only record the preference before that; on_open applies and consumes it. An unref made by net.ts bookkeeping on an established socket (FIN, backpressure) no longer leaks into a later reconnect through the same wrapper, so the live-handle initSocketHandle call and the TLS-upgrade copy are dropped again. Field name reverted to ref_pollref_on_connect, which describes it. - pause()/resume() on a socket whose connect has not completed are no-ops natively (as they already were for the DNS `connecting` arm): usockets re-arms reads on open without clearing its paused bit, so a pause latched there made every later pause() a no-op and a paused stream buffered without bound once connected. on_open also clears a stale IS_PAUSED from a previous connection on the same wrapper. - Test for the above (bytesRead stops growing once past highWaterMark). --- src/js/node/net.ts | 6 ++--- src/runtime/node/node_net_binding.rs | 2 +- src/runtime/socket/Listener.rs | 12 ++++----- src/runtime/socket/socket_body.rs | 20 ++++++++------ src/uws_sys/socket.rs | 7 +++-- test/js/node/net/node-net.test.ts | 40 ++++++++++++++++++++++++++++ 6 files changed, 66 insertions(+), 21 deletions(-) diff --git a/src/js/node/net.ts b/src/js/node/net.ts index dbb7cbc6f85b..c891d8559ade 100644 --- a/src/js/node/net.ts +++ b/src/js/node/net.ts @@ -2132,8 +2132,8 @@ Socket.prototype.connect = function connect(...args) { if (!this._handle) { this._handle = newDetachedSocket(typeof this[bunTlsSymbol] === "function"); + initSocketHandle(this); } - initSocketHandle(this); if (!pipe) { lookupAndConnect(this, options); @@ -4270,10 +4270,8 @@ function initSocketHandle(self) { const handle = self._handle; if (handle) { handle[owner_symbol] = self; - // The new connection (fresh handle, autoSelectFamily retry, or reconnect through a - // live handle) inherits a prior unref()/pause(), not the previous connection's hold. + // A fresh handle (e.g. an autoSelectFamily retry) inherits a prior unref()/pause(). if (self[kUserUnrefed] || self[kPausedUnref]) handle.unref?.(); - else handle.ref?.(); } } diff --git a/src/runtime/node/node_net_binding.rs b/src/runtime/node/node_net_binding.rs index d7e6f32468ce..23c77f09f614 100644 --- a/src/runtime/node/node_net_binding.rs +++ b/src/runtime/node/node_net_binding.rs @@ -147,7 +147,7 @@ pub(crate) fn new_detached_socket(global: &JSGlobalObject, frame: &CallFrame) -> flags: Cell::new(SocketFlags::default() | SocketFlags::DEFERS_SERVER_IDENTITY), this_value: JsCell::new(jsc::JsRef::empty()), poll_ref: JsCell::new(KeepAlive::init()), - user_wants_ref: Cell::new(true), + ref_pollref_on_connect: Cell::new(true), connection: JsCell::new(None), server_name: JsCell::new(None), buffered_data_for_node_net: Default::default(), diff --git a/src/runtime/socket/Listener.rs b/src/runtime/socket/Listener.rs index 3f2d4f466e46..992aaf4a1a54 100644 --- a/src/runtime/socket/Listener.rs +++ b/src/runtime/socket/Listener.rs @@ -626,7 +626,7 @@ impl Listener { owned_ssl_ctx: Cell::new(None), this_value: JsCell::new(jsc::JsRef::empty()), poll_ref: JsCell::new(KeepAlive::init()), - user_wants_ref: Cell::new(true), + ref_pollref_on_connect: Cell::new(true), connection: JsCell::new(None), local_binding: JsCell::new(None), server_name: JsCell::new(None), @@ -672,7 +672,7 @@ impl Listener { owned_ssl_ctx: Cell::new(None), this_value: JsCell::new(jsc::JsRef::empty()), poll_ref: JsCell::new(KeepAlive::init()), - user_wants_ref: Cell::new(true), + ref_pollref_on_connect: Cell::new(true), connection: JsCell::new(None), local_binding: JsCell::new(None), server_name: JsCell::new(None), @@ -1256,7 +1256,7 @@ impl Listener { flags: Cell::new(SocketFlags::default()), this_value: JsCell::new(jsc::JsRef::empty()), poll_ref: JsCell::new(KeepAlive::init()), - user_wants_ref: Cell::new(true), + ref_pollref_on_connect: Cell::new(true), buffered_data_for_node_net: Default::default(), bytes_written: Cell::new(0), native_callback: JsCell::new(crate::socket::NativeCallbacks::None), @@ -1342,7 +1342,7 @@ impl Listener { flags: Cell::new(SocketFlags::default()), this_value: JsCell::new(jsc::JsRef::empty()), poll_ref: JsCell::new(KeepAlive::init()), - user_wants_ref: Cell::new(true), + ref_pollref_on_connect: Cell::new(true), buffered_data_for_node_net: Default::default(), bytes_written: Cell::new(0), native_callback: JsCell::new(crate::socket::NativeCallbacks::None), @@ -1584,7 +1584,7 @@ fn connect_finish( flags: Cell::new(SocketFlags::default()), this_value: JsCell::new(jsc::JsRef::empty()), poll_ref: JsCell::new(KeepAlive::init()), - user_wants_ref: Cell::new(true), + ref_pollref_on_connect: Cell::new(true), buffered_data_for_node_net: Default::default(), bytes_written: Cell::new(0), native_callback: JsCell::new(crate::socket::NativeCallbacks::None), @@ -1619,7 +1619,7 @@ fn connect_finish( f.set(SocketFlags::ALLOW_HALF_OPEN, allow_half_open); socket_ref.flags.set(f); } - // Held for the connect attempt regardless of `user_wants_ref`; `on_open` applies that. + // Held for the connect attempt regardless of `ref_pollref_on_connect`; `on_open` applies that. socket_ref .poll_ref .with_mut(|p| p.ref_(bun_io::js_vm_ctx())); diff --git a/src/runtime/socket/socket_body.rs b/src/runtime/socket/socket_body.rs index c9bcea862fd8..d9c14986f315 100644 --- a/src/runtime/socket/socket_body.rs +++ b/src/runtime/socket/socket_body.rs @@ -308,7 +308,7 @@ pub struct NewSocket { /// downgraded to weak once the socket is closed/inactive so GC can reclaim it. pub this_value: JsCell, pub poll_ref: JsCell, - pub(crate) user_wants_ref: Cell, + pub(crate) ref_pollref_on_connect: Cell, pub(crate) connection: JsCell>, /// `localAddress`/`localPort` from the connect options: the socket is /// bound to this address before connecting. Always a literal IP. @@ -1423,7 +1423,9 @@ impl NewSocket { // update the internal socket instance to the one that was just connected // This socket must be replaced because the previous one is a connecting socket not a uSockets socket this.socket.set(socket); - if !this.user_wants_ref.get() { + // Stale if node:net reconnected through this wrapper while it was paused. + this.update_flags(|f| f.remove(Flags::IS_PAUSED)); + if !this.ref_pollref_on_connect.replace(true) { this.poll_ref.with_mut(|p| p.unref(js_loop_ctx())); } jsc::mark_binding!(); @@ -3227,10 +3229,11 @@ impl NewSocket { _frame: &CallFrame, ) -> JsResult { jsc::mark_binding!(); - this.user_wants_ref.set(true); - // Not yet established: `connect_finish` holds the loop and `on_open` applies this. if this.socket.get().is_established() { this.poll_ref.with_mut(|p| p.ref_(js_loop_ctx())); + } else { + // `connect_finish` holds the loop until then; `on_open` applies this. + this.ref_pollref_on_connect.set(true); } Ok(JSValue::UNDEFINED) } @@ -3242,9 +3245,10 @@ impl NewSocket { _frame: &CallFrame, ) -> JsResult { jsc::mark_binding!(); - this.user_wants_ref.set(false); if this.socket.get().is_established() { this.poll_ref.with_mut(|p| p.unref(js_loop_ctx())); + } else { + this.ref_pollref_on_connect.set(false); } Ok(JSValue::UNDEFINED) } @@ -3586,7 +3590,7 @@ impl NewSocket { flags: Cell::new(initial_flags), this_value: JsCell::new(JsRef::empty()), poll_ref: JsCell::new(KeepAlive::init()), - user_wants_ref: Cell::new(this.user_wants_ref.get()), + ref_pollref_on_connect: Cell::new(true), buffered_data_for_node_net: JsCell::new(Vec::new()), bytes_written: Cell::new(0), native_callback: JsCell::new(NativeCallbacks::None), @@ -3693,7 +3697,7 @@ impl NewSocket { flags: Cell::new(Flags::BYPASS_TLS | Flags::IS_ACTIVE | Flags::OWNED_PROTOS), this_value: JsCell::new(JsRef::empty()), poll_ref: JsCell::new(KeepAlive::init()), - user_wants_ref: Cell::new(true), + ref_pollref_on_connect: Cell::new(true), buffered_data_for_node_net: JsCell::new(Vec::new()), bytes_written: Cell::new(0), native_callback: JsCell::new(NativeCallbacks::None), @@ -4813,7 +4817,7 @@ pub fn js_upgrade_duplex_to_tls( flags: Cell::new(initial_flags), this_value: JsCell::new(JsRef::empty()), poll_ref: JsCell::new(KeepAlive::init()), - user_wants_ref: Cell::new(true), + ref_pollref_on_connect: Cell::new(true), buffered_data_for_node_net: JsCell::new(Vec::new()), bytes_written: Cell::new(0), native_callback: JsCell::new(NativeCallbacks::None), diff --git a/src/uws_sys/socket.rs b/src/uws_sys/socket.rs index 77e01f62721b..cbe2190c716b 100644 --- a/src/uws_sys/socket.rs +++ b/src/uws_sys/socket.rs @@ -504,9 +504,12 @@ impl NewSocketHandler { // ── flow control / sockopts ───────────────────────────────────────────── + /// A connect that has not completed yet is left alone (like the + /// `connecting` arm): the open re-arms reads, so latching a pause here + /// would only make the next real `pause()` a no-op. pub fn pause_stream(&self) -> bool { on_socket!(self.socket; - connected s => { s.pause(); true }, + connected s => if s.is_established() { s.pause(); true } else { false }, connecting _c => false, detached => true, duplex _d => false, // TODO: pause/resume upgraded duplex @@ -516,7 +519,7 @@ impl NewSocketHandler { pub fn resume_stream(&self) -> bool { on_socket!(self.socket; - connected s => { s.resume(); true }, + connected s => if s.is_established() { s.resume(); true } else { false }, connecting _c => false, detached => true, duplex _d => false, // TODO: pause/resume upgraded duplex diff --git a/test/js/node/net/node-net.test.ts b/test/js/node/net/node-net.test.ts index 6089a25e9a12..4c8d5fc5b531 100644 --- a/test/js/node/net/node-net.test.ts +++ b/test/js/node/net/node-net.test.ts @@ -787,6 +787,46 @@ describe.concurrent("unref()/pause() around connect()", () => { expect(stderr).toBe(""); expect(exitCode).toBe(0); }); + + // A pause() that reached the native socket mid-connect used to latch its paused bit while the + // open re-armed reads, so backpressure could never pause it again and the buffer grew unbounded. + it("pause() while the connect is in flight still lets backpressure stop reads", async () => { + const chunk = Buffer.alloc(64 * 1024, "x"); + await using server = createServer(c => { + const pump = () => { + while (!c.destroyed && c.write(chunk)) {} + }; + c.on("drain", pump); + c.on("error", () => {}); + pump(); + }); + await once(server.listen(0, "127.0.0.1"), "listening"); + const s = new Socket(); + try { + s.connect((server.address() as any).port, "127.0.0.1"); + await new Promise(resolve => + process.nextTick(() => + process.nextTick(() => { + s.pause(); + resolve(); + }), + ), + ); + await once(s, "connect"); + // Once the buffer passes the high-water mark reads must stop. Unpatched, every loop turn + // delivered another recv; allow at most a couple that were already in flight. + const deadline = performance.now() + 5000; + while (performance.now() < deadline && s.readableLength < s.readableHighWaterMark) + await new Promise(r => setTimeout(r, 1)); + expect(s.readableLength).toBeGreaterThanOrEqual(s.readableHighWaterMark); + const settled = s.bytesRead; + for (let i = 0; i < 100; i++) await new Promise(r => setTimeout(r, 0)); + const recvBuffer = 512 * 1024; + expect(s.bytesRead - settled).toBeLessThanOrEqual(2 * recvBuffer); + } finally { + s.destroy(); + } + }); }); it("socket should keep process alive if unref is not called", async () => { From bfbfe8e557f507aac04b12f9795dc619cb3e6d58 Mon Sep 17 00:00:00 2001 From: Jarred Sumner Date: Fri, 21 Aug 2026 06:38:07 +0000 Subject: [PATCH 4/4] test: don't require the paused client to cross its high-water mark before checking reads stopped No-Verification-Needed: test-only change (node-net.test.ts) --- test/js/node/net/node-net.test.ts | 11 +++++++---- 1 file changed, 7 insertions(+), 4 deletions(-) diff --git a/test/js/node/net/node-net.test.ts b/test/js/node/net/node-net.test.ts index 4c8d5fc5b531..08f95fd145f0 100644 --- a/test/js/node/net/node-net.test.ts +++ b/test/js/node/net/node-net.test.ts @@ -792,9 +792,11 @@ describe.concurrent("unref()/pause() around connect()", () => { // open re-armed reads, so backpressure could never pause it again and the buffer grew unbounded. it("pause() while the connect is in flight still lets backpressure stop reads", async () => { const chunk = Buffer.alloc(64 * 1024, "x"); + let serverBackedUp = false; await using server = createServer(c => { const pump = () => { while (!c.destroyed && c.write(chunk)) {} + serverBackedUp = !c.destroyed; }; c.on("drain", pump); c.on("error", () => {}); @@ -813,12 +815,13 @@ describe.concurrent("unref()/pause() around connect()", () => { ), ); await once(s, "connect"); - // Once the buffer passes the high-water mark reads must stop. Unpatched, every loop turn - // delivered another recv; allow at most a couple that were already in flight. + // Once the client stops reading (buffer past the high-water mark, or never started) the + // server backs up; from then on bytesRead must stay put. Unpatched, every loop turn + // delivered another recv; allow a couple that were already in flight. const deadline = performance.now() + 5000; - while (performance.now() < deadline && s.readableLength < s.readableHighWaterMark) + while (performance.now() < deadline && !serverBackedUp && s.readableLength < s.readableHighWaterMark) await new Promise(r => setTimeout(r, 1)); - expect(s.readableLength).toBeGreaterThanOrEqual(s.readableHighWaterMark); + expect(serverBackedUp || s.readableLength >= s.readableHighWaterMark).toBeTrue(); const settled = s.bytesRead; for (let i = 0; i < 100; i++) await new Promise(r => setTimeout(r, 0)); const recvBuffer = 512 * 1024;