Skip to content
Open
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
31 changes: 29 additions & 2 deletions src/js/node/net.ts
Original file line number Diff line number Diff line change
Expand Up @@ -395,6 +395,25 @@ function endNT(socket, callback, err) {
function emitCloseNT(self, hasError) {
self.emit("close", hasError);
}
// A write that waits for 'connect' or for the TLS handle still holds its chunk in _pendingData: the native handle never had it.
function takeInFlightWrite(self) {
const callback = self[kwriteCallback];
if (!callback || self._pendingData != null) return null;
self[kwriteCallback] = null;
return callback;
}
// uv_close() fails a queued write with UV_ECANCELED: https://github.com/nodejs/node/blob/v26.3.0/deps/uv/src/unix/stream.c#L464
function cancelWriteNT(callback) {
callback(new ErrnoException(uv().UV_ECANCELED, "write"));
}
// The cancel follows the 'error' tick that callback(err) queues. A destroy(err, cb) callback that throws must not lose it.
function finishDestroy(callback, err, canceledWrite) {
try {
callback(err);
} finally {
if (canceledWrite && err) process.nextTick(cancelWriteNT, canceledWrite);
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}
}
// Shared-fd TLS pair teardown: mirrors node's close ordering, where the
// close-callbacks phase runs after the check phase (lib/net.js close path in
// node v26.3.0), so destroy()-time setImmediates still see the pair alive.
Expand Down Expand Up @@ -2357,6 +2376,11 @@ Socket.prototype._destroy = function _destroy(err, callback) {
$debug("Socket.prototype._destroy");

this.connecting = false;
// Taken before anything closes the handle: the native close handler fails a write it still finds with ERR_SOCKET_CLOSED.
const canceledWrite = takeInFlightWrite(this);
// Node: after 'error', before 'close'. With no error it goes first, so the stream is errored before the EOF that the native close handler pushes can emit 'end'.
if (canceledWrite && !err) process.nextTick(cancelWriteNT, canceledWrite);

// Tear down a wrapped generic duplex with this socket: the native handle's
// close only flushes close_notify and lets the wrapper drain; without an
// explicit destroy here a late RST on the underlying transport can surface
Expand Down Expand Up @@ -2429,10 +2453,10 @@ Socket.prototype._destroy = function _destroy(err, callback) {
this._handle = null;
this._sockname = null;
}
callback(err);
finishDestroy(callback, err, canceledWrite);
closeOwedRaw(this, upgraded);
} else {
callback(err);
finishDestroy(callback, err, canceledWrite);
closeOwedRaw(this, upgraded);
process.nextTick(emitCloseNT, this, err ? true : false);
}
Expand Down Expand Up @@ -2958,6 +2982,9 @@ Socket.prototype._write = function _write(chunk, encoding, callback) {
this._pendingData = chunk;
this._pendingEncoding = encoding;
function onClose() {
// A wrapped socket opens without 'connect', so this listener outlives the wait and the teardown may have settled the write.
if (this[kwriteCallback] !== callback) return;
this[kwriteCallback] = null;
callback($ERR_SOCKET_CLOSED_BEFORE_CONNECTION());
}
this.once("connect", function connect() {
Expand Down
27 changes: 27 additions & 0 deletions test/harness.ts
Original file line number Diff line number Diff line change
Expand Up @@ -134,6 +134,33 @@ export function nodeExe(): string | null {
return which("node") || null;
}

let systemNodeMajor: number | undefined;

/** Major version of the system Node.js, or 0 when there is none. Spawns it one time. */
export function nodeMajorVersion(): number {
if (systemNodeMajor === undefined) {
const node = nodeExe();
const version = node ? Bun.spawnSync({ cmd: [node, "-p", "process.versions.node"], env: bunEnv }).stdout : "";
systemNodeMajor = parseInt(version.toString(), 10) || 0;
}
return systemNodeMajor;
}

/**
* A `describe.each` table of `[name, executable]` for a fixture that must give the same output
* under Bun and under Node.js. The system Node.js is in it when its major version is at least
* `minNodeMajor`. `nodeOnWindows: false` leaves it out on Windows, for a fixture whose
* precondition Node cannot build there.
*/
export function runtimesWithNode(
minNodeMajor: number,
{ nodeOnWindows = true }: { nodeOnWindows?: boolean } = {},
): [name: string, exe: string][] {
const runtimes: [string, string][] = [["bun", bunExe()]];
if ((nodeOnWindows || !isWindows) && nodeMajorVersion() >= minNodeMajor) runtimes.push(["node", nodeExe()!]);
return runtimes;
}

let abiMatchingNode: Promise<string> | undefined;

/**
Expand Down
128 changes: 128 additions & 0 deletions test/js/node/net/node-net.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,9 @@ import {
gc,
isASAN,
isDebug,
isLinux,
isWindows,
runtimesWithNode,
tempDir,
tls as tlsCert,
tmpdirSync,
Expand Down Expand Up @@ -1973,6 +1975,132 @@ describe("paused socket whose peer sends RST", () => {
});
});

// libuv cancels a write that is still in flight when its handle closes. Node gives
// the callback UV_ECANCELED, after 'error' and before 'close'.
// https://github.com/nodejs/node/blob/v26.3.0/lib/internal/stream_base_commons.js#L86-L90
describe("socket torn down with a write still in flight", () => {
// `holder` writes to a peer that never reads until the kernel stops taking the bytes, and
// queues one more write behind that one. Then the teardown runs. Prints what `holder`
// reports from there on, in order, up to 'close'.
const fixture = /* js */ `
const net = require("node:net");
const { SIDE, TEARDOWN, BATCH } = process.env;
const shape = err => (err ? [err.code, err.syscall].filter(Boolean).join(" ") : "ok");
const boom = Object.assign(new Error("boom"), { code: "EBOOM" });
const teardowns = {
"peer reset": (holder, peer) => peer.resetAndDestroy(),
"destroy()": holder => holder.destroy(),
"destroy(err)": holder => holder.destroy(boom),
"destroy(err, cb) with a cb that throws": holder =>
holder.destroy(boom, () => {
throw new Error("from the destroy callback");
}),
};

let client, accepted, ready = 0;
// A write made inside the 'connect' or 'connection' dispatch is flushed on another native path.
const onReady = () => ++ready === 2 && setImmediate(run);
const server = net.createServer(socket => {
accepted = socket;
onReady();
});
server.listen(0, "127.0.0.1", () => {
client = net.connect(server.address().port, "127.0.0.1", onReady);
});

function run() {
const [holder, peer] = SIDE === "client" ? [client, accepted] : [accepted, client];
const events = [];
peer.on("error", () => {});
holder.on("error", err => events.push("error " + shape(err)));
// Node emits no 'end' here. One that ran before the write callback made node:http's
// client report 'socket hang up' first and drop the request's 'finish'.
holder.on("end", () => events.push("end"));
holder.on("close", hadError => {
events.push("close " + hadError);
console.log(JSON.stringify(events));
peer.destroy();
server.close();
});
// A plain TCP write that the kernel takes whole is done when write() returns, so
// writableLength stays 0 until one write, or one corked batch, is left in flight: the last one.
// It is 64 MB so that the kernel cannot finish it: libuv sends two times inside one write(),
// and on macOS the second send can take the rest of a 1 MB chunk before the teardown runs.
const chunk = Buffer.alloc((BATCH ? 32 : 64) * 1024 * 1024);
let writes = 0;
while (writes < 8 && holder.writableLength === 0) {
const nth = ++writes;
const callback = err => {
if (nth === writes) events.push("write " + shape(err));
};
if (BATCH) holder.cork();
holder.write(chunk, callback);
if (BATCH) {
holder.write(chunk, callback);
holder.uncork();
}
}
if (holder.writableLength === 0) throw new Error("no write stayed in flight");
holder.write("b", err => events.push("queued " + shape(err)));
teardowns[TEARDOWN](holder, peer);
}
`;

// Node 22.0 still called a canceled write's callback without an error, so only Node 24 or later is a reference.
// Not on Windows: there a 64 MB write to a peer that never reads completes, and a write that is still
// pending when write() returns can complete a moment later, so Node cannot pin a write in flight.
describe.each(runtimesWithNode(24, { nodeOnWindows: false }))("%s", (_, exe) => {
async function run(env: Record<string, string>) {
await using proc = Bun.spawn({
cmd: [exe, "-e", fixture],
env: { ...bunEnv, ...env },
stdout: "pipe",
stderr: "pipe",
});
const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]);
return { events: stdout.trim() ? JSON.parse(stdout) : stdout, stderr, exitCode };
}
const reports = (events: string[]) => ({ events, stderr: "", exitCode: 0 });

describe.each(["client", "server"])("%s", SIDE => {
// Which comes first, the reset or a writable event, is only known on Linux: a loopback RST
// arrives before the close() that sends it returns. On Windows the writable event wins and
// the write path settles the write.
it.concurrent.skipIf(!isLinux)("peer reset", async () => {
expect(await run({ SIDE, TEARDOWN: "peer reset" })).toEqual(
reports(["error ECONNRESET read", "write ECANCELED write", "queued ECONNRESET read", "close true"]),
);
});

it.concurrent("destroy()", async () => {
expect(await run({ SIDE, TEARDOWN: "destroy()" })).toEqual(
reports(["write ECANCELED write", "queued ECANCELED write", "close false"]),
);
});

it.concurrent("destroy(err)", async () => {
expect(await run({ SIDE, TEARDOWN: "destroy(err)" })).toEqual(
reports(["error EBOOM", "write ECANCELED write", "queued EBOOM", "close true"]),
);
});

// The stream swallows the throw before it queues 'error'. The writes still settle.
it.concurrent("destroy(err, cb) with a cb that throws", async () => {
expect(await run({ SIDE, TEARDOWN: "destroy(err, cb) with a cb that throws" })).toEqual(
reports(["write ECANCELED write", "queued EBOOM", "close true"]),
);
});

// A corked batch reaches the socket through _writev. Both of its chunks report.
it.concurrent("destroy() with a corked batch in flight", async () => {
expect(await run({ SIDE, TEARDOWN: "destroy()", BATCH: "1" })).toEqual(
reports(["write ECANCELED write", "write ECANCELED write", "queued ECANCELED write", "close false"]),
);
});
});
});
});

// Node stops kernel reads when push() returns false, not on pause(): a paused
// socket keeps reading into its buffer, so it still sees the peer's FIN.
// https://github.com/nodejs/node/blob/v26.3.0/lib/net.js#L817-L827
Expand Down
127 changes: 126 additions & 1 deletion test/js/node/tls/node-tls-connect.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@ import { heapStats } from "bun:jsc";
import { describe, expect, it } from "bun:test";
import { once } from "events";
import { writeFileSync } from "fs";
import { bunEnv, bunExe, tls as COMMON_CERT_, isASAN, nodeExe, tempDir } from "harness";
import { bunEnv, bunExe, tls as COMMON_CERT_, isASAN, isLinux, nodeExe, runtimesWithNode, tempDir } from "harness";
import https from "https";
import net from "net";
import { join } from "path";
Expand Down Expand Up @@ -1938,6 +1938,131 @@ describe.concurrent("a client paused while connecting", () => {
});
});

// Node cancels a write that is still in flight when the socket is destroyed: its callback gets
// UV_ECANCELED, one time, after 'error' and before 'close'.
describe("TLS socket torn down with a write still in flight", () => {
// Prints the events of the client and every call of the callback that Writable hands to _write.
const fixture = /* js */ `
const CERT = ${JSON.stringify(COMMON_CERT_)};
const net = require("node:net");
const tls = require("node:tls");
const shape = err => (err ? [err.code, err.syscall].filter(Boolean).join(" ") : "ok");
const events = [];
const calls = [];
const sockets = [];
let client;

// Counts the calls of the callback that Writable hands to _write. _write can pass it to itself again.
function countWriteCallbackCalls(socket) {
const counted = new WeakSet();
const _write = socket._write;
socket._write = function (chunk, encoding, callback) {
if (!counted.has(callback)) {
const inner = callback;
callback = function (err) {
calls.push(shape(err));
return inner.apply(this, arguments);
};
counted.add(callback);
}
return _write.call(this, chunk, encoding, callback);
};
}

function watch(server) {
countWriteCallbackCalls(client);
client.on("error", err => events.push("error " + shape(err)));
client.on("end", () => events.push("end"));
client.on("close", hadError => {
events.push("close " + hadError);
// The other 'close' listeners run first: _write adds one for a write that waits for 'connect'.
setImmediate(() => {
console.log(JSON.stringify({ events, calls }));
for (const socket of sockets) socket.destroy();
server.close();
});
});
}

if (process.env.CASE === "established") {
const server = tls.createServer(CERT, peer => {
sockets.push(peer);
peer.on("error", () => {});
});
server.listen(0, "127.0.0.1", () => {
client = tls.connect({ port: server.address().port, host: "127.0.0.1", rejectUnauthorized: false });
watch(server);
client.on("secureConnect", () =>
setImmediate(() => {
// The peer never reads, so the kernel cannot take all of this.
client.write(Buffer.alloc(64 * 1024 * 1024, "a"), err => events.push("write " + shape(err)));
// A TLS write that the kernel took whole calls back on the next tick.
setImmediate(() => {
if (events.length > 0) throw new Error("the write did not stay in flight: " + events);
client.destroy();
});
}),
);
});
} else {
// This peer resets when it has the ClientHello, so the write is still behind the handshake.
const server = net.createServer(peer => {
sockets.push(peer);
peer.on("error", () => {});
peer.once("data", () => peer.resetAndDestroy());
});
server.listen(0, "127.0.0.1", () => {
const raw = net.connect(server.address().port, "127.0.0.1");
sockets.push(raw);
raw.on("error", () => {});
client = tls.connect({ socket: raw, rejectUnauthorized: false });
watch(server);
if (!client.connecting) throw new Error("the TLS socket is not connecting");
client.write("hello", err => events.push("write " + shape(err)));
});
}
`;

// Node 22.0 still called a canceled write's callback without an error, so only Node 24 or later is a reference.
// Not on Windows: there a 64 MB write to a peer that never reads completes, and a write that is still
// pending when write() returns can complete a moment later, so Node cannot pin a write in flight.
describe.each(runtimesWithNode(24, { nodeOnWindows: false }))("%s", (_, exe) => {
async function run(CASE: string) {
await using proc = Bun.spawn({
cmd: [exe, "-e", fixture],
env: { ...bunEnv, CASE },
stdout: "pipe",
stderr: "pipe",
});
const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]);
return { report: stdout.trim() ? JSON.parse(stdout) : stdout, stderr, exitCode };
}

// The native TLS close can finish after 'close' was emitted, so the write cannot wait for it.
it.concurrent("destroy() on an established connection", async () => {
expect(await run("established")).toEqual({
report: { events: ["write ECANCELED write", "close false"], calls: ["ECANCELED write"] },
stderr: "",
exitCode: 0,
});
});

// The write was made while the wrapped socket was still connecting. A wrapped socket opens
// without 'connect', so the 'close' listener of that wait is still there at teardown.
// Linux only: see "peer reset" in node-net.test.ts.
it.concurrent.skipIf(!isLinux)("peer reset during the handshake", async () => {
expect(await run("connecting")).toEqual({
report: {
events: ["error ECONNRESET read", "write ECANCELED write", "close true"],
calls: ["ECANCELED write"],
},
stderr: "",
exitCode: 0,
});
});
});
});

// #40653: a TLS 1.3 client must send its final handshake flight and the first
// write issued from the 'secureConnect' callback in ONE TCP segment, like
// Node does through its memory BIO. Two segments let a server that tears the
Expand Down
Loading