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
56 changes: 24 additions & 32 deletions src/js/node/_http_server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -843,9 +843,10 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort
if (!http_req[kReqShouldKeepAlive]) {
http_res[kMustCloseConnection] = true;
}
// on(), not once(): "finish" fires at most once per response and once()
// allocates a wrapper closure per request.
http_res.on("finish", emitResponseFinishHandleSocket);
// One plain on() listener (once() allocates a wrapper, a second listener
// deoptimizes every 'finish' emit), registered before the 'request' event
// like Node's resOnFinish so res.on-replacing middleware cannot swallow it.
http_res.on("finish", emitResponseFinish);

if (hasObserver("http")) {
startPerf(http_res, kServerResponseStatistics, {
Expand Down Expand Up @@ -1072,19 +1073,16 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort
if (handle.finished || didFinish) {
handle = undefined;
http_res[kCloseCallback] = undefined;
// Set in time only because end() defers the 'finish' emit to a
// process.nextTick (see ServerResponse.prototype.end) and nothing
// between the 'request' emit and here drains the tick queue.
http_res[kDispatcherDetached] = true;
http_res.detachSocket(socket);
if (socket[kPipelinedResponses] !== undefined) {
advanceResponsePipeline(server, socket);
}
return;
}
if (http_res.socket) {
// Detach the socket first, then advance the response pipeline on
// the connection it was on (the same order as the two listeners
// this replaces). One shared function instead of two per-request
// bind() closures.
http_res.on("finish", emitAsyncResponseFinish);
}

const { resolve, promise } = $newPromiseCapability(Promise);
resolveFunction = resolve;
Expand Down Expand Up @@ -1412,6 +1410,9 @@ const kPipelinedQueuedState = Symbol("kPipelinedQueuedState");
const kOutgoingData = Symbol("kOutgoingData");
const kReplayingPipelinedOps = Symbol("kReplayingPipelinedOps");
const kStopParsingOnCloseListener = Symbol("kStopParsingOnCloseListener");
// Set when the dispatcher already detached a synchronously-finished response,
// so the 'finish' listener does not detach/advance the pipeline a second time.
const kDispatcherDetached = Symbol("kDispatcherDetached");

// https://github.com/nodejs/node/blob/v26.3.0/lib/_http_server.js (socketOnError)
const badRequestResponse = Buffer.from(`HTTP/1.1 400 Bad Request\r\nConnection: close\r\n\r\n`, "latin1");
Expand Down Expand Up @@ -2424,29 +2425,18 @@ function stopServerResponsePerf(this: any) {
}
}

// `on("finish", fn.bind(server, socket, ...))` allocated a bound closure per
// request. Inside a "finish" listener `this` is the response; the connection
// socket is `this.req.socket`, with `this.socket` (assigned by assignSocket,
// cleared only by detachSocket) as the fallback for requests the stream
// destroyer already detached. A single shared function needs no per-request
// state at all.
function emitResponseFinishHandleSocket() {
// req.socket is nulled by the stream destroyer (pipeline/compose cleanup
// does `stream.socket = null` for server requests); the response's own
// socket - set by assignSocket and cleared only by detachSocket - still
// references the connection then.
// Node.js's resOnFinish as one shared listener: connection handling (close or
// arm keep-alive) runs first because onResponseFinishHandleSocket's guards
// read pre-detach state, then detach the socket and advance the pipeline.
function emitResponseFinish() {
// req.socket is nulled by the stream destroyer (pipeline/compose cleanup);
// the response's own socket (set by assignSocket, cleared only by
// detachSocket) still references the connection then.
const socket = this.req?.socket ?? this.socket;
onResponseFinishHandleSocket(socket?.server, socket, this);
}

// The async-response half of Node.js's resOnFinish. Detach before advancing,
// in the same order as the two separate listeners this replaces.
// advanceResponsePipeline already bails on a missing socket.
function emitAsyncResponseFinish() {
// Same destroyer-null fallback as emitResponseFinishHandleSocket: without
// it a destroyed request leaves socket._httpMessage assigned and the next
// kept-alive request fails with ERR_HTTP_SOCKET_ASSIGNED.
const socket = this.req?.socket ?? this.socket;
// The dispatcher detached a synchronously-finished response itself;
// advancing the pipeline again here would skip a queued response.
if (this[kDispatcherDetached]) return;
if (socket != null) this.detachSocket(socket);
advanceResponsePipeline(socket?.server, socket);
}
Expand Down Expand Up @@ -2586,7 +2576,6 @@ function advanceResponsePipeline(server, socket) {
res.assignSocket(socket);
}
socket[kRequest] = res.req;
res.on("finish", emitAsyncResponseFinish);

// Replay the writes buffered while the response was queued.
// The buffered bytes are handed to the native handle below, so they no
Expand Down Expand Up @@ -3225,6 +3214,9 @@ ServerResponse.prototype.end = function (chunk, encoding, callback) {
this.emit("prefinish");
this._callPendingCallbacks();

// Deferring the 'finish' emit to nextTick is load-bearing: the dispatcher
// sets kDispatcherDetached only after a sync-finished handler returns, so
// an emit before that would detach and advance the pipeline twice.
if (callback) {
process.nextTick(
function (callback, self) {
Expand Down
173 changes: 173 additions & 0 deletions test/regression/issue/34485.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,173 @@
import { expect, test } from "bun:test";
import { bunEnv, bunExe, tempDir } from "harness";

// https://github.com/oven-sh/bun/issues/34485
// Middleware like @polka/compression (used by `vite preview`) replaces res.on
// and ends the response asynchronously from a zlib stream event. The server's
// internal detach-on-finish listener must be registered before the 'request'
// event so the middleware cannot swallow it; otherwise the next keep-alive
// request on the same socket throws ERR_HTTP_SOCKET_ASSIGNED and the server
// process dies.
test.concurrent("keep-alive requests survive middleware that wraps res.on/write/end", async () => {
using dir = tempDir("issue-34485", {
"server.js": `
const http = require("http");
const zlib = require("zlib");

const BODY = Buffer.alloc(4096, "x").toString();

// The relevant parts of @polka/compression: route writes through a gzip
// stream, call the original end() from the gzip 'end' event, and divert
// later res.on() registrations onto the gzip stream.
function middleware(req, res) {
const { end, write, on } = res;
const compress = zlib.createGzip();
res.setHeader("Content-Encoding", "gzip");
compress.on("data", chunk => write.call(res, chunk) || compress.pause());
on.call(res, "drain", () => compress.resume());
compress.on("end", () => end.call(res));
res.write = function (chunk, enc) {
return compress.write(chunk, enc);
};
res.end = function (chunk, enc) {
return compress.end(chunk, enc);
};
res.on = function (type, listener) {
compress.on(type, listener);
return this;
};
}

const server = http.createServer((req, res) => {
middleware(req, res);
res.writeHead(200, { "content-type": "text/plain" });
res.end(BODY);
});

server.listen(0, "127.0.0.1", async () => {
const port = server.address().port;
const agent = new http.Agent({ keepAlive: true, maxSockets: 1 });
const sockets = new Set();

function get() {
return new Promise((resolve, reject) => {
http
.get({ host: "127.0.0.1", port, agent, headers: { "accept-encoding": "gzip" } }, res => {
sockets.add(res.socket);
const chunks = [];
res.on("data", c => chunks.push(c));
res.on("end", () => resolve([res.statusCode, Buffer.concat(chunks)]));
res.on("error", reject);
})
.on("error", reject);
});
}

for (let i = 0; i < 3; i++) {
const [status, body] = await get();
const text = zlib.gunzipSync(body).toString();
console.log(\`request \${i}: \${status} \${text.length} bytes ok=\${text === BODY}\`);
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
console.log("sockets used:", sockets.size);
agent.destroy();
server.close();
});
`,
});

await using proc = Bun.spawn({
cmd: [bunExe(), "server.js"],
env: bunEnv,
cwd: String(dir),
stderr: "pipe",
});

const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]);

expect({ stdout, stderr, exitCode }).toEqual({
stdout:
"request 0: 200 4096 bytes ok=true\n" +
"request 1: 200 4096 bytes ok=true\n" +
"request 2: 200 4096 bytes ok=true\n" +
"sockets used: 1\n",
stderr: "",
exitCode: 0,
});
Comment thread
coderabbitai[bot] marked this conversation as resolved.
});

// Same middleware pattern, but the requests are pipelined (all sent in one
// packet). Queued responses rely on the same pre-registered detach-on-finish
// listener to advance the pipeline; pre-fix this crashed identically.
test.concurrent("pipelined requests survive middleware that wraps res.on/write/end", async () => {
using dir = tempDir("issue-34485-pipelined", {
"server.js": `
const http = require("http");
const net = require("net");
const { PassThrough } = require("stream");

// Same shape as @polka/compression, with an identity stream so the
// response bytes stay directly assertable: writes routed through a
// stream, end() called from its 'end' event, res.on() diverted onto it.
function middleware(req, res) {
const { end, write, on } = res;
const compress = new PassThrough();
compress.on("data", chunk => write.call(res, chunk) || compress.pause());
on.call(res, "drain", () => compress.resume());
compress.on("end", () => end.call(res));
res.write = function (chunk, enc) {
return compress.write(chunk, enc);
};
res.end = function (chunk, enc) {
return compress.end(chunk, enc);
};
res.on = function (type, listener) {
compress.on(type, listener);
return this;
};
}

const server = http.createServer((req, res) => {
middleware(req, res);
res.writeHead(200, { "content-type": "text/plain" });
res.end("body-of" + req.url + "|");
});

server.listen(0, "127.0.0.1", () => {
const port = server.address().port;
const sock = net.connect(port, "127.0.0.1", () => {
sock.write(
"GET /1 HTTP/1.1\\r\\nHost: a\\r\\n\\r\\n" +
"GET /2 HTTP/1.1\\r\\nHost: a\\r\\n\\r\\n" +
"GET /3 HTTP/1.1\\r\\nHost: a\\r\\nConnection: close\\r\\n\\r\\n",
);
});
let data = "";
sock.setEncoding("latin1");
sock.on("data", c => (data += c));
sock.on("end", () => {
const statuses = data.split("HTTP/1.1 200").length - 1;
const i1 = data.indexOf("body-of/1|");
const i2 = data.indexOf("body-of/2|");
const i3 = data.indexOf("body-of/3|");
console.log("statuses:", statuses, "ordered:", i1 >= 0 && i1 < i2 && i2 < i3);
server.close();
});
});
`,
});

await using proc = Bun.spawn({
cmd: [bunExe(), "server.js"],
env: bunEnv,
cwd: String(dir),
stderr: "pipe",
});

const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]);

expect({ stdout, stderr, exitCode }).toEqual({
stdout: "statuses: 3 ordered: true\n",
stderr: "",
exitCode: 0,
});
});
Loading