Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
ba93cd5
node:http: answer every request past maxRequestsPerSocket with its ow…
robobun Jul 6, 2026
20f7169
test: cover over-limit requests that carry a body
robobun Jul 6, 2026
3bf1fc8
test: move maxRequestsPerSocket coverage into its own file
robobun Jul 6, 2026
0a8414e
test: drop an inaccurate file header
robobun Jul 6, 2026
8802411
test: pin the node:http proxy test to 127.0.0.1
robobun Jul 6, 2026
3adf374
test: keep maxRequestsPerSocket coverage in node-http.test.ts
robobun Jul 6, 2026
e2b4ca4
test: fold the two socket-collection loops into one helper
robobun Jul 6, 2026
52ab594
node:http: end the connection after the dropped requests are answered
robobun Jul 6, 2026
3a9ffbc
test: drop the now-dead sequential mode from sendRequests
robobun Jul 6, 2026
45698bc
ci: retrigger
robobun Jul 6, 2026
b1c8c2b
Merge origin/main into farm/65f58994/http-max-requests-per-socket-pip…
robobun Jul 17, 2026
a12166a
test: cover over-limit 503s queued behind an async response
robobun Jul 17, 2026
29730a3
[autofix.ci] apply automated fixes
autofix-ci[bot] Jul 17, 2026
ab4a04c
node:http: close over-limit connections on pipeline drain, not on a m…
robobun Jul 18, 2026
7cf9f9b
Merge remote-tracking branch 'origin/main' into farm/65f58994/http-ma…
robobun Jul 18, 2026
480216e
test: fail fast and clean up when the later-read test's socket closes…
robobun Jul 18, 2026
49e01f9
style: cap the kEndAfterDroppedRequests comment at 3 lines
robobun Jul 18, 2026
1721bc8
style: cap remaining comments in the diff at 3 lines
robobun Jul 18, 2026
30c1c77
style: cap the sendRequests doc comment at 3 lines
robobun Jul 18, 2026
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
62 changes: 37 additions & 25 deletions src/js/node/_http_server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -879,12 +879,9 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort
const requestCount = (socket._requestCount || 0) + 1;
socket._requestCount = requestCount;
http_res._maxRequestsPerSocket = server.maxRequestsPerSocket;
// At (or beyond) the limit the response advertises Connection:
// close, like Node.js - including the over-limit 503 dropRequest
// answer, which would otherwise claim keep-alive right before the
// socket is destroyed. Closing the socket here instead would race
// already-pipelined requests, which still need to be dispatched so
// they can be answered with 503 via dropRequest.
// At (or beyond) the limit the response advertises Connection: close,
// like Node.js - including the over-limit 503 dropRequest answer,
// which would otherwise claim keep-alive.
http_res.maxRequestsOnConnectionReached = server.maxRequestsPerSocket <= requestCount;
if (server.maxRequestsPerSocket < requestCount) {
reachedRequestsLimit = true;
Expand Down Expand Up @@ -975,15 +972,15 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort
}

if (reachedRequestsLimit) {
// Like Node.js's parserOnIncoming: every request past the limit gets
// its own 'dropRequest' + 503. Closing the socket now (or on this
// response's "finish") would drop the requests pipelined behind it.
server.emit("dropRequest", http_req, socket);
http_res.writeHead(503);
http_res.end();
Comment thread
claude[bot] marked this conversation as resolved.
if (isPipelined) {
// The 503 is queued behind the in-flight responses; the connection
// closes once it has been written, like Node.js (Connection: close).
http_res[kMustCloseConnection] = true;
} else {
socket.destroy();
if (!socket[kEndAfterDroppedRequests]) {
socket[kEndAfterDroppedRequests] = true;
setImmediate(endSocketAfterDroppedRequests, socket);
}
Comment thread
claude[bot] marked this conversation as resolved.
} else if (is_upgrade) {
// Hand the raw socket over to the 'upgrade' listener, like Node.js.
Expand Down Expand Up @@ -1079,10 +1076,8 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort
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.
// Detach then advance the response pipeline on this connection (same
// order as the two listeners this replaces, no per-request bind()s).
http_res.on("finish", emitAsyncResponseFinish);
}

Expand Down Expand Up @@ -2410,6 +2405,23 @@ function renderNativeHeaders(res) {
}

const kMustCloseConnection = Symbol("kMustCloseConnection");
// `true` = deferred end scheduled; the symbol itself = it fired (read fully
// parsed), so end on pipeline drain. Two states stop a 503's own "finish" in
// per-request drainMicrotasks() from closing the socket mid-dispatch.
const kEndAfterDroppedRequests = Symbol("kEndAfterDroppedRequests");

// Runs once the current read has been fully parsed. With nothing left in flight
// or queued, every pipelined request has been answered: end the connection.
// Otherwise onResponseFinishHandleSocket ends it once the pipeline drains.
function endSocketAfterDroppedRequests(socket) {
if (!socket || socket.destroyed || socket.writableEnded) return;
socket[kEndAfterDroppedRequests] = kEndAfterDroppedRequests;
if (socket[kPipelinedResponses]?.length) return;
const inFlight = socket._httpMessage;
if (inFlight && !inFlight.finished) return;
socket.end();
}

function stopServerResponsePerf(this: any) {
if (this[kServerResponseStatistics] && hasObserver("http")) {
stopPerf(this, kServerResponseStatistics, {
Expand All @@ -2424,16 +2436,9 @@ 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
// req.socket is nulled by the stream destroyer (pipeline/compose cleanup);
// this.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);
Expand Down Expand Up @@ -2473,6 +2478,13 @@ function onResponseFinishHandleSocket(server, socket, res) {
if (socket[kPipelinedResponses]?.length) {
return;
}
// A request past maxRequestsPerSocket was dropped on this connection, the
// current read has been fully parsed, and the response pipeline has now
// drained: end it, deferred so every queued 503 could be written first.
if (socket[kEndAfterDroppedRequests] === kEndAfterDroppedRequests) {
socket.end();
return;
}
const rawKeepAliveTimeout = server.keepAliveTimeout;
const keepAliveTimeout = Number.isFinite(rawKeepAliveTimeout) && rawKeepAliveTimeout >= 0 ? rawKeepAliveTimeout : 0;
const rawKeepAliveBuffer = server.keepAliveTimeoutBuffer;
Expand Down
7 changes: 5 additions & 2 deletions test/js/node/http/node-http-proxy.js
Original file line number Diff line number Diff line change
Expand Up @@ -32,12 +32,15 @@ export async function run() {
req.pipe(proxyRequest); // Use pipe instead of manual data handling
});

proxyServer.listen(0, "localhost", async () => {
// 127.0.0.1 rather than "localhost": where the resolver prefers IPv6 the
// server binds ::1 while the client dials 127.0.0.1 and the connect is
// refused. exampleSite() already binds 127.0.0.1.
proxyServer.listen(0, "127.0.0.1", async () => {
const address = proxyServer.address();

const options = {
protocol: "http:",
hostname: "localhost",
hostname: "127.0.0.1",
port: address.port,
path: "/", // Change path to /
headers: {
Expand Down
221 changes: 208 additions & 13 deletions test/js/node/http/node-http.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3657,29 +3657,54 @@ it("Expect: 100-Continue matches case-insensitively like Node.js", async () => {
}
});

it("the over-limit 503 advertises Connection: close, not keep-alive", async () => {
const rawGet = (path: string) => `GET ${path} HTTP/1.1\r\nHost: x\r\n\r\n`;
const rawPost = (path: string, body: string) =>
`POST ${path} HTTP/1.1\r\nHost: x\r\nContent-Length: ${body.length}\r\n\r\n${body}`;

// Responses whose header block has fully arrived.
const responsesIn = (raw: string) => Math.min(raw.split("HTTP/1.1 ").length - 1, raw.split("\r\n\r\n").length - 1);

// Pipelines every request in one segment. "responses" hangs up once all arrive;
// "serverClose" waits for the server to. Resolving on "close" either way turns
// a mid-pipeline teardown into a truncated-bytes assertion instead of a hang.
function sendRequests(
port: number,
requests: string[],
until: "responses" | "serverClose" = "responses",
): Promise<{ raw: string; serverEnded: boolean }> {
const { promise, resolve, reject } = Promise.withResolvers<{ raw: string; serverEnded: boolean }>();
const socket = connect(port, "127.0.0.1");
let data = "";
let serverEnded = false;
socket.on("data", chunk => {
data += chunk;
if (until === "responses" && responsesIn(data) >= requests.length) {
socket.destroy();
resolve({ raw: data, serverEnded });
}
});
socket.on("end", () => (serverEnded = true));
socket.on("close", () => resolve({ raw: data, serverEnded }));
socket.on("error", reject);
socket.on("connect", () => socket.write(requests.join("")));
return promise;
}

it.concurrent("the over-limit 503 advertises Connection: close, not keep-alive", async () => {
// Node sets maxRequestsOnConnectionReached unconditionally
// (maxRequestsPerSocket <= count), so the dropRequest 503 carries
// Connection: close instead of advertising keep-alive right before the
// socket is destroyed.
// Connection: close instead of advertising keep-alive.
const server = createServer((req, res) => res.end("ok"));
server.maxRequestsPerSocket = 1;
try {
server.listen(0, "127.0.0.1");
await once(server, "listening");
const { port } = server.address() as AddressInfo;

const out = await new Promise<string>((resolve, reject) => {
const socket = connect(port, "127.0.0.1");
let data = "";
socket.on("data", chunk => (data += chunk));
socket.on("close", () => resolve(data));
socket.on("error", reject);
// Two pipelined requests: the second exceeds maxRequestsPerSocket.
socket.write("GET / HTTP/1.1\r\nHost: x\r\n\r\n" + "GET / HTTP/1.1\r\nHost: x\r\n\r\n");
});
// Two pipelined requests: the second exceeds maxRequestsPerSocket.
const { raw } = await sendRequests(port, [rawGet("/a"), rawGet("/b")]);

const second = out.slice(out.indexOf("HTTP/1.1 503"));
const second = raw.slice(raw.indexOf("HTTP/1.1 503"));
expect(second).toContain("HTTP/1.1 503");
expect(second).toContain("Connection: close");
expect(second).not.toContain("keep-alive");
Expand All @@ -3688,6 +3713,176 @@ it("the over-limit 503 advertises Connection: close, not keep-alive", async () =
}
});

it.concurrent("every pipelined request past maxRequestsPerSocket gets its own 503 and dropRequest", async () => {
// Node answers each over-limit request with a 503 and emits 'dropRequest'
// for it, so a client that pipelined several requests at once is not left
// with a torn-down connection and no response.
const drops: { url: string; isIncomingMessage: boolean; socket: unknown }[] = [];
const server = createServer((req, res) => {
req.resume();
res.end("ok");
});
server.maxRequestsPerSocket = 1;
server.on("dropRequest", (req, socket) =>
drops.push({ url: req.url, isIncomingMessage: req instanceof IncomingMessage, socket }),
);
try {
server.listen(0, "127.0.0.1");
await once(server, "listening");
const { port } = server.address() as AddressInfo;

// Three requests pipelined into one segment: /b and /c are both over the
// limit of 1, so both must be answered.
const { raw } = await sendRequests(port, [rawGet("/a"), rawGet("/b"), rawGet("/c")]);

expect([...raw.matchAll(/HTTP\/1\.1 (\d{3})/g)].map(m => m[1])).toEqual(["200", "503", "503"]);
expect(drops.map(({ url, isIncomingMessage }) => ({ url, isIncomingMessage }))).toEqual([
{ url: "/b", isIncomingMessage: true },
{ url: "/c", isIncomingMessage: true },
]);
// Both drops report the one connection they arrived on.
expect(drops[0].socket).toBeDefined();
expect(drops[1].socket).toBe(drops[0].socket);
} finally {
server.close();
}
});

it.concurrent("a dropped request carrying a body does not stall the rest of the pipeline", async () => {
// The 503'd request's body still has to come off the wire, otherwise the
// parser never reaches the request pipelined behind it.
const drops: string[] = [];
const server = createServer((req, res) => {
req.resume();
res.end("ok");
});
server.maxRequestsPerSocket = 1;
server.on("dropRequest", req => drops.push(req.url));
try {
server.listen(0, "127.0.0.1");
await once(server, "listening");
const { port } = server.address() as AddressInfo;

const body = "hello=world";
const requests = [rawPost("/a", body), rawPost("/b", body), rawPost("/c", body)];
const { raw } = await sendRequests(port, requests);

expect([...raw.matchAll(/HTTP\/1\.1 (\d{3})/g)].map(m => m[1])).toEqual(["200", "503", "503"]);
expect(drops).toEqual(["/b", "/c"]);
} finally {
server.close();
}
});

it.concurrent("the connection is ended only after the pipelined requests have been answered", async () => {
// Node leans on keepAliveTimeout to reap the connection; Bun hangs up once
// the current read has been answered - but not before the requests already
// sitting in the read buffer have each gotten their own 503.
const drops: string[] = [];
const server = createServer((req, res) => {
req.resume();
res.end("ok");
});
server.maxRequestsPerSocket = 1;
server.on("dropRequest", req => drops.push(req.url));
try {
server.listen(0, "127.0.0.1");
await once(server, "listening");
const { port } = server.address() as AddressInfo;

// The client never hangs up, so reaching "close" at all means the server did.
const { raw, serverEnded } = await sendRequests(port, [rawGet("/a"), rawGet("/b"), rawGet("/c")], "serverClose");

expect(serverEnded).toBe(true);
expect([...raw.matchAll(/HTTP\/1\.1 (\d{3})/g)].map(m => m[1])).toEqual(["200", "503", "503"]);
expect(drops).toEqual(["/b", "/c"]);
} finally {
server.close();
}
});

it.concurrent("over-limit 503s queued behind an async response are written before the connection ends", async () => {
// The first response finishes two check phases out, so the deferred end runs
// while the over-limit 503s are still queued behind it; closing then (or on
// the still-in-flight response's "finish") would drop them.
const drops: string[] = [];
const server = createServer((req, res) => {
req.resume();
setImmediate(() => setImmediate(() => res.end("ok")));
});
server.maxRequestsPerSocket = 1;
server.on("dropRequest", req => drops.push(req.url));
try {
server.listen(0, "127.0.0.1");
await once(server, "listening");
const { port } = server.address() as AddressInfo;

const { raw } = await sendRequests(port, [rawGet("/a"), rawGet("/b"), rawGet("/c")], "serverClose");

expect([...raw.matchAll(/HTTP\/1\.1 (\d{3})/g)].map(m => m[1])).toEqual(["200", "503", "503"]);
expect(drops).toEqual(["/b", "/c"]);
} finally {
server.close();
}
});

it.concurrent("an over-limit request arriving in a later read is answered before the connection ends", async () => {
// The deferred end must not pin the close to a response that was the
// queue tail at the time: a later read can grow the queue before the
// in-flight response finishes.
const drops: string[] = [];
const { promise: firstResponseBarrier, resolve: unblockFirstResponse } = Promise.withResolvers<void>();
const { promise: sawFirstDrop, resolve: onFirstDrop } = Promise.withResolvers<void>();
const { promise: sawSecondDrop, resolve: onSecondDrop } = Promise.withResolvers<void>();
const server = createServer(async (req, res) => {
req.resume();
await firstResponseBarrier;
res.end("ok");
});
server.maxRequestsPerSocket = 1;
server.on("dropRequest", req => {
drops.push(req.url);
(drops.length === 1 ? onFirstDrop : onSecondDrop)();
});
let socket: ReturnType<typeof connect> | undefined;
try {
server.listen(0, "127.0.0.1");
await once(server, "listening");
const { port } = server.address() as AddressInfo;

const { promise: raw, resolve: gotRaw } = Promise.withResolvers<string>();
const { promise: socketClosed, resolve: onSocketClosed } = Promise.withResolvers<void>();
socket = connect(port, "127.0.0.1");
let data = "";
socket.on("data", chunk => (data += chunk));
socket.on("close", () => {
gotRaw(data);
onSocketClosed();
});
socket.on("error", () => {});
await once(socket, "connect");
// Read #1: /a is handled, /b is over the limit and queued behind it.
socket.write(rawGet("/a") + rawGet("/b"));
await sawFirstDrop;
// Let the server's deferred-end callback (same event loop) run first.
await new Promise<void>(resolve => setImmediate(resolve));
// Read #2: /c is over the limit too, queued behind /b. Racing against the
// socket closing turns a regression that ends the connection too early
// into an assertion failure instead of a hang on the missing drop.
socket.write(rawGet("/c"));
await Promise.race([sawSecondDrop, socketClosed]);
unblockFirstResponse();

const out = await raw;
expect([...out.matchAll(/HTTP\/1\.1 (\d{3})/g)].map(m => m[1])).toEqual(["200", "503", "503"]);
expect(drops).toEqual(["/b", "/c"]);
} finally {
unblockFirstResponse();
socket?.destroy();
server.close();
}
});

it("a non-200 CONNECT through a proxy that holds the connection open is destroyed client-side", async () => {
// cleanupAndPropagate deliberately defers destroy to req.onSocket for
// status-code tunnel failures; oncreate must forward the socket so
Expand Down