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
19 changes: 18 additions & 1 deletion src/js/node/_http_server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -755,6 +755,9 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort
// Node.js's parserOnIncoming: req.upgrade is true for CONNECT
// regardless of shouldUpgradeCallback.
http_req.upgrade = true;
// llhttp completes a CONNECT request at the end of its headers, so
// Node's 'connect' listener already sees req.complete === true.
http_req.complete = true;
// Node frees the parser before handing the raw socket to 'connect'.
releaseServerParserShim(socket, http_req);
server.emit("connect", http_req, socket, head);
Expand Down Expand Up @@ -970,6 +973,11 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort
if (hasBody) {
socket[kUpgradeIncoming] = http_req;
http_req.once("end", clearUpgradeIncoming.bind(undefined, socket));
} else {
// llhttp completes an Upgrade request without a body at the end
// of its headers, so Node's 'upgrade' listener already sees
// req.complete === true.
http_req.complete = true;
}
const upgradeHead = !hasBody && connectHead ? connectHead : kEmptyBuffer;
let upgradeHandled;
Expand Down Expand Up @@ -2369,8 +2377,17 @@ 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;
const req = this.req;
const socket = req?.socket ?? this.socket;
onResponseFinishHandleSocket(socket?.server, socket, this);
// Like Node's clearIncoming: a request that already ended (for example one
// that optimizeEmptyRequests pre-dumped, which never reaches
// emitEOFIncomingMessageOuter) must not stay reachable from socket.parser
// while the kept-alive connection idles.
if (req != null && req.readableEnded) {
const parser = socket?.parser;
if (parser != null && parser.incoming === req) parser.incoming = null;
}
// The dispatcher detached a synchronously-finished response itself;
// advancing the pipeline again here would skip a queued response.
if (this[kDispatcherDetached]) return;
Expand Down
32 changes: 15 additions & 17 deletions src/runtime/server/NodeHTTPResponse.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,6 @@ use bstr::BStr;

use bun_collections::VecExt;
use bun_core::scoped_log;
use bun_http::Method as HttpMethod;
use bun_jsc::JsCell;
use bun_ptr::AsCtxPtr;
use bun_uws as uws;
Expand Down Expand Up @@ -2596,23 +2595,22 @@ pub(crate) unsafe extern "C" fn NodeHTTPResponse__createForJS(
let request_ref = bun_opaque::opaque_deref(request.cast_const());

let vm = bun_vm_mut(global_object);
let method = HttpMethod::which(request_ref.method()).unwrap_or(HttpMethod::OPTIONS);
// GET in node.js can have a body
if method.has_request_body() || method == HttpMethod::GET {
let req_len: usize = 'brk: {
if let Some(content_length) = request_ref.header(b"content-length") {
scoped_log!(
NodeHTTPResponse,
"content-length: {}",
BStr::new(content_length)
);
break 'brk bun_http_types::parse_content_length(content_length);
}
break 'brk 0;
};
// Like llhttp, the framing headers decide whether a request has a body
// for every method: node delivers the body of a HEAD or TRACE request
// that declares one.
let req_len: usize = 'brk: {
if let Some(content_length) = request_ref.header(b"content-length") {
scoped_log!(
NodeHTTPResponse,
"content-length: {}",
BStr::new(content_length)
);
break 'brk bun_http_types::parse_content_length(content_length);
}
break 'brk 0;
};

*has_body = req_len > 0 || request_ref.has_transfer_encoding();
}
*has_body = req_len > 0 || request_ref.has_transfer_encoding();

let raw_response = if is_ssl != 0 {
uws::AnyResponse::SSL(response_ptr.cast())
Expand Down
91 changes: 91 additions & 0 deletions test/js/node/http/node-http.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4668,3 +4668,94 @@ it("connectionListener pauses reads when queued pipelined responses back up", as
clientSide.destroy();
serverSide.destroy();
});

describe("request completion (req.complete, socket.parser.incoming)", () => {
function rawClient(server: Server, request: string) {
const { port, address } = server.address() as AddressInfo;
const client = connect(port, address, () => client.write(request));
client.on("error", () => {});
client.on("data", () => {});
return client;
}

test("req.complete is true inside a 'connect' listener", async () => {
await using server = createServer();
server.listen(0, "127.0.0.1");
await once(server, "listening");
let sawComplete: boolean | undefined;
server.on("connect", (req, socket) => {
sawComplete = req.complete;
socket.destroy();
});
const client = rawClient(server, "CONNECT example.com:443 HTTP/1.1\r\nHost: example.com:443\r\n\r\n");
const [req] = await once(server, "connect");
expect(sawComplete).toBe(true);
await new Promise(resolve => process.nextTick(resolve));
expect(req.complete).toBe(true);
await once(client, "close");
});

test("req.complete is true inside an 'upgrade' listener for a request without a body", async () => {
await using server = createServer();
server.listen(0, "127.0.0.1");
await once(server, "listening");
let sawComplete: boolean | undefined;
server.on("upgrade", (req, socket) => {
sawComplete = req.complete;
socket.destroy();
});
const client = rawClient(server, "GET / HTTP/1.1\r\nHost: x\r\nConnection: Upgrade\r\nUpgrade: websocket\r\n\r\n");
const [req] = await once(server, "upgrade");
expect(sawComplete).toBe(true);
await new Promise(resolve => process.nextTick(resolve));
expect(req.complete).toBe(true);
await once(client, "close");
});

for (const method of ["HEAD", "TRACE"]) {
test(`a ${method} request that declares a body delivers it and completes only once it arrived`, async () => {
await using server = createServer();
server.listen(0, "127.0.0.1");
await once(server, "listening");
const chunks: Buffer[] = [];
const body = Promise.withResolvers<string>();
server.on("request", (req, res) => {
req.on("data", c => chunks.push(c));
req.on("end", () => {
body.resolve(Buffer.concat(chunks).toString());
res.end();
});
});
// The headers go first. The declared body follows only once the
// request has been dispatched.
const client = rawClient(server, `${method} / HTTP/1.1\r\nHost: x\r\nContent-Length: 5\r\n\r\n`);
const [req] = await once(server, "request");
await new Promise(resolve => process.nextTick(resolve));
expect(req.complete).toBe(false);
client.write("hello");
expect(await body.promise).toBe("hello");
expect(req.complete).toBe(true);
client.destroy();
await once(client, "close");
});
}

test("optimizeEmptyRequests: socket.parser.incoming does not keep the request once the response finished", async () => {
const closed = Promise.withResolvers<{ atRequest: boolean; atClose: boolean }>();
await using server = createServer({ optimizeEmptyRequests: true }, (req, res) => {
const socket = req.socket;
const atRequest = socket.parser.incoming === req;
// The keep-alive connection idles after this response. Node cleared
// parser.incoming in resOnFinish because the pre-dumped request had
// already ended.
res.on("close", () => closed.resolve({ atRequest, atClose: socket.parser.incoming === req }));
res.end("ok");
});
server.listen(0, "127.0.0.1");
await once(server, "listening");
const client = rawClient(server, "GET / HTTP/1.1\r\nHost: x\r\n\r\n");
expect(await closed.promise).toEqual({ atRequest: true, atClose: false });
client.destroy();
await once(client, "close");
});
});
Loading