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
41 changes: 35 additions & 6 deletions src/js/node/http2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2083,6 +2083,13 @@ function markStreamClosed(stream: Http2Stream) {
markWritableDone(stream);
}
}
// writeStream callback: a frame that settles after the stream is destroyed was dropped during
// teardown, so report an error to the user's write callback (onwriteError never emits 'drain',
// and errorOrDestroy is a no-op on a destroyed stream).
function onWriteStreamDone(this: Http2Stream, callback, err) {
if (err == null && this.destroyed) return callback($ERR_HTTP2_INVALID_STREAM());
callback(err);
}
function rstNextTick(id: number, rstCode: number) {
const session = this as Http2Session;
session[bunHTTP2Native]?.rstStream(id, rstCode);
Expand Down Expand Up @@ -2523,15 +2530,17 @@ class Http2Stream extends Duplex {
}
}
const chunk = Buffer.concat(chunks || []);
native.writeStream(this.#id, chunk, undefined, false, callback);
native.writeStream(this.#id, chunk, undefined, false, onWriteStreamDone.bind(this, callback));
if (onClientStreamBodyChunkSentChannel.hasSubscribers && this instanceof ClientHttp2Stream) {
onClientStreamBodyChunkSentChannel.publish({ stream: this, writev: true, data, encoding: "" });
}
return;
}
}
// No session/native: the session was destroyed. Reporting success would emit 'drain' and make
// write() return true forever, so a backpressured producer never stops.
if (typeof callback == "function") {
callback();
callback($ERR_HTTP2_INVALID_STREAM());
}
}
_write(chunk, encoding, callback) {
Expand All @@ -2552,15 +2561,17 @@ class Http2Stream extends Duplex {
wireChunk = Buffer.from(chunk, encoding);
wireEncoding = undefined;
}
native.writeStream(this.#id, wireChunk, wireEncoding, false, callback);
native.writeStream(this.#id, wireChunk, wireEncoding, false, onWriteStreamDone.bind(this, callback));
if (onClientStreamBodyChunkSentChannel.hasSubscribers && this instanceof ClientHttp2Stream) {
onClientStreamBodyChunkSentChannel.publish({ stream: this, writev: false, data: chunk, encoding });
}
return;
}
}
// No session/native: the session was destroyed. Reporting success would emit 'drain' and make
// write() return true forever, so a backpressured producer never stops.
if (typeof callback == "function") {
callback();
callback($ERR_HTTP2_INVALID_STREAM());
}
}

Expand Down Expand Up @@ -4066,7 +4077,13 @@ class ServerHttp2Session extends Http2Session {
if (parser) {
// Like Node's Http2Stream._destroy: a received GOAWAY's code takes
// precedence over the destroy code when streams are torn down.
parser.emitErrorToAllStreams(this[kGoawayCode] || code || constants.NGHTTP2_NO_ERROR);
const streamRstCode = this[kGoawayCode] || code || constants.NGHTTP2_NO_ERROR;
// A non-numeric code is rejected by emitErrorToAllStreams below; only run the synchronous
// per-stream destroy when the code is valid so that rejection is still observable.
if (typeof streamRstCode === "number") {
parser.forEachStream(sessionDestroyStream.bind(this, streamRstCode));
}
parser.emitErrorToAllStreams(streamRstCode);
parser.detach();
this.#parser = null;
}
Expand All @@ -4089,6 +4106,12 @@ function destroySelfOnEnd(this: Http2Stream) {
function streamCancel(stream: Http2Stream) {
stream.close(NGHTTP2_CANCEL);
}
// session.destroy(): destroy still-open streams before native drops their queued DATA frames so
// the dropped-frame Writable callback sees kDestroyed and afterWrite skips 'drain'. Streams
// already marked closed completed normally and are not surfaced as errored here.
function sessionDestroyStream(this: Http2Session, rstCode: number, stream: Http2Stream) {
if (stream && !stream.destroyed && !stream.closed) emitStreamErrorNT(this, stream, rstCode, true, false);
}

// After the socket is gone a graceful close can never complete — the parser
// is detached, so the stream's writable side has nothing left to flush
Expand Down Expand Up @@ -4881,7 +4904,13 @@ class ClientHttp2Session extends Http2Session {
}
// Like Node's Http2Stream._destroy: a received GOAWAY's code takes
// precedence over the destroy code when streams are torn down.
parser.emitErrorToAllStreams(this[kGoawayCode] || (code !== undefined ? code : constants.NGHTTP2_CANCEL));
const streamRstCode = this[kGoawayCode] || (code !== undefined ? code : constants.NGHTTP2_CANCEL);
// A non-numeric code is rejected by emitErrorToAllStreams below; only run the synchronous
// per-stream destroy when the code is valid so that rejection is still observable.
if (typeof streamRstCode === "number") {
parser.forEachStream(sessionDestroyStream.bind(this, streamRstCode));
}
parser.emitErrorToAllStreams(streamRstCode);
parser.detach();
}
this.#parser = null;
Expand Down
112 changes: 112 additions & 0 deletions test/js/node/http2/node-http2-session-destroy-backpressure.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,112 @@
import { expect, test } from "bun:test";
import { bunEnv, bunExe } from "harness";

// session.destroy() with a DATA frame queued behind flow control must fail the held write
// callback so the stream does not emit 'drain' and later write() does not report success.
test("http2 client session.destroy() with a flow-control-blocked write does not emit 'drain' or accept further writes", async () => {
await using proc = Bun.spawn({
cmd: [
bunExe(),
"-e",
`
const http2 = require("node:http2");
const net = require("node:net");
const frame = (t, fl, sid, p = Buffer.alloc(0)) => {
const b = Buffer.alloc(9 + p.length);
b.writeUIntBE(p.length, 0, 3);
b[3] = t;
b[4] = fl;
b.writeUInt32BE(sid >>> 0, 5);
p.copy(b, 9);
return b;
};
const PREFACE = Buffer.from("PRI * HTTP/2.0\\r\\n\\r\\nSM\\r\\n\\r\\n");
// A raw TCP peer that speaks just enough HTTP/2 to establish the connection and then
// withholds WINDOW_UPDATE so the client's DATA stays flow-control blocked.
const server = net.createServer(s => {
let buf = Buffer.alloc(0);
let seenPreface = false;
s.on("error", () => {});
s.on("data", d => {
buf = Buffer.concat([buf, d]);
if (!seenPreface) {
if (buf.length < PREFACE.length) return;
seenPreface = true;
buf = buf.slice(PREFACE.length);
s.write(frame(4, 0, 0));
}
while (buf.length >= 9) {
const len = buf.readUIntBE(0, 3);
if (buf.length < 9 + len) break;
const t = buf[3];
const fl = buf[4];
const pay = buf.slice(9, 9 + len);
buf = buf.slice(9 + len);
if (t === 4 && (fl & 1) === 0) s.write(frame(4, 1, 0));
else if (t === 6 && (fl & 1) === 0) s.write(frame(6, 1, 0, pay));
}
});
});
server.listen(0, "127.0.0.1", () => {
const session = http2.connect("http://127.0.0.1:" + server.address().port);
session.on("error", () => {});
session.once("remoteSettings", () => {
const stream = session.request({ ":path": "/u", ":method": "POST" });
stream.on("error", () => {});
// One write larger than the 65535-byte initial window: the first 65535 bytes go out
// immediately and the remainder is queued natively, holding the _write callback until
// the peer reopens the window (which never happens here).
let writeCbErrorCode = null;
const backpressured = stream.write(Buffer.alloc(65535 + 32768, 0x41), err => {
writeCbErrorCode = err ? err.code : "none";
}) === false;
let drainsAfterDestroy = 0;
let writeOkAfterDestroy = 0;
stream.once("drain", function onDrain() {
drainsAfterDestroy++;
// The canonical backpressured producer: on 'drain', keep writing until write()
// returns false again. Budget capped so the broken case terminates.
let budget = 32;
while (budget-- > 0 && stream.write(Buffer.alloc(16384))) writeOkAfterDestroy++;
stream.once("drain", onDrain);
});
// Destroy from outside the native dispatch that delivered remoteSettings so the
// session's stream teardown runs its own dispatch (where the queued-frame callback
// fires before the stream is destroyed).
process.nextTick(() => {
session.destroy();
stream.once("close", () => {
server.close();
console.log(
JSON.stringify({
backpressured,
drainsAfterDestroy,
writeOkAfterDestroy,
writeCbErrorCode,
}),
);
process.exit(0);
});
});
});
});
`,
],
env: bunEnv,
stdout: "pipe",
stderr: "pipe",
});

const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]);

expect({ result: JSON.parse(stdout.trim() || "null"), stderr, exitCode }).toEqual({
result: {
backpressured: true,
drainsAfterDestroy: 0,
writeOkAfterDestroy: 0,
writeCbErrorCode: expect.stringMatching(/^(ERR_HTTP2_INVALID_STREAM|ECANCELED|ERR_STREAM_DESTROYED)$/),
},
stderr: expect.anything(),
exitCode: 0,
});
}, 30_000);
Loading