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
85 changes: 28 additions & 57 deletions src/js/node/http2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -444,11 +444,22 @@ function emitEventNT(self: any, event: string, ...args: any[]) {
// The frame is passed in: destroy() clears it off the session before emitting
// 'error', so a throwing listener on either event cannot leave a retained
// session pinning the store.
function emitSessionCloseNT(self: Http2Session, frame) {
function emitSessionCloseNT(self: Http2Session, error: Error | null | undefined, frame) {
if (error) {
runInFrame(frame, self.emit, self, "error", error);
}
if (self.listenerCount("close") > 0) {
runInFrame(frame, self.emit, self, "close");
}
}
// Node's emitClose: the session reports 'error' and 'close' once its socket has closed.
function emitSessionCloseAfterSocket(self: Http2Session, socket, error: Error | null | undefined, frame) {
if (socket && !socket.destroyed) {
socket.once("close", () => emitSessionCloseNT(self, error, frame));
} else {
process.nextTick(emitSessionCloseNT, self, error, frame);
}
}
function emitErrorNT(self: any, error: any, destroy: boolean) {
if (destroy) {
if (self.listenerCount("error") > 0) {
Expand Down Expand Up @@ -4750,7 +4761,8 @@ class ServerHttp2Session extends Http2Session {
return;
}
this.#destroying = true;
emitHttp2SessionPerf(this, this.#parser, this[bunHTTP2Socket]);
const socket = this[bunHTTP2Socket];
emitHttp2SessionPerf(this, this.#parser, socket);
try {
const server = this[kServer];
if (server) {
Expand All @@ -4774,7 +4786,6 @@ class ServerHttp2Session extends Http2Session {
this[kSessionDestroyError] = error;
}

const socket = this[bunHTTP2Socket];
if (!this.#connected) return;
Comment thread
robobun marked this conversation as resolved.
this.#closed = true;
this.#connected = false;
Expand All @@ -4785,24 +4796,7 @@ class ServerHttp2Session extends Http2Session {
// a destroy(err) after close() must still put the error GOAWAY on the wire.
this.goaway(code || constants.NGHTTP2_NO_ERROR, 0, Buffer.alloc(0));
}
if (error) {
// node's finishSessionClose destroys the socket when the session dies
// with an error (a misbehaving peer must observe the connection going
// away) - but it still ends first and destroys a tick later, so the
// final GOAWAY flushes behind a FIN instead of an abortive close (see
// endThenDestroySessionSocket).
endThenDestroySessionSocket(socket, error);
} else {
// Node's finishSessionClose: "If we're gracefully closing the socket,
// call resume() so we can detect the peer closing in case
// binding.Http2Session is already gone." Without a reader, unread
// inbound bytes (a late GOAWAY from the peer) turn the close into an
// RST, which the peer surfaces as read ECONNRESET (routine on Windows
// loopback - the same reason Node delays the error-path destroy).
// https://github.com/nodejs/node/blob/v26.3.0/lib/internal/http2/core.js#L1188
socket.resume();
socket.end();
}
closeSessionSocket(socket, this.#closeCalled && !error, error);
}
const parser = this.#parser;
if (parser) {
Expand All @@ -4828,16 +4822,9 @@ class ServerHttp2Session extends Http2Session {
}
this[bunHTTP2Socket] = null;

// Read-and-clear the frame first: emitting 'error' with no listener throws,
// which would skip the clear and leave a retained session pinning the store.
const asyncFrame = this[bunHTTP2AsyncContextFrame];
this[bunHTTP2AsyncContextFrame] = undefined;
if (error) {
runInFrame(asyncFrame, this.emit, this, "error", error);
}
// node emits the session 'close' event asynchronously (a listener attached right after
// close()/destroy() returns must still observe it).
process.nextTick(emitSessionCloseNT, this, asyncFrame);
emitSessionCloseAfterSocket(this, socket, error, asyncFrame);
}
}
function emitTimeout(session: ClientHttp2Session) {
Expand Down Expand Up @@ -4880,19 +4867,21 @@ function setSessionTimeout(this: Http2Session, msecs, callback) {
return this;
}

// Node's finishSessionClose error path: socket.end() flushes and sends the FIN
// first, and the hard destroy runs a tick later - "If session.destroy() was
// called, destroy the underlying socket. Delay it a bit to try to avoid
// ECONNRESET on Windows" - so the peer reads our final frames off a FIN'd
// socket instead of observing an abortive close.
// https://github.com/nodejs/node/blob/v26.3.0/lib/internal/http2/core.js#L1188
function destroySessionSocketDelayedNT(socket, error) {
if (!socket.destroyed) {
socket.destroy(error);
}
}
function endThenDestroySessionSocket(socket, error) {
socket.end(() => setImmediate(destroySessionSocketDelayedNT, socket, error));
// Node's finishSessionClose: end(), then destroy() unless close() asked for a graceful shutdown.
function closeSessionSocket(socket, graceful: boolean, error: Error | null | undefined) {
if (socket.destroyed) return;
if (graceful) {
// Unread inbound bytes would turn the close into an RST.
socket.resume();
socket.end();
} else {
Comment thread
robobun marked this conversation as resolved.
socket.end(() => setImmediate(destroySessionSocketDelayedNT, socket, error));
}
}
// node callTimeout (lib/internal/http2/core.js): when the timer expires while writes are still in
// flight and bytes have reached the wire since the previous expiry, the session is not idle —
Expand Down Expand Up @@ -5896,18 +5885,7 @@ class ClientHttp2Session extends Http2Session {
// a destroy(err) after close() must still put the error GOAWAY on the wire.
this.goaway(code || constants.NGHTTP2_NO_ERROR, 0, Buffer.alloc(0));
}
if (error) {
// See the client session: end first, destroy a tick later (node's
// finishSessionClose Windows-ECONNRESET avoidance).
endThenDestroySessionSocket(socket, error);
} else {
// See the client session's destroy: Node's finishSessionClose resumes
// the socket on a graceful close so unread inbound bytes cannot turn
// the FIN teardown into an RST.
// https://github.com/nodejs/node/blob/v26.3.0/lib/internal/http2/core.js#L1188
socket.resume();
socket.end();
}
closeSessionSocket(socket, this.#closeCalled && !error, error);
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}
const parser = this.#parser;
if (parser) {
Expand Down Expand Up @@ -5937,16 +5915,9 @@ class ClientHttp2Session extends Http2Session {
this.#parser = null;
this[bunHTTP2Socket] = null;

// Read-and-clear the frame first: emitting 'error' with no listener throws,
// which would skip the clear and leave a retained session pinning the store.
const asyncFrame = this[bunHTTP2AsyncContextFrame];
this[bunHTTP2AsyncContextFrame] = undefined;
if (error) {
runInFrame(asyncFrame, this.emit, this, "error", error);
}
// node emits the session 'close' event asynchronously (a listener attached right after
// close()/destroy() returns must still observe it).
process.nextTick(emitSessionCloseNT, this, asyncFrame);
emitSessionCloseAfterSocket(this, socket, error, asyncFrame);
}

request(headers?: HeadersObject | any[] | null, options?: ClientRequestOptions) {
Expand Down
28 changes: 19 additions & 9 deletions test/js/node/async_hooks/AsyncLocalStorage.test.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
import { AsyncLocalStorage, AsyncResource } from "async_hooks";
import { heapStats } from "bun:jsc";
import { describe, expect, test } from "bun:test";
import { bunEnv, bunExe } from "harness";
import { bunEnv, bunExe, isASAN, isDebug } from "harness";
import http2 from "http2";

describe("AsyncLocalStorage", () => {
Expand Down Expand Up @@ -280,7 +280,10 @@ test("re-entering a storage inside run() does not grow the context", () => {
};
als.run(0, () => {
const before = objects();
for (let i = 0; i < 100_000; i++) {
// A context that grows per re-entry adds at least one object per iteration, so 10k
// iterations still blow past the threshold. The full count is too slow under ASAN.
const iterations = isASAN || isDebug ? 10_000 : 100_000;
for (let i = 0; i < iterations; i++) {
als.run(1, () => {
using _ = als.withScope(2);
});
Expand Down Expand Up @@ -1180,9 +1183,11 @@ describe("async context passes through", () => {
expect(stderr).not.toContain("AssertionError");
});

// destroy(err) with no 'error' listener throws out of the emit, which must
// not skip the frame clear (the 'close' tick after it never runs).
test("http2 clears the session frame when destroy(err) throws past the emit", async () => {
// Like node, destroy(err) emits 'error' later (once the socket has closed), so
// with no listener it surfaces as an uncaught exception rather than a throw out
// of destroy(); the frame must already be clear when destroy() returns, not
// only once the deferred emit has run.
test("http2 clears the session frame when destroy(err) has no 'error' listener", async () => {
await using proc = Bun.spawn({
cmd: [
bunExe(),
Expand All @@ -1196,13 +1201,18 @@ describe("async context passes through", () => {
als.run({ marker: true }, () => {
client = http2.connect("http://127.0.0.1:" + server.address().port);
});
// The throw IS the condition: the unlistened 'error' can only come from
// the deferred emit, which runs strictly after the read-and-clear.
process.on("uncaughtException", err => {
console.log("UNCAUGHT " + err.message);
// Drain rather than process.exit(), like the siblings.
server.close();
});
client.on("connect", () => {
try { client.destroy(new Error("boom")); } catch {}
client.destroy(new Error("boom"));
const sym = Object.getOwnPropertySymbols(client)
.find(x => x.description === "::bunhttp2asynccontextframe::");
console.log(sym === undefined ? "SYMBOL-MISSING" : client[sym] === undefined ? "CLEARED" : "PINNED");
// Drain rather than process.exit(), like the siblings.
server.close();
});
});`,
],
Expand All @@ -1211,7 +1221,7 @@ describe("async context passes through", () => {
stderr: "pipe",
});
const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]);
expect({ stdout: stdout.trim(), exitCode }).toEqual({ stdout: "CLEARED", exitCode: 0 });
expect({ stdout: stdout.trim(), exitCode }).toEqual({ stdout: "CLEARED\nUNCAUGHT boom", exitCode: 0 });
expect(stderr).not.toContain("AssertionError");
});

Expand Down
22 changes: 15 additions & 7 deletions test/js/node/http2/h2-conformance.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1943,10 +1943,18 @@ describe("stream-reset floods (CVE-2023-44487 rapid reset, CVE-2025-8671 MadeYou
const calm = (f: Frame) => f.type === FrameType.GOAWAY && goawayErrorCode(f) === ErrorCode.ENHANCE_YOUR_CALM;

function respondingServer(options: Record<string, unknown> = {}, rejectUploads = false) {
const state = { handlers: 0, sessionErrorCode: undefined as string | undefined };
// A promise, not a value sampled when the GOAWAY arrives: like node, the session reports its
// error once its socket has closed, which is after the GOAWAY has reached the peer.
const sessionError = Promise.withResolvers<string>();
sessionError.promise.catch(() => {}); // the tests that expect no session error never read it
const state = { handlers: 0, sessionErrorCode: sessionError.promise };
const server = http2.createServer(options);
server.on("sessionError", (e: any) => (state.sessionErrorCode = e.code));
server.on("session", s => s.on("error", () => {}));
server.on("sessionError", (e: any) => sessionError.resolve(e.code));
Comment thread
robobun marked this conversation as resolved.
server.on("session", s => {
s.on("error", () => {});
// 'error' comes before 'close', so this only settles the promise when there was no error.
s.on("close", () => sessionError.reject(new Error("the session closed without a 'sessionError'")));
});
server.on("stream", (stream: any, headers: any) => {
state.handlers++;
stream.on("error", () => {});
Expand Down Expand Up @@ -1992,7 +2000,7 @@ describe("stream-reset floods (CVE-2023-44487 rapid reset, CVE-2025-8671 MadeYou
expect(goaway.payload.subarray(8).toString()).toBe("too many stream resets");
// Like node, every request up to the bucket's edge still reaches the handler.
expect(handlers).toBeGreaterThan(900);
expect(sessionErrorCode).toBe("ERR_HTTP2_ERROR");
expect(await sessionErrorCode).toBe("ERR_HTTP2_ERROR");
});

test("streamResetBurst sets where the flood is detected", async () => {
Expand Down Expand Up @@ -2049,7 +2057,7 @@ describe("stream-reset floods (CVE-2023-44487 rapid reset, CVE-2025-8671 MadeYou
expect(goawayErrorCode(goaway)).toBe(ErrorCode.ENHANCE_YOUR_CALM);
expect(c.frames.filter(f => f.type === FrameType.RST_STREAM).length).toBeGreaterThanOrEqual(1000);
expect(handlers).toBeGreaterThanOrEqual(1000);
expect(sessionErrorCode).toBe("ERR_HTTP2_ERROR");
expect(await sessionErrorCode).toBe("ERR_HTTP2_ERROR");
});
}

Expand Down Expand Up @@ -2150,7 +2158,7 @@ describe("stream-reset floods (CVE-2023-44487 rapid reset, CVE-2025-8671 MadeYou
c.send(pairs(1200, 1, rstStream));
const goaway = await c.waitFor(calm, 10_000);
expect(goaway.payload.subarray(8).toString()).toBe("too many stream resets");
expect(state.sessionErrorCode).toBe("ERR_HTTP2_ERROR");
expect(await state.sessionErrorCode).toBe("ERR_HTTP2_ERROR");
});
});

Expand All @@ -2164,7 +2172,7 @@ describe("stream-reset floods (CVE-2023-44487 rapid reset, CVE-2025-8671 MadeYou
c.send(Buffer.concat(ids.map(madeYouReset["WINDOW_UPDATE with a 0 increment"])));
const goaway = await c.waitFor(calm, 10_000);
expect(goaway.payload.subarray(8).toString()).toBe("too many stream resets");
expect(state.sessionErrorCode).toBe("ERR_HTTP2_ERROR");
expect(await state.sessionErrorCode).toBe("ERR_HTTP2_ERROR");
});
});
});
Loading
Loading