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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 7 additions & 9 deletions src/js/node/net.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
Expand Down Expand Up @@ -4272,8 +4270,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?.();
}
}

Expand Down
12 changes: 4 additions & 8 deletions src/runtime/socket/Listener.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1619,6 +1619,10 @@ fn connect_finish<const IS_SSL: bool>(
f.set(SocketFlags::ALLOW_HALF_OPEN, allow_half_open);
socket_ref.flags.set(f);
}
// 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()));
// 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
Expand Down Expand Up @@ -1687,14 +1691,6 @@ fn connect_finish<const IS_SSL: bool>(
}
};

// 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 {
Expand Down
33 changes: 14 additions & 19 deletions src/runtime/socket/socket_body.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1423,6 +1423,11 @@ impl<const SSL: bool> NewSocket<SSL> {
// 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);
// 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!();

// Add SNI support for TLS (mongodb and others requires this)
Expand Down Expand Up @@ -3220,41 +3225,31 @@ impl<const SSL: bool> NewSocket<SSL> {
#[bun_jsc::host_fn(method)]
pub(crate) fn js_ref(
this: &Self,
global: &JSGlobalObject,
_global: &JSGlobalObject,
_frame: &CallFrame,
) -> JsResult<JSValue> {
jsc::mark_binding!();
if this.socket.get().is_detached() {
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);
}
if this.socket.get().is_detached() {
return Ok(JSValue::UNDEFINED);
}
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<JSValue> {
jsc::mark_binding!();
if this.socket.get().is_detached() {
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);
}
Comment thread
claude[bot] marked this conversation as resolved.
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)
}

Expand Down
7 changes: 5 additions & 2 deletions src/uws_sys/socket.rs
Original file line number Diff line number Diff line change
Expand Up @@ -504,9 +504,12 @@ impl<const IS_SSL: bool> NewSocketHandler<IS_SSL> {

// ── 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
Expand All @@ -516,7 +519,7 @@ impl<const IS_SSL: bool> NewSocketHandler<IS_SSL> {

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
Expand Down
119 changes: 108 additions & 11 deletions test/js/node/net/node-net.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -699,9 +699,9 @@
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<void>(resolve => server.listen(0, "127.0.0.1", resolve));
try {
Expand All @@ -714,27 +714,124 @@
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",
stderr: "pipe",
});
// 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"]);
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);
Comment thread
coderabbitai[bot] marked this conversation as resolved.
} 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: "pipe",
});
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() 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) => {
Comment thread
coderabbitai[bot] marked this conversation as resolved.
const { stdout, stderr, exitCode } = await run(client);
expect(stdout).toBe("connected\n");
expect(stderr).toBe("");
expect(exitCode).toBe(0);
});

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();")}`],

Check warning on line 783 in test/js/node/net/node-net.test.ts

View check run for this annotation

Claude / Claude Code Review

Missing 'unref(); ref()' before-connect() row (robobun feedback unaddressed)

The `ref() after unref()` it.each only covers 'right after connect()' and 'while the connect is in flight' — robobun's last comment (closing #37087) asked for the missing `s.unref(); s.ref();` *before* `connect()` row, which exercises the distinct `_handle === null` path (kUserUnrefed + once('connect') handlers + initSocketHandle) rather than js_ref/js_unref on a native handle. One more row: `["before connect()", 's.unref(); s.ref(); s.resume(); s.connect(port, "127.0.0.1");']`.
Comment on lines +780 to +783

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 The ref() after unref() it.each only covers 'right after connect()' and 'while the connect is in flight' — robobun's last comment (closing #37087) asked for the missing s.unref(); s.ref(); before connect() row, which exercises the distinct _handle === null path (kUserUnrefed + once('connect') handlers + initSocketHandle) rather than js_ref/js_unref on a native handle. One more row: ["before connect()", 's.unref(); s.ref(); s.resume(); s.connect(port, "127.0.0.1");'].

Extended reasoning...

What's missing

The ref() after unref() matrix at test/js/node/net/node-net.test.ts:781-784 has two rows:

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", ...);

There is no before connect() row. The sibling matrix immediately above (lines 767-774) covers all three timings — before/after/in-flight — for both unref() and pause(), so this matrix is asymmetric.

Why it was requested

robobun's comment at 2026-08-21T06:33:41Z is the last entry in the PR timeline (after the final commit bfbfe8e) and explicitly asks for this row:

One test timing from that PR is not in the unref()/pause() around connect() block here: s.unref(); s.ref(); with both calls before connect(), where the server ends the connection. The expected output is connected then closed, since the later ref() has to win over the earlier unref(). The two ref() after unref() cases here both run after connect(). The before-connect() timing could be one more row in that it.each list.

This comment is unaddressed in the final diff.

Why the before-connect timing is a distinct code path

The two existing rows both run after connect() has been called, so _handle already exists (detached right-after, or a non-established native socket in-flight). In both cases unref()/ref() reach the native js_unref/js_ref (socket_body.rs:3225-3252), which — after this PR's change — toggle ref_pollref_on_connect on the not-yet-established handle and on_open applies it.

Before connect(), _handle is null. Socket.prototype.unref() and Socket.prototype.ref() in net.ts instead:

  • toggle this[kUserUnrefed] (true then false),
  • each register a once('connect', ...) handler,
  • and later initSocketHandle reads the final kUserUnrefed state when the handle is created (net.ts:4273-4274 — if (self[kUserUnrefed] || self[kPausedUnref]) handle.unref?.()), then both once-handlers fire in registration order on 'connect'.

None of that JS-side bookkeeping is exercised by the after-connect / in-flight rows.

Step-by-step through the missing row

With s.unref(); s.ref(); s.resume(); s.connect(port, "127.0.0.1"); and the server calling c.end():

  1. s.unref(): _handle is null → kUserUnrefed = true, once('connect', h1) registered.
  2. s.ref(): _handle is null → kUserUnrefed = false, once('connect', h2) registered.
  3. s.connect(...): creates _handle, calls initSocketHandle → kUserUnrefed is false and kPausedUnref is false, so handle.unref() is not called. connect_finish refs poll_ref for the connect attempt.
  4. Connection opens → on_open sees ref_pollref_on_connect == true (never touched) → poll_ref stays held. 'connect' fires → h1 (unref-era) and h2 (ref-era) run in order; the last one wins.
  5. Server ends → client emits 'close' → stdout is connected\nclosed\n, exit 0.

That's the expected output the existing assertion already checks — so adding the row is a one-line change.

Impact and fix

Per REVIEW.md "Cover the variant matrix, not just the repro": every sibling entry point receiving the same fix should be covered, and this timing variant was explicitly requested by a reviewer. The code path likely works correctly (the final kUserUnrefed state is what initSocketHandle reads), but leaving the row out means the pre-handle ref()-wins-over-unref() ordering has no coverage.

Add one row to the it.each at line 781:

["before connect()", `s.unref(); s.ref(); s.resume(); s.connect(port, "127.0.0.1");`],

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@robobun add this in a follow-up PR

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Added the before connect() row in #40134.

])("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);
});

// 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");
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", () => {});
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<void>(resolve =>
process.nextTick(() =>
process.nextTick(() => {
s.pause();
resolve();
}),
),
);
await once(s, "connect");
// 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 && !serverBackedUp && s.readableLength < s.readableHighWaterMark)
await new Promise(r => setTimeout(r, 1));
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;
expect(s.bytesRead - settled).toBeLessThanOrEqual(2 * recvBuffer);
} finally {
s.destroy();
}
});
});

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")],
Expand Down
Loading