Skip to content
Closed
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
6 changes: 6 additions & 0 deletions packages/bun-uws/src/HttpContext.h
Original file line number Diff line number Diff line change
Expand Up @@ -844,6 +844,12 @@ struct HttpContext {
&& httpResponseData->onWritable == nullptr) {
responseDone = true;
}
/* socket.destroySoon() issued while bytes were queued: Node's
* destroy() on 'finish' closes once they are out, whether or
* not the response in flight ever ends. */
if (httpResponseData->state & HttpResponseData<SSL>::HTTP_NODE_CLOSE_AFTER_DRAIN) {
responseDone = true;
}
}
if (responseDone && asyncSocket->hasFullyDrained()) {
asyncSocket->shutdown();
Expand Down
9 changes: 7 additions & 2 deletions packages/bun-uws/src/HttpResponseData.h
Original file line number Diff line number Diff line change
Expand Up @@ -150,13 +150,18 @@ struct HttpResponseData : AsyncSocketData<SSL>, HttpParser {
* shutdown sweep; the shouldCloseConnection() gates act on it once the
* in-flight work completes. */
HTTP_CLOSE_WHEN_IDLE = 1 << 17,
/* node:http socket.destroySoon() with outgoing bytes still queued: shut
* down and close as soon as they have flushed, whether or not the
* response in flight has ended (Node's destroy() on 'finish'). */
HTTP_NODE_CLOSE_AFTER_DRAIN = 1 << 18,

/* Bits that describe the connection rather than the response in flight.
* There is one HttpResponseData per socket, reused by every request on a
* keep-alive connection, so starting a new response clears the rest of the
* word (resetResponseState) - these have to survive that. */
HTTP_CONNECTION_SCOPED = HTTP_NODE_PARSING_STOPPED | HTTP_NODE_READS_PAUSED
| HTTP_NODE_TUNNEL_AFTER_BODY | HTTP_NODE_RECEIVED_FIN | HTTP_CLOSE_WHEN_IDLE,
| HTTP_NODE_TUNNEL_AFTER_BODY | HTTP_NODE_RECEIVED_FIN | HTTP_CLOSE_WHEN_IDLE
| HTTP_NODE_CLOSE_AFTER_DRAIN,
};

/* Begin a new response on this connection. Clearing the word in one go is
Expand Down Expand Up @@ -228,7 +233,7 @@ struct HttpResponseData : AsyncSocketData<SSL>, HttpParser {
/* Whether the connection should be torn down once the in-flight response (if
* any) has completed and all buffered outgoing data has been flushed. */
bool shouldCloseConnection() const {
return (state & HTTP_CONNECTION_CLOSE)
return (state & (HTTP_CONNECTION_CLOSE | HTTP_NODE_CLOSE_AFTER_DRAIN))
|| ((state & HTTP_NODE_RECEIVED_FIN) && nodeHttpQueuedPipelinedCount == 0)
|| ((state & HTTP_CLOSE_WHEN_IDLE) && this->isIdle);
}
Expand Down
28 changes: 25 additions & 3 deletions src/js/node/_http_server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1298,9 +1298,10 @@ function onServerClientError(ssl: boolean, socket: unknown, errorCode: number, r
// reaches socketOnError and 'clientError' never fires for it.
function replyMissingHostHeader(socket) {
if (!socket.writable) return;
socket.end(
socket.write(
`HTTP/1.1 400 Bad Request\r\nConnection: close\r\nDate: ${new Date().toUTCString()}\r\nTransfer-Encoding: chunked\r\n\r\n0\r\n\r\n`,
);
socket.destroySoon();
}

const kBytesWritten = Symbol("kBytesWritten");
Expand Down Expand Up @@ -1333,6 +1334,8 @@ function resolveHandoffPromise(promise) {
}
const kSocketTimeoutTimer = Symbol("socketTimeoutTimer");
const kStreamingEnabled = Symbol("kStreamingEnabled");
// destroySoon() was called: _final closes the connection right behind the FIN.
const kDestroySoon = Symbol("kDestroySoon");
// Scratch options object for the builtin ServerResponse (see the dispatcher).
const scratchResponseOptions = {
[kHandle]: undefined,
Expand Down Expand Up @@ -1485,6 +1488,7 @@ function getNodeHTTPServerSocket() {
[kBytesWritten] = 0;
[kHandle];
[kUpgradeIncoming] = undefined;
[kDestroySoon] = false;
server: Server;
_httpMessage;
_secureEstablished = false;
Expand Down Expand Up @@ -1765,10 +1769,22 @@ function getNodeHTTPServerSocket() {
callback();
return;
}
handle.end();
handle.end(this[kDestroySoon]);
callback();
}

// Not destroy() on 'finish' like net.Socket: 'finish' does not wait for bytes uWS still queues.
destroySoon() {
if (this[kDestroySoon]) return;
const handle = this[kHandle];
if (this.writable && handle && !handle.closed) {
this[kDestroySoon] = true;
this.end();
return;
}
super.destroySoon();
}
Comment thread
robobun marked this conversation as resolved.

get localAddress() {
return this[kHandle]?.localAddress?.address;
}
Expand Down Expand Up @@ -2405,7 +2421,13 @@ function emitResponseFinish() {
// is eventually closed.
function onResponseFinishHandleSocket(server, socket, res) {
if (res[kMustCloseConnection]) {
socket?.end();
if (socket != null) {
if (typeof socket.destroySoon === "function") {
socket.destroySoon();
} else {
socket.end();
}
}
return;
}
if (!socket || socket.destroyed || typeof socket.setTimeout !== "function") {
Expand Down
12 changes: 8 additions & 4 deletions src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -243,7 +243,7 @@ bool JSNodeHTTPServerSocket::isClosed() const
}

template<bool SSL>
static bool deferShutdownUntilResponseDrains(us_socket_t* socket)
static bool deferShutdownUntilResponseDrains(us_socket_t* socket, bool thenClose)
{
if (reinterpret_cast<uWS::AsyncSocket<SSL>*>(socket)->getBufferedAmount() == 0) {
return false;
Expand All @@ -253,18 +253,22 @@ static bool deferShutdownUntilResponseDrains(us_socket_t* socket)
* is sequenced after the response bytes (like Node's destroySoon). */
auto* httpResponseData = reinterpret_cast<uWS::HttpResponseData<SSL>*>(us_socket_ext(socket));
httpResponseData->state |= uWS::HttpResponseData<SSL>::HTTP_CONNECTION_CLOSE;
if (thenClose) {
/* destroySoon(): and closes it there, even if the response never ends. */
httpResponseData->state |= uWS::HttpResponseData<SSL>::HTTP_NODE_CLOSE_AFTER_DRAIN;
}
return true;
}

bool JSNodeHTTPServerSocket::shutdownAfterResponseDrains()
bool JSNodeHTTPServerSocket::shutdownAfterResponseDrains(bool thenClose)
{
if (!socket || upgraded || us_socket_is_closed(socket) || us_socket_is_shut_down(socket)) {
return false;
}
if (is_ssl) {
return deferShutdownUntilResponseDrains<true>(socket);
return deferShutdownUntilResponseDrains<true>(socket, thenClose);
}
return deferShutdownUntilResponseDrains<false>(socket);
return deferShutdownUntilResponseDrains<false>(socket, thenClose);
}

template<bool SSL>
Expand Down
2 changes: 1 addition & 1 deletion src/jsc/bindings/node/JSNodeHTTPServerSocket.h
Original file line number Diff line number Diff line change
Expand Up @@ -106,7 +106,7 @@ class JSNodeHTTPServerSocket : public JSC::JSDestructibleObject {
/* node:http socket.end(): when the in-flight response still has bytes in
* uWS's send buffer, a shutdown now would put the FIN ahead of them and
* truncate the response. Returns true after handing the close to uWS. */
bool shutdownAfterResponseDrains();
bool shutdownAfterResponseDrains(bool thenClose);

/* Switch the connection into CONNECT-style tunnel mode after an accepted
* Upgrade: subsequent bytes bypass the HTTP parser and stream to the
Expand Down
15 changes: 14 additions & 1 deletion src/jsc/bindings/node/JSNodeHTTPServerSocketPrototype.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ extern "C" uint64_t uws_res_get_remote_address_info(void* res, const char** dest
extern "C" uint64_t uws_res_get_local_address_info(void* res, const char** dest, int* port, bool* is_ipv6);
extern "C" void us_socket_resume(us_socket_t*);
extern "C" void us_socket_pause(us_socket_t*);
extern "C" void us_socket_shutdown(us_socket_t*);

namespace Bun {

Expand Down Expand Up @@ -216,12 +217,24 @@ JSC_DEFINE_HOST_FUNCTION(jsFunctionNodeHTTPServerSocketEnd, (JSC::JSGlobalObject
}

thisObject->ended = true;
// end(true) is Node's destroySoon(): close right behind the FIN instead of waiting for the peer's.
bool destroySoon = callFrame->argument(0).toBoolean(globalObject);
// The response's buffered body must reach the kernel before the FIN; uWS
// performs the shutdown after its send buffer drains.
if (thisObject->shutdownAfterResponseDrains()) {
if (thisObject->shutdownAfterResponseDrains(destroySoon)) {
return JSValue::encode(JSC::jsUndefined());
}
auto bufferedSize = thisObject->streamBuffer.bufferedSize();
if (destroySoon) {
// One flush for raw socket.write() bytes; the close drops the rest, as destroy() did.
if (bufferedSize == 0) {
us_socket_shutdown(thisObject->socket);
} else {
us_socket_buffered_js_write(thisObject->socket, thisObject->is_ssl, thisObject->ended, &thisObject->streamBuffer, globalObject, JSValue::encode(JSC::jsUndefined()), JSValue::encode(JSC::jsUndefined()));
}
thisObject->close();
return JSValue::encode(JSC::jsUndefined());
}
if (bufferedSize == 0) {
// onNodeHTTPRequest no longer pauses at dispatch; pause here so the
// shutdown+resume below still cycles kqueue's EVFILT_READ (delete then
Expand Down
154 changes: 131 additions & 23 deletions test/js/node/http/node-http.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2708,34 +2708,142 @@ it("removing only Content-Length falls back to chunked encoding and keeps the co
}
});

it("an explicit Connection: close response header closes the server-side socket after finish", async () => {
// Node.js's matchHeader sets _last for a user-set Connection: close, and
// resOnFinish then ends the socket; the transport must match the header.
const server = createServer((req, res) => {
res.setHeader("Connection", "close");
res.end("ok");
});
try {
// Node.js's resOnFinish destroySoon()s the socket after a response that ends the
// connection: the server sends its FIN and releases the socket right behind it
// ('close' on the server-side socket, fd gone) without waiting for the peer to
// close its side. A half-close that waits for the client's FIN lets a client
// that never closes pin one socket per completed request.
describe("a response that closes the connection releases the server-side socket without waiting for the peer", () => {
const cases = [
{
name: "Connection: close response header",
request: "GET / HTTP/1.1\r\nHost: localhost\r\n\r\n",
prepare: (res: ServerResponse) => res.setHeader("Connection", "close"),
expected: /^HTTP\/1\.1 200 OK\r\n[^]*ok$/,
},
{
name: "Connection: close request header",
request: "GET / HTTP/1.1\r\nHost: localhost\r\nConnection: close\r\n\r\n",
expected: /^HTTP\/1\.1 200 OK\r\n[^]*ok$/,
},
{
name: "HTTP/1.0 request",
request: "GET / HTTP/1.0\r\nHost: localhost\r\n\r\n",
expected: /^HTTP\/1\.1 200 OK\r\n[^]*ok$/,
},
{
// Node answers this itself with res.writeHead(400, ['Connection', 'close']).
name: "HTTP/1.1 request without a Host header",
request: "GET / HTTP/1.1\r\n\r\n",
expected: /^HTTP\/1\.1 400 Bad Request\r\n[^]*\r\n0\r\n\r\n$/,
},
];
for (const { name, request, prepare, expected } of cases) {
it.concurrent(name, async () => {
const { promise: serverSocketClosed, resolve: onServerSocketClose } = Promise.withResolvers<void>();
const server = createServer((req, res) => {
prepare?.(res);
res.end("ok");
});
server.on("connection", socket => socket.on("close", onServerSocketClose));
server.listen(0, "127.0.0.1");
await once(server, "listening");
const { port } = server.address() as AddressInfo;
// allowHalfOpen: the client keeps its side open after the server's FIN, so
// only the server can release the server-side socket.
const client = connect({ port, host: "127.0.0.1", allowHalfOpen: true });
try {
const out = await new Promise<string>((resolve, reject) => {
let data = "";
client.setEncoding("latin1");
client.on("data", chunk => (data += chunk));
client.on("end", () => resolve(data));
client.on("error", reject);
client.write(request);
});
expect(out).toMatch(expected);
expect(out).toContain("\r\nConnection: close\r\n");
// The client has the whole response and the server's FIN, and has not
// sent its own. The server-side socket must close on its own now.
await serverSocketClosed;
expect(client.writableEnded).toBe(false);
} finally {
client.destroy();
server.close();
}
});
}

// net.Socket.destroySoon() is what Node's resOnFinish uses: end(), then
// destroy() once the write side is done, without waiting for the peer.
it.concurrent("socket.destroySoon() from a handler", async () => {
const { promise: serverSocketClosed, resolve: onServerSocketClose } = Promise.withResolvers<void>();
const server = createServer((req, res) => {
req.socket.write("raw bytes\r\n");
req.socket.destroySoon();
});
server.on("connection", socket => socket.on("close", onServerSocketClose));
server.listen(0, "127.0.0.1");
await once(server, "listening");
const { port } = server.address() as AddressInfo;
const client = connect({ port, host: "127.0.0.1", allowHalfOpen: true });
try {
const out = await new Promise<string>((resolve, reject) => {
let data = "";
client.setEncoding("latin1");
client.on("data", chunk => (data += chunk));
client.on("end", () => resolve(data));
client.on("error", reject);
client.write("GET / HTTP/1.1\r\nHost: localhost\r\n\r\n");
});
expect(out).toBe("raw bytes\r\n");
await serverSocketClosed;
expect(client.writableEnded).toBe(false);
} finally {
client.destroy();
server.close();
}
});

const out = await new Promise<string>((resolve, reject) => {
const socket = connect(port, "127.0.0.1");
let data = "";
socket.on("data", chunk => (data += chunk));
// The server must send FIN on its own; the client never half-closes.
socket.on("end", () => resolve(data));
socket.on("error", reject);
socket.write("GET / HTTP/1.1\r\nHost: localhost\r\n\r\n");
// destroySoon() with response bytes still queued in the transport (the client
// is not reading yet) and no res.end(): the queued bytes are delivered first,
// then the connection closes, like Node's destroy() on the socket's 'finish'.
it.concurrent("socket.destroySoon() behind a backed-up response that never ends", async () => {
const CHUNK = Buffer.alloc(16 * 1024, "a");
const CHUNKS = 1024; // 16 MB: more than the loopback socket buffers absorb
const { promise: serverSocketClosed, resolve: onServerSocketClose } = Promise.withResolvers<void>();
const { promise: destroyed, resolve: onDestroySoonCalled } = Promise.withResolvers<void>();
const server = createServer((req, res) => {
res.on("error", () => {});
for (let i = 0; i < CHUNKS; i++) res.write(CHUNK);
req.socket.destroySoon();
onDestroySoonCalled();
});

expect(out).toContain("HTTP/1.1 200");
expect(out).toContain("Connection: close");
expect(out).toEndWith("ok");
} finally {
server.close();
}
server.on("connection", socket => socket.on("close", onServerSocketClose));
server.listen(0, "127.0.0.1");
await once(server, "listening");
const { port } = server.address() as AddressInfo;
const client = connect({ port, host: "127.0.0.1", allowHalfOpen: true });
try {
let received = 0;
const ended = new Promise<void>((resolve, reject) => {
client.on("data", chunk => (received += chunk.length));
client.on("end", resolve);
client.on("error", reject);
});
client.pause();
client.write("GET / HTTP/1.1\r\nHost: localhost\r\n\r\n");
await destroyed;
client.resume();
await ended;
expect(received).toBeGreaterThan(CHUNK.length * CHUNKS);
await serverSocketClosed;
expect(client.writableEnded).toBe(false);
} finally {
client.destroy();
server.close();
}
});
});

it("a pipelined request behind Connection: close is never dispatched (clientError HPE_CLOSED_CONNECTION)", async () => {
Expand Down
Loading