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
23 changes: 15 additions & 8 deletions src/js/node/http2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2021,6 +2021,12 @@ function pushToStream(stream, data) {
}
}

// Like node's onSessionHeaders, a HEADERS frame with END_STREAM ends the readable before its event.
function endInboundHalf(stream: Http2Stream) {
if (!stream.rstCode) stream.rstCode = 0;
pushToStream(stream, null);
}

enum StreamState {
EndedCalled = 1 << 0, // 00001 = 1
WantTrailer = 1 << 1, // 00010 = 2
Expand Down Expand Up @@ -3266,6 +3272,7 @@ class ServerHttp2Stream extends Http2Stream {
if (pushedStream && pushedStream[bunHTTP2Headers] == null) {
pushedStream[bunHTTP2Headers] = headers;
}
if (pushedStream) endInboundHalf(pushedStream);
if (onServerStreamCreatedChannel.hasSubscribers) {
onServerStreamCreatedChannel.publish({ stream: pushedStream, headers });
}
Expand Down Expand Up @@ -4057,10 +4064,7 @@ class ServerHttp2Session extends Http2Session {
}
if (state == 6 || state == 7) {
if (stream.readable) {
if (!stream.rstCode) {
stream.rstCode = 0;
}
pushToStream(stream, null);
endInboundHalf(stream);

// If the user hasn't tried to consume the stream then dump the incoming data so the
// stream can finish — but at half-close only when nothing is buffered: a consumer may
Expand Down Expand Up @@ -4149,7 +4153,9 @@ class ServerHttp2Session extends Http2Session {
stream[kHeadRequest] = true;
}
const status = stream[bunHTTP2StreamStatus];
const endOfStream = (flags & constants.NGHTTP2_FLAG_END_STREAM) !== 0;
if ((status & StreamState.StreamResponded) !== 0) {
if (endOfStream) endInboundHalf(stream);
stream.emit("trailers", headers, flags, rawheaders);
} else {
// Set the StreamResponded bit BEFORE dispatching the 'stream' event
Expand All @@ -4166,6 +4172,7 @@ class ServerHttp2Session extends Http2Session {
if (onServerStreamStartChannel.hasSubscribers) {
onServerStreamStartChannel.publish({ stream, headers });
}
if (endOfStream) endInboundHalf(stream);
// performServerHandshake() sessions have no owning server.
self[kServer]?.emit("stream", stream, headers, flags, rawheaders);
self.emit("stream", stream, headers, flags, rawheaders);
Expand Down Expand Up @@ -5056,12 +5063,9 @@ class ClientHttp2Session extends Http2Session {
}
if (state == 6 || state == 7) {
if (stream.readable) {
if (!stream.rstCode) {
stream.rstCode = 0;
}
// Push a null so the stream can end whenever the client consumes
// it completely.
pushToStream(stream, null);
endInboundHalf(stream);
stream.read(0);
}
}
Expand Down Expand Up @@ -5129,11 +5133,13 @@ class ClientHttp2Session extends Http2Session {
}
const status = stream[bunHTTP2StreamStatus];
const header_status = headers[HTTP2_HEADER_STATUS];
const endOfStream = (flags & constants.NGHTTP2_FLAG_END_STREAM) !== 0;
if (header_status === HTTP_STATUS_CONTINUE) {
stream.emit("continue");
}

if ((status & StreamState.StreamResponded) !== 0) {
if (endOfStream) endInboundHalf(stream);
stream.emit("trailers", headers, flags, rawheaders);
} else {
if (header_status >= 100 && header_status < 200) {
Expand All @@ -5151,6 +5157,7 @@ class ClientHttp2Session extends Http2Session {
if (onClientStreamFinishChannel.hasSubscribers) {
onClientStreamFinishChannel.publish({ stream, headers, flags });
}
if (endOfStream) endInboundHalf(stream);
if (stream[kPush]) {
// A pushed stream delivers its response via 'push'; the session 'stream' event already
// fired (with the promised request headers) when the PUSH_PROMISE arrived.
Expand Down
149 changes: 147 additions & 2 deletions test/js/node/http2/node-http2-client-close.test.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,9 @@
/**
* ClientHttp2Stream.close(code) event contract:
* Http2Stream.close(code) event contract, client and server, while the peer's half is still open:
* NO_ERROR / CANCEL -> 'end', 'close' (documented 'error' exemption)
* any other code -> 'error', 'close' (no 'end': the body was killed by RST_STREAM)
* Once the peer's half has ended (END_STREAM on a HEADERS frame, or a server push, which has no
* inbound half) the readable already has its EOF, so 'end' comes first for every code.
*
* Works with both:
* bun bd test test/js/node/http2/node-http2-client-close.test.ts
Expand Down Expand Up @@ -183,10 +185,153 @@ describe("ClientHttp2Stream.close(code) while pending", () => {
}
});

// The listeners below call close(code) once the peer's half of the stream has ended. Each one is
// the event of a HEADERS frame that carried END_STREAM (or the pushStream() callback: a server push
// has no inbound half). "server 'stream', request open" is the contrast: no END_STREAM yet.
type Site =
| "server 'stream'"
| "server 'stream' after respond()"
| "server 'stream', request open"
| "server 'trailers'"
| "pushStream() callback"
| "client 'response'"
| "client 'trailers'";

async function closeIn(site: Site, code: number, defer: boolean): Promise<string[]> {
const events: string[] = [];
const { promise, resolve, reject } = Promise.withResolvers<string[]>();
const watch = (stream: http2.Http2Stream) => {
stream.on("end", () => events.push("end"));
stream.on("error", e => events.push("error:" + (e as NodeJS.ErrnoException).code));
stream.on("close", () => {
events.push("close:" + stream.rstCode);
resolve(events);
});
stream.resume();
};
const close = (stream: http2.Http2Stream) => {
if (defer) process.nextTick(() => stream.close(code));
else stream.close(code);
};
const ignoreErrors = (stream: http2.Http2Stream) => stream.on("error", () => {});

const server = http2.createServer();
server.on("stream", stream => {
switch (site) {
case "server 'stream'":
case "server 'stream', request open":
watch(stream);
close(stream);
break;
case "server 'stream' after respond()":
watch(stream);
stream.respond({ ":status": 200 }, { endStream: true });
close(stream);
break;
case "server 'trailers'":
watch(stream);
stream.respond({ ":status": 200 });
stream.on("trailers", () => close(stream));
break;
case "pushStream() callback":
ignoreErrors(stream);
stream.pushStream({ ":path": "/pushed" }, (err, pushed) => {
if (err) return reject(err);
watch(pushed);
close(pushed);
});
stream.respond({ ":status": 200 });
stream.end();
break;
case "client 'response'":
ignoreErrors(stream);
stream.respond({ ":status": 204 }, { endStream: true });
break;
case "client 'trailers'":
ignoreErrors(stream);
stream.respond({ ":status": 200 }, { waitForTrailers: true });
stream.on("wantTrailers", () => stream.sendTrailers({ "x-trailer": "1" }));
stream.end("body");
break;
}
});
const port = await listen(server);
const client = http2.connect(`http://127.0.0.1:${port}`);
client.on("error", reject);
client.on("stream", pushed => ignoreErrors(pushed).resume());
try {
switch (site) {
case "server 'trailers'": {
const req = client.request({ ":path": "/", ":method": "POST" }, { waitForTrailers: true });
req.on("wantTrailers", () => req.sendTrailers({ "x-trailer": "1" }));
ignoreErrors(req).resume();
req.end("body");
break;
}
case "client 'response'":
case "client 'trailers'": {
// The request stays open, so nothing but close(code) can close the stream.
const req = client.request({ ":path": "/", ":method": "POST" });
watch(req);
req.on(site === "client 'response'" ? "response" : "trailers", () => close(req));
break;
}
case "server 'stream', request open":
ignoreErrors(client.request({ ":path": "/", ":method": "POST" })).resume();
break;
default:
ignoreErrors(client.request({ ":path": "/" }, { endStream: true })).resume();
}
return await promise;
} finally {
client.destroy();
server.close();
}
}

const endThenClose = [
[NGHTTP2_NO_ERROR, ["end", "close:0"]],
[NGHTTP2_CANCEL, ["end", "close:8"]],
[NGHTTP2_INTERNAL_ERROR, ["end", "error:ERR_HTTP2_STREAM_ERROR", "close:2"]],
[NGHTTP2_ENHANCE_YOUR_CALM, ["end", "error:ERR_HTTP2_STREAM_ERROR", "close:11"]],
] as const;
const noEndForErrorCodes = [
[NGHTTP2_NO_ERROR, ["end", "close:0"]],
[NGHTTP2_CANCEL, ["end", "close:8"]],
[NGHTTP2_INTERNAL_ERROR, ["error:ERR_HTTP2_STREAM_ERROR", "close:2"]],
[NGHTTP2_ENHANCE_YOUR_CALM, ["error:ERR_HTTP2_STREAM_ERROR", "close:11"]],
] as const;

for (const [site, table] of [
["server 'stream'", endThenClose],
["server 'stream' after respond()", endThenClose],
["server 'trailers'", endThenClose],
["pushStream() callback", endThenClose],
["client 'response'", endThenClose],
["client 'trailers'", endThenClose],
["server 'stream', request open", noEndForErrorCodes],
] as const) {
for (const defer of [false, true]) {
describe(`close(code) ${defer ? "a tick after" : "in"} ${site}`, () => {
for (const [code, expected] of table) {
test(`close(${code})`, async () => {
Comment thread
coderabbitai[bot] marked this conversation as resolved.
assert.deepStrictEqual(await closeIn(site, code, defer), expected);
});
}
});
}
}

if (typeof Bun !== "undefined") {
const node = Bun.which("node");
// Alpine's node segfaults at a random point of this file (alpine 3.23 aarch64, on main too), and
// the CI runner fails a file for any new core dump, whichever process wrote it. The glibc, macOS
// and Windows lanes keep the cross-check.
const isMusl =
process.platform === "linux" &&
!(process.report.getReport() as { header: { glibcVersionRuntime?: string } }).header.glibcVersionRuntime;
Comment thread
robobun marked this conversation as resolved.
describe("Node.js compatibility", () => {
test("tests should run on node.js", { skip: !node }, async () => {
test("tests should run on node.js", { skip: !node || isMusl }, async () => {
await using proc = Bun.spawn({
cmd: [node as string, "--test", import.meta.filename],
stdout: "inherit",
Expand Down