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
10 changes: 10 additions & 0 deletions packages/bun-uws/src/HttpContext.h
Original file line number Diff line number Diff line change
Expand Up @@ -657,6 +657,16 @@ struct HttpContext {
httpResponseData->inStream = nullptr;
}
}
if constexpr (IsNodeHttp) {
/* The kind check: upgrade() from the body handler turns the ext into a WebSocketData. */
if (switchToTunnelAfterThisChunk && us_socket_kind((struct us_socket_t *) user) == socketKind()) {
/* The response cannot resume reads in tunnel mode: lift the pause the body left. */
Bun__NodeHTTP__onReadsResumable(SSL, (struct us_socket_t *) user);
if (us_socket_is_closed((struct us_socket_t *) user)) {
return nullptr;
}
}
}
return user;
Comment thread
robobun marked this conversation as resolved.
});

Expand Down
16 changes: 15 additions & 1 deletion src/js/node/_http_server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1319,6 +1319,18 @@ function clearUpgradeIncoming(socket) {
socket[kUpgradeIncoming] = undefined;
}

// Deferred: the 'upgrade' listener runs before the rest of its read, which can complete the body.
function resumePausedUpgradeIncoming(socket) {
const req = socket[kUpgradeIncoming];
if (req === undefined) return;
const response = socket[kHandle]?.response;
if (response && (response.hasBody & NodeHTTPBodyReadState.done) !== 0) {
socket[kUpgradeIncoming] = undefined;
} else {
req.resume();
Comment thread
robobun marked this conversation as resolved.
}
}

// Node.js hands the connection over to 'connect'/'upgrade' listeners with the
// connection-listener set removed (onParserExecuteCommon removes its data/end/
// close/drain/error/timeout listeners) and only net.Socket's own 'end' listener
Expand Down Expand Up @@ -1779,7 +1791,9 @@ function getNodeHTTPServerSocket() {
upgradeIncoming.push(resumed);
}
}
upgradeIncoming.resume();
// Not paused: req.complete is false in the listener, so it can wait for 'end' with no reader.
if (upgradeIncoming.readableFlowing !== false) upgradeIncoming.resume();
else setImmediate(resumePausedUpgradeIncoming, this);
Comment thread
robobun marked this conversation as resolved.
return;
}
if (response) {
Expand Down
193 changes: 191 additions & 2 deletions test/js/node/http/node-http-req-socket-pause.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,10 +2,11 @@
* All tests in this file should also run in Node.js.
*/
import { describe, expect, it } from "bun:test";
import { once } from "node:events";
import { Agent, createServer, request, type Server } from "node:http";
import { once, type EventEmitter } from "node:events";
import { Agent, createServer, request, type IncomingMessage, type Server } from "node:http";
import type { AddressInfo, Socket } from "node:net";
import { connect } from "node:net";
import type { Duplex } from "node:stream";

it("req.socket emits 'pause' once an unread request body fills the IncomingMessage buffer", async () => {
// Node's test-http-no-read-no-dump: a handler that never reads the body sees
Expand Down Expand Up @@ -294,3 +295,191 @@ it("upgrade request whose whole body arrived while it was paused still hands the
if (server.listening) server.close();
}
});

describe("upgrade request whose whole body arrived with its head", () => {
const upgradeHeaders = "Upgrade: test\r\nConnection: Upgrade\r\n";
const switchingProtocols = `HTTP/1.1 101 Switching Protocols\r\n${upgradeHeaders}\r\n`;
const body = Buffer.alloc(100, "y").toString();
const fixedLengthPost = `POST /upgrade HTTP/1.1\r\nHost: a\r\n${upgradeHeaders}Content-Length: ${body.length}\r\n\r\n${body}`;

/** Collects what the upgrade socket receives. */
function tunnelReader() {
let bytes = "";
let wake: (() => void) | undefined;
return {
onData(chunk: Buffer) {
bytes += chunk;
wake?.();
},
get bytes() {
return bytes;
},
async receives(expected: string) {
while (bytes.length < expected.length) {
const { promise, resolve } = Promise.withResolvers<void>();
wake = resolve;
await promise;
}
},
};
}

/** `orFail(p)` settles like `p`, or rejects when a watched socket errors or closes first. */
function failureWatcher() {
const { promise: failed, reject } = Promise.withResolvers<never>();
failed.catch(() => {}); // the teardown closes the sockets after the last race
return {
watch(what: string, socket: EventEmitter) {
socket.on("error", reject);
socket.on("close", () => reject(new Error(`${what} closed`)));
},
orFail: <T>(promise: Promise<T>) => Promise.race([promise, failed]),
};
}

it.each([
{ when: "in", readInListener: true },
{ when: "after", readInListener: false },
])(
"a paused request keeps its body when the upgrade socket is read $when the listener",
async ({ readInListener }) => {
const tunnel = tunnelReader();
const { watch, orFail } = failureWatcher();
const { promise: handedOff, resolve: onUpgrade } = Promise.withResolvers<[IncomingMessage, Duplex]>();
const server = createServer();
server.on("upgrade", (req, socket) => {
req.pause();
watch("the upgrade socket", socket);
if (readInListener) socket.on("data", tunnel.onData);
socket.write(switchingProtocols);
onUpgrade([req, socket]);
});
let client: Awaited<ReturnType<typeof connectTo>> | undefined;
try {
client = await connectTo(server);
watch("the client socket", client.socket);
client.socket.write(fixedLengthPost);
const [req, socket] = await orFail(handedOff);
// The 101 is a round trip: the server has parsed the whole first read.
Comment thread
coderabbitai[bot] marked this conversation as resolved.
await orFail(client.receive("101 Switching Protocols"));
client.socket.write("ping-1;");
if (!readInListener) socket.on("data", tunnel.onData);
// The read of the socket does not resume the request, at once or a turn later.
expect(req.readableFlowing).toBe(false);
await orFail(tunnel.receives("ping-1;"));
expect(req.readableFlowing).toBe(false);

let received = "";
req.on("data", chunk => (received += chunk));
req.resume();
await orFail(once(req, "end"));
expect(received).toBe(body);

client.socket.write("ping-2;");
await orFail(tunnel.receives("ping-1;ping-2;"));
expect(tunnel.bytes).toBe("ping-1;ping-2;");
} finally {
client?.socket.destroy();
server.closeAllConnections();
if (server.listening) server.close();
}
},
);

it("the upgrade socket gets the bytes after a body that paused the connection", async () => {
// As above, the first chunk fills the request's buffer, and the end of the body is
// received while the connection is paused. Node.js v26.3.0 delivers the body, but
// never reads the socket again.
const tunnel = tunnelReader();
const { watch, orFail } = failureWatcher();
const { promise: handedOff, resolve: onUpgrade } = Promise.withResolvers<[IncomingMessage, Duplex]>();
const server = createServer({ highWaterMark: 1024 });
server.on("upgrade", (req, socket) => {
watch("the upgrade socket", socket);
socket.write(switchingProtocols);
onUpgrade([req, socket]);
});
let client: Awaited<ReturnType<typeof connectTo>> | undefined;
try {
client = await connectTo(server);
watch("the client socket", client.socket);
client.socket.write(chunkedPost("/upgrade", upgradeHeaders));
const [req, socket] = await orFail(handedOff);
await orFail(client.receive("101 Switching Protocols"));

let received = "";
req.on("data", chunk => (received += chunk));
await orFail(once(req, "end"));
expect(received).toBe(BODY_HEAD + BODY_TAIL);

client.socket.write("ping;");
socket.on("data", tunnel.onData);
await orFail(tunnel.receives("ping;"));
expect(tunnel.bytes).toBe("ping;");
} finally {
client?.socket.destroy();
server.closeAllConnections();
if (server.listening) server.close();
}
});

it("a read of the upgrade socket still resumes a paused request whose body is incomplete", async () => {
// Like Node.js's UpgradeStream._read: the listener never resumes the request, and the part
// of the body that came with the head fills its buffer. The socket still gets its bytes.
const tunnel = tunnelReader();
const { watch, orFail } = failureWatcher();
const server = createServer({ highWaterMark: 1024 });
server.on("upgrade", (req, socket) => {
req.pause();
watch("the upgrade socket", socket);
socket.on("data", tunnel.onData);
socket.write(switchingProtocols);
});
const request = chunkedPost("/upgrade", upgradeHeaders);
const cut = request.lastIndexOf(`${BODY_TAIL.length.toString(16)}\r\n${BODY_TAIL}`);
let client: Awaited<ReturnType<typeof connectTo>> | undefined;
try {
client = await connectTo(server);
watch("the client socket", client.socket);
client.socket.write(request.slice(0, cut));
await orFail(client.receive("101 Switching Protocols"));
client.socket.write(request.slice(cut) + "ping;");
await orFail(tunnel.receives("ping;"));
expect(tunnel.bytes).toBe("ping;");
} finally {
client?.socket.destroy();
server.closeAllConnections();
if (server.listening) server.close();
}
});

it("a read of the upgrade socket still resumes a request that is not paused", async () => {
// Inside the listener req.complete is still false here (Node.js has true), so a listener
// that waits for the end of the message depends on this resume.
const { watch, orFail } = failureWatcher();
const { promise: completeWhenAccepted, resolve: onAccept } = Promise.withResolvers<boolean>();
const server = createServer();
server.on("upgrade", (req, socket) => {
watch("the upgrade socket", socket);
socket.on("data", () => {});
const accept = () => {
socket.write(switchingProtocols);
onAccept(req.complete);
};
if (req.complete) accept();
else req.once("end", accept);
});
let client: Awaited<ReturnType<typeof connectTo>> | undefined;
try {
client = await connectTo(server);
watch("the client socket", client.socket);
client.socket.write(fixedLengthPost);
expect(await orFail(completeWhenAccepted)).toBe(true);
await orFail(client.receive("101 Switching Protocols"));
} finally {
client?.socket.destroy();
server.closeAllConnections();
if (server.listening) server.close();
}
});
});
Loading