From 4277fa003cf1bc33a4a8f93dc23e988a85886528 Mon Sep 17 00:00:00 2001 From: robobun Date: Wed, 22 Apr 2026 08:30:47 +0000 Subject: [PATCH 01/10] http: publish to http.server.* diagnostics_channel channels MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Node's node:http server publishes to three diagnostics_channel channels at request/response lifecycle points, which observability libraries (OpenTelemetry, Prometheus, APM tooling) rely on for instrumentation. Bun had not implemented these, so subscribers silently collected nothing. Adds publishes for: - http.server.response.created — { request, response } - http.server.request.start — { request, response, socket, server } - http.server.response.finish — { request, response, socket, server } Channel placement mirrors Node's lib/_http_server.js so all observable edges match Node: - response.created is published from inside the ServerResponse constructor, so it fires for direct `new http.ServerResponse(req)` (the pattern used by light-my-request / fastify.inject() / mocks) in addition to the live-server path. Node does the same. - request.start is published once per non-upgrade request before the dispatch chain, so it fires on the normal, checkContinue, checkExpectation, 417, and dropRequest/503 branches — not only the two that reach server.emit('request', ...). - response.finish is wired via an always-attached res.on('finish') listener (matching Node's resOnFinish) with the hasSubscribers guard inside the listener, so subscribers that register between request arrival and response finish are still observed. Synchronous publishes keep a call-site hasSubscribers guard so no payload object is allocated when nobody is listening. Ordering matches Node: response.created fires before the 'request' event so user handlers can't have mutated the response yet, and response.finish fires after the response is sent. Fixes #29586 --- src/js/node/_http_server.ts | 45 +++ .../diagnostics_channel.test.ts | 292 ++++++++++++++++++ 2 files changed, 337 insertions(+) diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index 8e515ee41c0e..569a27127f4d 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -91,6 +91,26 @@ const MathFloor = Math.floor; let cluster; +// diagnostics_channel channels for the HTTP server. Mirrors Node's +// lib/_http_server.js. Inactive channels are no-ops until someone subscribes. +const dc = require("node:diagnostics_channel"); +const onRequestStartChannel = dc.channel("http.server.request.start"); +const onResponseCreatedChannel = dc.channel("http.server.response.created"); +const onResponseFinishChannel = dc.channel("http.server.response.finish"); + +function emitResponseFinishChannel(this: { req; res; socket; server }) { + // Re-checked here (not at listener-attach) so subscribers that attach + // between request-arrival and response-finish are still observed. + if (onResponseFinishChannel.hasSubscribers) { + onResponseFinishChannel.publish({ + request: this.req, + response: this.res, + socket: this.socket, + server: this.server, + }); + } +} + function emitCloseServer(self: Server) { callCloseCallback(self); self.emit("close"); @@ -660,6 +680,13 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort http_res.once("finish", stopServerResponsePerf); } + // Attached unconditionally to match Node's resOnFinish; the + // hasSubscribers check happens in emitResponseFinishChannel. + http_res.on( + "finish", + emitResponseFinishChannel.bind({ req: http_req, res: http_res, socket, server }), + ); + setIsNextIncomingMessageHTTPS(prevIsNextIncomingMessageHTTPS); handle.onabort = onServerRequestEvent.bind(socket); // start buffering data if any, the user will need to resume() or .on("data") to read it @@ -737,6 +764,18 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort http_req._dumpAndCloseReadable(); } + // Node publishes http.server.request.start once per non-upgrade + // request in parserOnIncoming, before any branching, so it fires on + // the dropRequest/503, checkContinue, checkExpectation, 417 and + // normal paths alike. Mirror that here. + if (!is_upgrade && onRequestStartChannel.hasSubscribers) { + onRequestStartChannel.publish({ + request: http_req, + response: http_res, + socket, + server, + }); + } if (reachedRequestsLimit) { server.emit("dropRequest", http_req, socket); http_res.writeHead(503); @@ -1482,6 +1521,12 @@ function ServerResponse(req, options): void { this.statusCode = 200; this.statusMessage = undefined; this.chunkedEncoding = false; + + // Publish response.created from the constructor (matches Node) so it also + // fires for direct `new ServerResponse(req)` use — light-my-request etc. + if (onResponseCreatedChannel.hasSubscribers) { + onResponseCreatedChannel.publish({ request: req, response: this }); + } } $toClass(ServerResponse, "ServerResponse", OutgoingMessage); diff --git a/test/js/node/diagnostics_channel/diagnostics_channel.test.ts b/test/js/node/diagnostics_channel/diagnostics_channel.test.ts index 37dfd54d7a8f..b79918293d40 100644 --- a/test/js/node/diagnostics_channel/diagnostics_channel.test.ts +++ b/test/js/node/diagnostics_channel/diagnostics_channel.test.ts @@ -2,6 +2,8 @@ import { gc } from "bun"; import { beforeEach, describe, expect, mock, test } from "bun:test"; import { AsyncLocalStorage } from "node:async_hooks"; import { channel, Channel, hasSubscribers, subscribe, unsubscribe } from "node:diagnostics_channel"; +import http from "node:http"; +import net from "node:net"; describe("Channel", () => { // test-diagnostics-channel-has-subscribers.js @@ -345,6 +347,296 @@ describe("TracingChannel", () => { test.todo("TODO"); }); +describe("http server channels (#29586)", () => { + test("publishes http.server.request.start, response.created, response.finish", async () => { + const channelNames = [ + "http.server.response.created", + "http.server.request.start", + "http.server.response.finish", + ]; + const events: { channel: string; payload: Record }[] = []; + const subs: [ReturnType, (msg: unknown) => void][] = []; + + // Hold direct channel refs so subscriptions can't be silently lost to GC, + // and also so we can cleanly unsubscribe in the finally. + for (const name of channelNames) { + const ch = channel(name); + const sub = (msg: unknown) => { + events.push({ channel: ch.name as string, payload: msg as Record }); + }; + ch.subscribe(sub); + subs.push([ch, sub]); + } + + try { + await using server = http.createServer((_req, res) => res.end("ok")); + await new Promise(resolve => server.listen(0, resolve)); + const { port } = server.address(); + await (await fetch(`http://127.0.0.1:${port}/`)).text(); + // Wait for the "finish" nextTick and any tail events to drain. + await new Promise(resolve => setImmediate(resolve)); + await new Promise(resolve => setImmediate(resolve)); + + expect(events.map(e => e.channel)).toEqual([ + "http.server.response.created", + "http.server.request.start", + "http.server.response.finish", + ]); + + // response.created: { request, response } only + const created = events[0].payload; + expect(Object.keys(created).sort()).toEqual(["request", "response"]); + expect(created.request).toBeInstanceOf(http.IncomingMessage); + expect(created.response).toBeInstanceOf(http.ServerResponse); + + // request.start: { request, response, socket, server } + const reqStart = events[1].payload; + expect(Object.keys(reqStart).sort()).toEqual(["request", "response", "server", "socket"]); + expect(reqStart.request).toBe(created.request); + expect(reqStart.response).toBe(created.response); + expect(reqStart.server).toBe(server); + expect(reqStart.socket).toBe((reqStart.request as any).socket); + + // response.finish: { request, response, socket, server } + const resFinish = events[2].payload; + expect(Object.keys(resFinish).sort()).toEqual(["request", "response", "server", "socket"]); + expect(resFinish.request).toBe(created.request); + expect(resFinish.response).toBe(created.response); + expect(resFinish.server).toBe(server); + expect(resFinish.socket).toBe(reqStart.socket); + } finally { + for (const [ch, sub] of subs) ch.unsubscribe(sub); + } + }); + + // Node's contract: response.created fires before any user handler can mutate + // the response. Mirror of Node's test-diagnostic-channel-http-response-created.js. + test("response.created fires before the request handler runs", async () => { + const created = channel("http.server.response.created"); + const finish = channel("http.server.response.finish"); + const snapshots: { event: string; baz: unknown }[] = []; + + const onCreated = (msg: any) => { + snapshots.push({ event: "created", baz: msg.response.getHeader("baz") }); + }; + const onFinish = (msg: any) => { + snapshots.push({ event: "finish", baz: msg.response.getHeader("baz") }); + }; + created.subscribe(onCreated); + finish.subscribe(onFinish); + + try { + await using server = http.createServer((_req, res) => { + res.setHeader("baz", "bar"); + res.end("done"); + }); + await new Promise(resolve => server.listen(0, resolve)); + const { port } = server.address(); + await (await fetch(`http://127.0.0.1:${port}/`)).text(); + await new Promise(resolve => setImmediate(resolve)); + await new Promise(resolve => setImmediate(resolve)); + + expect(snapshots).toEqual([ + { event: "created", baz: undefined }, // fired before handler set the header + { event: "finish", baz: "bar" }, // fired after the handler completed + ]); + } finally { + created.unsubscribe(onCreated); + finish.unsubscribe(onFinish); + } + }); + + // Subscribing after the request arrived but before the response finished + // must still deliver response.finish — the 'finish' listener is attached + // unconditionally, matching Node's resOnFinish, so this works no matter + // when subscription happens. + test("response.finish delivers to subscribers added after the request arrived", async () => { + const finish = channel("http.server.response.finish"); + const received: any[] = []; + const onFinish = (msg: any) => { + received.push(msg); + }; + + let subscribed = false; + const { promise: handlerEntered, resolve: onHandlerEntered } = Promise.withResolvers(); + const { promise: mayFinish, resolve: continueFinish } = Promise.withResolvers(); + + try { + await using server = http.createServer(async (_req, res) => { + onHandlerEntered(); + // Hold the response open until we've subscribed. + await mayFinish; + res.end("done"); + }); + await new Promise(resolve => server.listen(0, resolve)); + const { port } = server.address(); + const fetched = fetch(`http://127.0.0.1:${port}/`); + await handlerEntered; + + // Subscribe *after* request arrival, *before* response finish. + finish.subscribe(onFinish); + subscribed = true; + continueFinish(); + + await (await fetched).text(); + await new Promise(resolve => setImmediate(resolve)); + await new Promise(resolve => setImmediate(resolve)); + + expect(received).toHaveLength(1); + expect(Object.keys(received[0]).sort()).toEqual(["request", "response", "server", "socket"]); + } finally { + if (subscribed) finish.unsubscribe(onFinish); + } + }); + + // Node publishes response.created from inside the ServerResponse + // constructor (lib/_http_server.js), so `new http.ServerResponse(req)` + // — the pattern used by light-my-request, fastify.inject(), etc. — fires + // it too. Mirror that. + test("response.created fires for direct new http.ServerResponse()", () => { + const created = channel("http.server.response.created"); + const received: any[] = []; + const onCreated = (msg: any) => { + received.push(msg); + }; + created.subscribe(onCreated); + try { + const req = new http.IncomingMessage(new net.Socket()); + const res = new http.ServerResponse(req); + expect(received).toHaveLength(1); + expect(Object.keys(received[0]).sort()).toEqual(["request", "response"]); + expect(received[0].request).toBe(req); + expect(received[0].response).toBe(res); + } finally { + created.unsubscribe(onCreated); + } + }); + + // Node publishes request.start once per non-upgrade request in + // parserOnIncoming — before the branching — so it fires on the + // checkContinue, checkExpectation, 417 auto-response and dropRequest/503 + // paths too, not only the plain server.emit('request') path. + test("request.start fires on checkContinue path", async () => { + const requestStart = channel("http.server.request.start"); + const received: any[] = []; + const onStart = (msg: any) => { + received.push(msg); + }; + requestStart.subscribe(onStart); + try { + await using server = http.createServer((_req, res) => res.end("fallback")); + server.on("checkContinue", (_req, res) => { + res.writeContinue(); + res.end("cc"); + }); + await new Promise(resolve => server.listen(0, resolve)); + const { port } = server.address(); + await new Promise((resolve, reject) => { + const req = http.request({ port, method: "POST", headers: { expect: "100-continue" } }); + req.on("response", res => { + res.resume(); + res.on("end", resolve); + }); + req.on("error", reject); + req.end(); + }); + await new Promise(resolve => setImmediate(resolve)); + expect(received).toHaveLength(1); + expect(Object.keys(received[0]).sort()).toEqual(["request", "response", "server", "socket"]); + expect(received[0].server).toBe(server); + } finally { + requestStart.unsubscribe(onStart); + } + }); + + test("request.start fires on 417 Expectation Failed path", async () => { + const requestStart = channel("http.server.request.start"); + const received: any[] = []; + const onStart = (msg: any) => { + received.push(msg); + }; + requestStart.subscribe(onStart); + let client: any; + try { + await using server = http.createServer((_req, res) => res.end("x")); + await new Promise(resolve => server.listen(0, resolve)); + const { port } = server.address(); + // Raw socket — we want `Expect: weird` to reach the 417 branch. + await new Promise((resolve, reject) => { + client = net.connect(port, () => { + client.write("GET / HTTP/1.1\r\nHost: x\r\nExpect: weird\r\n\r\n"); + }); + let buf = ""; + client.on("data", d => { + buf += d; + if (buf.includes("\r\n\r\n")) { + // Full response received; don't wait for close (server keeps + // the connection alive). + buf.startsWith("HTTP/1.1 417") ? resolve() : reject(new Error(buf)); + } + }); + client.on("error", reject); + }); + await new Promise(resolve => setImmediate(resolve)); + expect(received).toHaveLength(1); + expect(Object.keys(received[0]).sort()).toEqual(["request", "response", "server", "socket"]); + expect(received[0].server).toBe(server); + } finally { + client?.destroy?.(); + requestStart.unsubscribe(onStart); + } + }); + + test("request.start does NOT fire on upgrade path", async () => { + const requestStart = channel("http.server.request.start"); + const receivedStart: any[] = []; + const onStart = (msg: any) => { + receivedStart.push(msg); + }; + requestStart.subscribe(onStart); + try { + await using server = http.createServer(); + server.on("upgrade", (_req, socket) => { + // Graceful FIN rather than RST — avoids a flaky ECONNRESET on the + // client side that could race the assertion. + socket.end(); + }); + await new Promise(resolve => server.listen(0, resolve)); + const { port } = server.address(); + await new Promise((resolve, reject) => { + const client = net.connect(port, () => { + client.write( + "GET / HTTP/1.1\r\nHost: x\r\nUpgrade: websocket\r\nConnection: Upgrade\r\n\r\n", + ); + }); + client.on("close", () => resolve()); + // Swallow ECONNRESET etc. — we only care that the connection + // terminates, not how. + client.on("error", () => {}); + }); + await new Promise(resolve => setImmediate(resolve)); + // Node gates request.start on `!is_upgrade` (it returns early in + // parserOnIncoming before constructing ServerResponse), and so do we. + // Note: Bun currently *does* construct ServerResponse for upgrades + // before checking is_upgrade, so response.created still fires — that + // is a pre-existing divergence out of scope for this PR. + expect(receivedStart).toHaveLength(0); + } finally { + requestStart.unsubscribe(onStart); + } + }); + + test("server works normally when nobody subscribed", async () => { + // No subscribers means no publish payload is allocated — just prove the + // normal request/response path still works with the channel plumbing in. + await using server = http.createServer((_req, res) => res.end("ok")); + await new Promise(resolve => server.listen(0, resolve)); + const { port } = server.address(); + const body = await (await fetch(`http://127.0.0.1:${port}/`)).text(); + expect(body).toBe("ok"); + }); +}); + const mocks = new Map(); function mustCall(fn: (...args: any[]) => T, expected?: number) { From 568a92c665179df5923c5a57454eece75849d5d8 Mon Sep 17 00:00:00 2001 From: "autofix-ci[bot]" <114827586+autofix-ci[bot]@users.noreply.github.com> Date: Wed, 22 Apr 2026 09:21:08 +0000 Subject: [PATCH 02/10] [autofix.ci] apply automated fixes --- .../diagnostics_channel/diagnostics_channel.test.ts | 10 ++-------- 1 file changed, 2 insertions(+), 8 deletions(-) diff --git a/test/js/node/diagnostics_channel/diagnostics_channel.test.ts b/test/js/node/diagnostics_channel/diagnostics_channel.test.ts index b79918293d40..99b556dca88a 100644 --- a/test/js/node/diagnostics_channel/diagnostics_channel.test.ts +++ b/test/js/node/diagnostics_channel/diagnostics_channel.test.ts @@ -349,11 +349,7 @@ describe("TracingChannel", () => { describe("http server channels (#29586)", () => { test("publishes http.server.request.start, response.created, response.finish", async () => { - const channelNames = [ - "http.server.response.created", - "http.server.request.start", - "http.server.response.finish", - ]; + const channelNames = ["http.server.response.created", "http.server.request.start", "http.server.response.finish"]; const events: { channel: string; payload: Record }[] = []; const subs: [ReturnType, (msg: unknown) => void][] = []; @@ -605,9 +601,7 @@ describe("http server channels (#29586)", () => { const { port } = server.address(); await new Promise((resolve, reject) => { const client = net.connect(port, () => { - client.write( - "GET / HTTP/1.1\r\nHost: x\r\nUpgrade: websocket\r\nConnection: Upgrade\r\n\r\n", - ); + client.write("GET / HTTP/1.1\r\nHost: x\r\nUpgrade: websocket\r\nConnection: Upgrade\r\n\r\n"); }); client.on("close", () => resolve()); // Swallow ECONNRESET etc. — we only care that the connection From d9935c2df0f0c742b2d3313eb649c976557a076d Mon Sep 17 00:00:00 2001 From: robobun Date: Wed, 22 Apr 2026 09:34:45 +0000 Subject: [PATCH 03/10] tighten http server channel tests and inline publish Pull the listen/port and setImmediate-drain boilerplate into two helpers in the describe block; drop the channelNames/subs ceremony in the multi-channel test in favor of three flat subscribe/unsubscribe pairs; collapse a multi-line request.start publish + trim its comment to a single line. --- src/js/node/_http_server.ts | 13 +--- .../diagnostics_channel.test.ts | 74 +++++++++---------- 2 files changed, 36 insertions(+), 51 deletions(-) diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index 569a27127f4d..d20cbe25ef02 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -764,17 +764,10 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort http_req._dumpAndCloseReadable(); } - // Node publishes http.server.request.start once per non-upgrade - // request in parserOnIncoming, before any branching, so it fires on - // the dropRequest/503, checkContinue, checkExpectation, 417 and - // normal paths alike. Mirror that here. + // Match Node's parserOnIncoming: publish once, before branching, for + // non-upgrade requests (fires on 503/checkContinue/417/normal paths). if (!is_upgrade && onRequestStartChannel.hasSubscribers) { - onRequestStartChannel.publish({ - request: http_req, - response: http_res, - socket, - server, - }); + onRequestStartChannel.publish({ request: http_req, response: http_res, socket, server }); } if (reachedRequestsLimit) { server.emit("dropRequest", http_req, socket); diff --git a/test/js/node/diagnostics_channel/diagnostics_channel.test.ts b/test/js/node/diagnostics_channel/diagnostics_channel.test.ts index 99b556dca88a..7752fd62bcea 100644 --- a/test/js/node/diagnostics_channel/diagnostics_channel.test.ts +++ b/test/js/node/diagnostics_channel/diagnostics_channel.test.ts @@ -348,30 +348,29 @@ describe("TracingChannel", () => { }); describe("http server channels (#29586)", () => { + const listen = (s: http.Server) => + new Promise(r => s.listen(0, () => r((s.address() as any).port))); + // Two ticks drains both the 'finish' nextTick and any tail events. + const drain = async () => { + await new Promise(r => setImmediate(r)); + await new Promise(r => setImmediate(r)); + }; + test("publishes http.server.request.start, response.created, response.finish", async () => { - const channelNames = ["http.server.response.created", "http.server.request.start", "http.server.response.finish"]; - const events: { channel: string; payload: Record }[] = []; - const subs: [ReturnType, (msg: unknown) => void][] = []; - - // Hold direct channel refs so subscriptions can't be silently lost to GC, - // and also so we can cleanly unsubscribe in the finally. - for (const name of channelNames) { - const ch = channel(name); - const sub = (msg: unknown) => { - events.push({ channel: ch.name as string, payload: msg as Record }); - }; - ch.subscribe(sub); - subs.push([ch, sub]); - } + const events: { channel: string; payload: any }[] = []; + const push = (name: string) => (payload: any) => events.push({ channel: name, payload }); + const onCreated = push("http.server.response.created"); + const onStart = push("http.server.request.start"); + const onFinish = push("http.server.response.finish"); + channel("http.server.response.created").subscribe(onCreated); + channel("http.server.request.start").subscribe(onStart); + channel("http.server.response.finish").subscribe(onFinish); try { await using server = http.createServer((_req, res) => res.end("ok")); - await new Promise(resolve => server.listen(0, resolve)); - const { port } = server.address(); + const port = await listen(server); await (await fetch(`http://127.0.0.1:${port}/`)).text(); - // Wait for the "finish" nextTick and any tail events to drain. - await new Promise(resolve => setImmediate(resolve)); - await new Promise(resolve => setImmediate(resolve)); + await drain(); expect(events.map(e => e.channel)).toEqual([ "http.server.response.created", @@ -401,7 +400,9 @@ describe("http server channels (#29586)", () => { expect(resFinish.server).toBe(server); expect(resFinish.socket).toBe(reqStart.socket); } finally { - for (const [ch, sub] of subs) ch.unsubscribe(sub); + channel("http.server.response.created").unsubscribe(onCreated); + channel("http.server.request.start").unsubscribe(onStart); + channel("http.server.response.finish").unsubscribe(onFinish); } }); @@ -426,11 +427,9 @@ describe("http server channels (#29586)", () => { res.setHeader("baz", "bar"); res.end("done"); }); - await new Promise(resolve => server.listen(0, resolve)); - const { port } = server.address(); + const port = await listen(server); await (await fetch(`http://127.0.0.1:${port}/`)).text(); - await new Promise(resolve => setImmediate(resolve)); - await new Promise(resolve => setImmediate(resolve)); + await drain(); expect(snapshots).toEqual([ { event: "created", baz: undefined }, // fired before handler set the header @@ -464,8 +463,7 @@ describe("http server channels (#29586)", () => { await mayFinish; res.end("done"); }); - await new Promise(resolve => server.listen(0, resolve)); - const { port } = server.address(); + const port = await listen(server); const fetched = fetch(`http://127.0.0.1:${port}/`); await handlerEntered; @@ -475,8 +473,7 @@ describe("http server channels (#29586)", () => { continueFinish(); await (await fetched).text(); - await new Promise(resolve => setImmediate(resolve)); - await new Promise(resolve => setImmediate(resolve)); + await drain(); expect(received).toHaveLength(1); expect(Object.keys(received[0]).sort()).toEqual(["request", "response", "server", "socket"]); @@ -525,8 +522,7 @@ describe("http server channels (#29586)", () => { res.writeContinue(); res.end("cc"); }); - await new Promise(resolve => server.listen(0, resolve)); - const { port } = server.address(); + const port = await listen(server); await new Promise((resolve, reject) => { const req = http.request({ port, method: "POST", headers: { expect: "100-continue" } }); req.on("response", res => { @@ -536,7 +532,7 @@ describe("http server channels (#29586)", () => { req.on("error", reject); req.end(); }); - await new Promise(resolve => setImmediate(resolve)); + await drain(); expect(received).toHaveLength(1); expect(Object.keys(received[0]).sort()).toEqual(["request", "response", "server", "socket"]); expect(received[0].server).toBe(server); @@ -555,8 +551,7 @@ describe("http server channels (#29586)", () => { let client: any; try { await using server = http.createServer((_req, res) => res.end("x")); - await new Promise(resolve => server.listen(0, resolve)); - const { port } = server.address(); + const port = await listen(server); // Raw socket — we want `Expect: weird` to reach the 417 branch. await new Promise((resolve, reject) => { client = net.connect(port, () => { @@ -573,7 +568,7 @@ describe("http server channels (#29586)", () => { }); client.on("error", reject); }); - await new Promise(resolve => setImmediate(resolve)); + await drain(); expect(received).toHaveLength(1); expect(Object.keys(received[0]).sort()).toEqual(["request", "response", "server", "socket"]); expect(received[0].server).toBe(server); @@ -597,8 +592,7 @@ describe("http server channels (#29586)", () => { // client side that could race the assertion. socket.end(); }); - await new Promise(resolve => server.listen(0, resolve)); - const { port } = server.address(); + const port = await listen(server); await new Promise((resolve, reject) => { const client = net.connect(port, () => { client.write("GET / HTTP/1.1\r\nHost: x\r\nUpgrade: websocket\r\nConnection: Upgrade\r\n\r\n"); @@ -608,7 +602,7 @@ describe("http server channels (#29586)", () => { // terminates, not how. client.on("error", () => {}); }); - await new Promise(resolve => setImmediate(resolve)); + await drain(); // Node gates request.start on `!is_upgrade` (it returns early in // parserOnIncoming before constructing ServerResponse), and so do we. // Note: Bun currently *does* construct ServerResponse for upgrades @@ -621,11 +615,9 @@ describe("http server channels (#29586)", () => { }); test("server works normally when nobody subscribed", async () => { - // No subscribers means no publish payload is allocated — just prove the - // normal request/response path still works with the channel plumbing in. + // No subscribers: prove the channel plumbing doesn't break the hot path. await using server = http.createServer((_req, res) => res.end("ok")); - await new Promise(resolve => server.listen(0, resolve)); - const { port } = server.address(); + const port = await listen(server); const body = await (await fetch(`http://127.0.0.1:${port}/`)).text(); expect(body).toBe("ok"); }); From 695396b94b0f440b3e49da954f9df1645e5471d2 Mon Sep 17 00:00:00 2001 From: "autofix-ci[bot]" <114827586+autofix-ci[bot]@users.noreply.github.com> Date: Wed, 22 Apr 2026 09:36:36 +0000 Subject: [PATCH 04/10] [autofix.ci] apply automated fixes --- test/js/node/diagnostics_channel/diagnostics_channel.test.ts | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/test/js/node/diagnostics_channel/diagnostics_channel.test.ts b/test/js/node/diagnostics_channel/diagnostics_channel.test.ts index 7752fd62bcea..67b1b2c1752e 100644 --- a/test/js/node/diagnostics_channel/diagnostics_channel.test.ts +++ b/test/js/node/diagnostics_channel/diagnostics_channel.test.ts @@ -348,8 +348,7 @@ describe("TracingChannel", () => { }); describe("http server channels (#29586)", () => { - const listen = (s: http.Server) => - new Promise(r => s.listen(0, () => r((s.address() as any).port))); + const listen = (s: http.Server) => new Promise(r => s.listen(0, () => r((s.address() as any).port))); // Two ticks drains both the 'finish' nextTick and any tail events. const drain = async () => { await new Promise(r => setImmediate(r)); From fb090374c4c89bb97c30c4c4277705751af501e7 Mon Sep 17 00:00:00 2001 From: "autofix-ci[bot]" <114827586+autofix-ci[bot]@users.noreply.github.com> Date: Thu, 18 Jun 2026 08:47:28 +0000 Subject: [PATCH 05/10] [autofix.ci] apply automated fixes (attempt 2/3) --- src/js/node/_http_server.ts | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index d20cbe25ef02..5c2ebecce718 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -682,10 +682,7 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort // Attached unconditionally to match Node's resOnFinish; the // hasSubscribers check happens in emitResponseFinishChannel. - http_res.on( - "finish", - emitResponseFinishChannel.bind({ req: http_req, res: http_res, socket, server }), - ); + http_res.on("finish", emitResponseFinishChannel.bind({ req: http_req, res: http_res, socket, server })); setIsNextIncomingMessageHTTPS(prevIsNextIncomingMessageHTTPS); handle.onabort = onServerRequestEvent.bind(socket); From d2594e333634cd003cefabc4beb16272c19ce4ae Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Thu, 18 Jun 2026 09:24:43 +0000 Subject: [PATCH 06/10] http: publish response.finish before socket end; wire channels into http2 allowHTTP1 fallback Register the http.server.response.finish listener before endSocketOnFinishIfNeeded so the publish observes the socket before its writable side is ended on the Connection: close path, matching Node's resOnFinish (which publishes before socket.destroySoon()). The response.created publish lives in the ServerResponse constructor, so it already fired on the http2 createSecureServer({ allowHTTP1: true }) HTTP/1 fallback path (connectionListenerHTTP1), but that path never published request.start or response.finish, leaving orphaned events. Publish all three there, like Node routes allowHTTP1 through the full http1 connectionListener. --- src/js/node/_http_server.ts | 12 ++- src/js/node/http2.ts | 22 ++++++ .../diagnostics_channel.test.ts | 77 +++++++++++++++++++ 3 files changed, 107 insertions(+), 4 deletions(-) diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index 5c2ebecce718..6f8326ce06fc 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -663,6 +663,14 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort if (!requestShouldKeepAlive(http_req)) { http_res[kMustCloseConnection] = true; } + + // Attached unconditionally to match Node's resOnFinish; the + // hasSubscribers check happens in emitResponseFinishChannel. Registered + // before endSocketOnFinishIfNeeded so the publish observes the socket + // before its writable side is ended (as Node's resOnFinish publishes + // before socket.destroySoon()). + http_res.on("finish", emitResponseFinishChannel.bind({ req: http_req, res: http_res, socket, server })); + http_res.once("finish", endSocketOnFinishIfNeeded.bind(undefined, socket, http_res)); if (hasObserver("http")) { @@ -680,10 +688,6 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort http_res.once("finish", stopServerResponsePerf); } - // Attached unconditionally to match Node's resOnFinish; the - // hasSubscribers check happens in emitResponseFinishChannel. - http_res.on("finish", emitResponseFinishChannel.bind({ req: http_req, res: http_res, socket, server })); - setIsNextIncomingMessageHTTPS(prevIsNextIncomingMessageHTTPS); handle.onabort = onServerRequestEvent.bind(socket); // start buffering data if any, the user will need to resume() or .on("data") to read it diff --git a/src/js/node/http2.ts b/src/js/node/http2.ts index daea2b34d26d..17726c14a265 100644 --- a/src/js/node/http2.ts +++ b/src/js/node/http2.ts @@ -5469,6 +5469,15 @@ function connectionListenerHTTP1(server, socket, options) { const ServerResponseClass = http1Options.ServerResponse || http.ServerResponse; const keepAliveTimeout = typeof server.keepAliveTimeout === "number" ? server.keepAliveTimeout : 5000; + // http.server.request.start / http.server.response.finish for the HTTP/1 + // fallback path (allowHTTP1). response.created is published by the + // ServerResponse constructor; publish the other two here so subscribers see + // the same three events Node fires on this path. Same channel objects as + // node:_http_server (keyed by name in diagnostics_channel's registry). + const dc = require("node:diagnostics_channel"); + const onRequestStartChannel = dc.channel("http.server.request.start"); + const onResponseFinishChannel = dc.channel("http.server.response.finish"); + const connections = (server[kHttp1Connections] ??= new SafeSet()); connections.add(socket); socket[kHttp1ActiveRequests] = 0; @@ -5524,6 +5533,9 @@ function connectionListenerHTTP1(server, socket, options) { }; const res = new ServerResponseClass(req); + // Stable reference for the diagnostics closure: the outer `req` is reused + // across pipelined requests on this connection. + const request = req; const handle = createHttp1FallbackResponseHandle(socket, shouldKeepAlive, keepAliveTimeout); handle.onfinished = function () { socket[kHttp1ActiveRequests] = Math.max(0, (socket[kHttp1ActiveRequests] || 1) - 1); @@ -5534,6 +5546,16 @@ function connectionListenerHTTP1(server, socket, options) { res[kHttp1ResponseHandle] = handle; res.assignSocket(socket); + // Attached unconditionally to match Node's resOnFinish; the hasSubscribers + // check happens inside. + res.on("finish", () => { + if (onResponseFinishChannel.hasSubscribers) { + onResponseFinishChannel.publish({ request, response: res, socket, server }); + } + }); + if (onRequestStartChannel.hasSubscribers) { + onRequestStartChannel.publish({ request, response: res, socket, server }); + } server.emit("request", req, res); return 0; }; diff --git a/test/js/node/diagnostics_channel/diagnostics_channel.test.ts b/test/js/node/diagnostics_channel/diagnostics_channel.test.ts index 67b1b2c1752e..1d46180bb29d 100644 --- a/test/js/node/diagnostics_channel/diagnostics_channel.test.ts +++ b/test/js/node/diagnostics_channel/diagnostics_channel.test.ts @@ -1,8 +1,11 @@ import { gc } from "bun"; import { beforeEach, describe, expect, mock, test } from "bun:test"; +import { tls as tlsCert } from "harness"; import { AsyncLocalStorage } from "node:async_hooks"; import { channel, Channel, hasSubscribers, subscribe, unsubscribe } from "node:diagnostics_channel"; import http from "node:http"; +import http2 from "node:http2"; +import https from "node:https"; import net from "node:net"; describe("Channel", () => { @@ -613,6 +616,80 @@ describe("http server channels (#29586)", () => { } }); + // Node's resOnFinish publishes before it ends the socket, so on the + // connection-close path a subscriber still sees a not-yet-ended socket. + // The finish listener is registered before endSocketOnFinishIfNeeded to + // preserve that ordering. + test("response.finish publishes before the socket is ended (Connection: close)", async () => { + const finish = channel("http.server.response.finish"); + let writableEndedAtPublish: boolean | undefined; + const onFinish = (msg: any) => { + writableEndedAtPublish = msg.socket.writableEnded; + }; + finish.subscribe(onFinish); + let client: any; + try { + await using server = http.createServer((_req, res) => res.end("ok")); + const port = await listen(server); + await new Promise((resolve, reject) => { + client = net.connect(port, () => { + client.write("GET / HTTP/1.1\r\nHost: x\r\nConnection: close\r\n\r\n"); + }); + let buf = ""; + client.on("data", d => (buf += d)); + client.on("close", () => resolve()); + client.on("error", reject); + }); + await drain(); + expect(writableEndedAtPublish).toBe(false); + } finally { + client?.destroy?.(); + finish.unsubscribe(onFinish); + } + }); + + // Node routes http2 allowHTTP1 through the full http1 connectionListener, so + // all three channels fire on the HTTP/1 fallback. Bun's fallback lives in + // http2.ts and must publish the same three. + test("all three channels fire on the http2 allowHTTP1 fallback", async () => { + const created = channel("http.server.response.created"); + const requestStart = channel("http.server.request.start"); + const finish = channel("http.server.response.finish"); + const seen: string[] = []; + const onCreated = () => seen.push("created"); + const onStart = () => seen.push("start"); + const onFinish = () => seen.push("finish"); + created.subscribe(onCreated); + requestStart.subscribe(onStart); + finish.subscribe(onFinish); + try { + await using server = http2.createSecureServer( + { allowHTTP1: true, key: tlsCert.key, cert: tlsCert.cert }, + (_req, res) => res.end("ok"), + ); + const port = await listen(server as unknown as http.Server); + // node:https client negotiates http/1.1 (no h2 ALPN) → HTTP/1 fallback. + await new Promise((resolve, reject) => { + const req = https.request( + { port, host: "127.0.0.1", rejectUnauthorized: false, ALPNProtocols: ["http/1.1"] }, + res => { + res.resume(); + res.on("end", resolve); + res.on("error", reject); + }, + ); + req.on("error", reject); + req.end(); + }); + await drain(); + expect(seen.sort()).toEqual(["created", "finish", "start"]); + } finally { + created.unsubscribe(onCreated); + requestStart.unsubscribe(onStart); + finish.unsubscribe(onFinish); + } + }); + test("server works normally when nobody subscribed", async () => { // No subscribers: prove the channel plumbing doesn't break the hot path. await using server = http.createServer((_req, res) => res.end("ok")); From 1b88cf2358f290cc59e814ef9e94fbd16e5abb10 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Thu, 18 Jun 2026 09:52:16 +0000 Subject: [PATCH 07/10] http2: hoist allowHTTP1 diagnostics channels to module scope Drop the redundant require("node:diagnostics_channel") inside connectionListenerHTTP1 (it shadowed the module-level dc) and resolve the two http.server channels once at module load alongside the existing http2.* channel constants, instead of re-resolving them on every HTTP/1 connection. --- src/js/node/http2.ts | 24 +++++++++++------------- 1 file changed, 11 insertions(+), 13 deletions(-) diff --git a/src/js/node/http2.ts b/src/js/node/http2.ts index 17726c14a265..27eba4e5de9f 100644 --- a/src/js/node/http2.ts +++ b/src/js/node/http2.ts @@ -60,6 +60,13 @@ const onServerStreamStartChannel = dc.channel("http2.server.stream.start"); const onServerStreamErrorChannel = dc.channel("http2.server.stream.error"); const onServerStreamFinishChannel = dc.channel("http2.server.stream.finish"); const onServerStreamCloseChannel = dc.channel("http2.server.stream.close"); +// node:_http_server's HTTP server channels, for the allowHTTP1 fallback path. +// response.created is published by the ServerResponse constructor; the other +// two are published in connectionListenerHTTP1 so subscribers see the same +// three events Node fires on this path. Same channel objects as +// node:_http_server (keyed by name in diagnostics_channel's registry). +const onHttp1RequestStartChannel = dc.channel("http.server.request.start"); +const onHttp1ResponseFinishChannel = dc.channel("http.server.response.finish"); const { Readable } = Stream; type Http2ConnectOptions = { settings?: Settings; @@ -5469,15 +5476,6 @@ function connectionListenerHTTP1(server, socket, options) { const ServerResponseClass = http1Options.ServerResponse || http.ServerResponse; const keepAliveTimeout = typeof server.keepAliveTimeout === "number" ? server.keepAliveTimeout : 5000; - // http.server.request.start / http.server.response.finish for the HTTP/1 - // fallback path (allowHTTP1). response.created is published by the - // ServerResponse constructor; publish the other two here so subscribers see - // the same three events Node fires on this path. Same channel objects as - // node:_http_server (keyed by name in diagnostics_channel's registry). - const dc = require("node:diagnostics_channel"); - const onRequestStartChannel = dc.channel("http.server.request.start"); - const onResponseFinishChannel = dc.channel("http.server.response.finish"); - const connections = (server[kHttp1Connections] ??= new SafeSet()); connections.add(socket); socket[kHttp1ActiveRequests] = 0; @@ -5549,12 +5547,12 @@ function connectionListenerHTTP1(server, socket, options) { // Attached unconditionally to match Node's resOnFinish; the hasSubscribers // check happens inside. res.on("finish", () => { - if (onResponseFinishChannel.hasSubscribers) { - onResponseFinishChannel.publish({ request, response: res, socket, server }); + if (onHttp1ResponseFinishChannel.hasSubscribers) { + onHttp1ResponseFinishChannel.publish({ request, response: res, socket, server }); } }); - if (onRequestStartChannel.hasSubscribers) { - onRequestStartChannel.publish({ request, response: res, socket, server }); + if (onHttp1RequestStartChannel.hasSubscribers) { + onHttp1RequestStartChannel.publish({ request, response: res, socket, server }); } server.emit("request", req, res); return 0; From 2c416b6e05c0e9751fd2ff62f41acdaa9e8f1445 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Wed, 15 Jul 2026 04:20:29 +0000 Subject: [PATCH 08/10] http: fold response.finish publish into the existing resOnFinish listener Avoids the second per-request 'finish' listener and its .bind({...}) allocation by publishing http.server.response.finish from the same listener that already closes the socket when needed, like Node's resOnFinish. hasSubscribers is still checked at call time so late subscribers are observed. --- src/js/node/_http_server.ts | 33 ++++++++++----------------------- 1 file changed, 10 insertions(+), 23 deletions(-) diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index e91fff39facc..2afb60447f65 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -113,19 +113,6 @@ const onRequestStartChannel = dc.channel("http.server.request.start"); const onResponseCreatedChannel = dc.channel("http.server.response.created"); const onResponseFinishChannel = dc.channel("http.server.response.finish"); -function emitResponseFinishChannel(this: { req; res; socket; server }) { - // Re-checked here (not at listener-attach) so subscribers that attach - // between request-arrival and response-finish are still observed. - if (onResponseFinishChannel.hasSubscribers) { - onResponseFinishChannel.publish({ - request: this.req, - response: this.res, - socket: this.socket, - server: this.server, - }); - } -} - function emitCloseServer(self: Server) { callCloseCallback(self); self.emit("close"); @@ -681,15 +668,7 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort if (!requestShouldKeepAlive(http_req)) { http_res[kMustCloseConnection] = true; } - - // Attached unconditionally to match Node's resOnFinish; the - // hasSubscribers check happens in emitResponseFinishChannel. Registered - // before endSocketOnFinishIfNeeded so the publish observes the socket - // before its writable side is ended (as Node's resOnFinish publishes - // before socket.destroySoon()). - http_res.on("finish", emitResponseFinishChannel.bind({ req: http_req, res: http_res, socket, server })); - - http_res.once("finish", endSocketOnFinishIfNeeded.bind(undefined, socket, http_res)); + http_res.once("finish", resOnFinish.bind(undefined, http_req, http_res, socket, server)); if (hasObserver("http")) { startPerf(http_res, kServerResponseStatistics, { @@ -1723,7 +1702,15 @@ function stopServerResponsePerf(this: any) { } } -function endSocketOnFinishIfNeeded(socket, res) { +function resOnFinish(req, res, socket, server) { + if (onResponseFinishChannel.hasSubscribers) { + onResponseFinishChannel.publish({ + request: req, + response: res, + socket, + server, + }); + } if (res[kMustCloseConnection]) { socket?.end(); } From 584c6155950a96787d73dcc7f1610f3382aef905 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Wed, 15 Jul 2026 04:41:33 +0000 Subject: [PATCH 09/10] http2: publish response.finish before socket.end() on the allowHTTP1 fallback Move the publish from a res 'finish' listener into handle.onfinished, before its socket.end() call, so subscribers observe the socket before its writable side is ended on non-keep-alive requests - the same ordering resOnFinish uses in node:_http_server. Also drops the per request closure. The allowHTTP1 test now sends Connection: close and asserts socket.writableEnded === false at publish time. Refresh a stale test comment that referenced the removed endSocketOnFinishIfNeeded two-listener mechanism. --- src/js/node/http2.ts | 12 +++++----- .../diagnostics_channel.test.ts | 22 ++++++++++++++----- 2 files changed, 22 insertions(+), 12 deletions(-) diff --git a/src/js/node/http2.ts b/src/js/node/http2.ts index 5a1cb7a55ab3..cddf6345b0e5 100644 --- a/src/js/node/http2.ts +++ b/src/js/node/http2.ts @@ -5575,6 +5575,11 @@ function connectionListenerHTTP1(server, socket, options) { const handle = createHttp1FallbackResponseHandle(socket, shouldKeepAlive, keepAliveTimeout); handle.onfinished = function () { socket[kHttp1ActiveRequests] = Math.max(0, (socket[kHttp1ActiveRequests] || 1) - 1); + // Publish before socket.end() so subscribers observe the socket before + // its writable side is ended, like resOnFinish in node:_http_server. + if (onHttp1ResponseFinishChannel.hasSubscribers) { + onHttp1ResponseFinishChannel.publish({ request, response: res, socket, server }); + } if (!shouldKeepAlive && !socket.destroyed) { socket.end(); } @@ -5582,13 +5587,6 @@ function connectionListenerHTTP1(server, socket, options) { res[kHttp1ResponseHandle] = handle; res.assignSocket(socket); - // Attached unconditionally to match Node's resOnFinish; the hasSubscribers - // check happens inside. - res.on("finish", () => { - if (onHttp1ResponseFinishChannel.hasSubscribers) { - onHttp1ResponseFinishChannel.publish({ request, response: res, socket, server }); - } - }); if (onHttp1RequestStartChannel.hasSubscribers) { onHttp1RequestStartChannel.publish({ request, response: res, socket, server }); } diff --git a/test/js/node/diagnostics_channel/diagnostics_channel.test.ts b/test/js/node/diagnostics_channel/diagnostics_channel.test.ts index 1d46180bb29d..bdf0ea01a02c 100644 --- a/test/js/node/diagnostics_channel/diagnostics_channel.test.ts +++ b/test/js/node/diagnostics_channel/diagnostics_channel.test.ts @@ -618,8 +618,7 @@ describe("http server channels (#29586)", () => { // Node's resOnFinish publishes before it ends the socket, so on the // connection-close path a subscriber still sees a not-yet-ended socket. - // The finish listener is registered before endSocketOnFinishIfNeeded to - // preserve that ordering. + // resOnFinish publishes to the channel before calling socket.end(). test("response.finish publishes before the socket is ended (Connection: close)", async () => { const finish = channel("http.server.response.finish"); let writableEndedAtPublish: boolean | undefined; @@ -650,15 +649,21 @@ describe("http server channels (#29586)", () => { // Node routes http2 allowHTTP1 through the full http1 connectionListener, so // all three channels fire on the HTTP/1 fallback. Bun's fallback lives in - // http2.ts and must publish the same three. + // http2.ts and must publish the same three. Connection: close also checks + // that the finish publish observes the socket before it is ended, like + // resOnFinish in node:_http_server. test("all three channels fire on the http2 allowHTTP1 fallback", async () => { const created = channel("http.server.response.created"); const requestStart = channel("http.server.request.start"); const finish = channel("http.server.response.finish"); const seen: string[] = []; + let writableEndedAtPublish: boolean | undefined; const onCreated = () => seen.push("created"); const onStart = () => seen.push("start"); - const onFinish = () => seen.push("finish"); + const onFinish = (msg: any) => { + seen.push("finish"); + writableEndedAtPublish = msg.socket.writableEnded; + }; created.subscribe(onCreated); requestStart.subscribe(onStart); finish.subscribe(onFinish); @@ -671,7 +676,13 @@ describe("http server channels (#29586)", () => { // node:https client negotiates http/1.1 (no h2 ALPN) → HTTP/1 fallback. await new Promise((resolve, reject) => { const req = https.request( - { port, host: "127.0.0.1", rejectUnauthorized: false, ALPNProtocols: ["http/1.1"] }, + { + port, + host: "127.0.0.1", + rejectUnauthorized: false, + ALPNProtocols: ["http/1.1"], + headers: { connection: "close" }, + }, res => { res.resume(); res.on("end", resolve); @@ -683,6 +694,7 @@ describe("http server channels (#29586)", () => { }); await drain(); expect(seen.sort()).toEqual(["created", "finish", "start"]); + expect(writableEndedAtPublish).toBe(false); } finally { created.unsubscribe(onCreated); requestStart.unsubscribe(onStart); From c11cc98803de9e85ef3e972756338dbd88b41383 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Wed, 15 Jul 2026 05:00:54 +0000 Subject: [PATCH 10/10] http2: publish allowHTTP1 response.finish from the 'finish' event Publishing from handle.onfinished ran synchronously inside handle.end(), before ServerResponse.prototype.end set finished = true, so subscribers observed response.finished === false where Node (and node:_http_server's resOnFinish) show true. Publish from a res 'finish' listener instead and defer the socket teardown into the same listener after the publish - both the non-keep-alive socket.end() and the close-delimited one, which handle.end() now records on handle.closeDelimited instead of ending the socket itself. Subscribers now observe finished === true and writableEnded === false, the same publish-time state as the plain http path. Both ordering tests now assert { finished, writableEnded } at publish time; the allowHTTP1 one fails against the previous placement. --- src/js/node/http2.ts | 25 +++++++++++-------- .../diagnostics_channel.test.ts | 15 ++++++----- 2 files changed, 24 insertions(+), 16 deletions(-) diff --git a/src/js/node/http2.ts b/src/js/node/http2.ts index cddf6345b0e5..ad29e0454155 100644 --- a/src/js/node/http2.ts +++ b/src/js/node/http2.ts @@ -5443,6 +5443,7 @@ function createHttp1FallbackResponseHandle(socket, shouldKeepAlive, keepAliveTim ended: false, finished: false, aborted: false, + closeDelimited: false, bufferedAmount: 0, shouldKeepAlive, onfinished: null, @@ -5480,15 +5481,15 @@ function createHttp1FallbackResponseHandle(socket, shouldKeepAlive, keepAliveTim if (chunked && !noBody) socket.write("0\r\n\r\n"); this.ended = true; this.finished = true; + // A close-delimited body ends at EOF, so the response must end the + // connection; the 'finish' listener in connectionListenerHTTP1 does + // that after the diagnostics publish. + this.closeDelimited = closeDelimited; const onfinished = this.onfinished; if (onfinished) { this.onfinished = null; onfinished(); } - // A close-delimited body ends at EOF, so the response ends the connection. - if (closeDelimited && !socket.destroyed) { - socket.end(); - } return length; }, abort() { @@ -5575,17 +5576,21 @@ function connectionListenerHTTP1(server, socket, options) { const handle = createHttp1FallbackResponseHandle(socket, shouldKeepAlive, keepAliveTimeout); handle.onfinished = function () { socket[kHttp1ActiveRequests] = Math.max(0, (socket[kHttp1ActiveRequests] || 1) - 1); - // Publish before socket.end() so subscribers observe the socket before - // its writable side is ended, like resOnFinish in node:_http_server. + }; + res[kHttp1ResponseHandle] = handle; + res.assignSocket(socket); + + // Like resOnFinish in node:_http_server: publish from the 'finish' event + // (so subscribers observe the finished response) and only then end the + // socket on non-keep-alive / close-delimited responses. + res.once("finish", () => { if (onHttp1ResponseFinishChannel.hasSubscribers) { onHttp1ResponseFinishChannel.publish({ request, response: res, socket, server }); } - if (!shouldKeepAlive && !socket.destroyed) { + if ((!shouldKeepAlive || handle.closeDelimited) && !socket.destroyed) { socket.end(); } - }; - res[kHttp1ResponseHandle] = handle; - res.assignSocket(socket); + }); if (onHttp1RequestStartChannel.hasSubscribers) { onHttp1RequestStartChannel.publish({ request, response: res, socket, server }); diff --git a/test/js/node/diagnostics_channel/diagnostics_channel.test.ts b/test/js/node/diagnostics_channel/diagnostics_channel.test.ts index bdf0ea01a02c..eab913a239dd 100644 --- a/test/js/node/diagnostics_channel/diagnostics_channel.test.ts +++ b/test/js/node/diagnostics_channel/diagnostics_channel.test.ts @@ -621,9 +621,9 @@ describe("http server channels (#29586)", () => { // resOnFinish publishes to the channel before calling socket.end(). test("response.finish publishes before the socket is ended (Connection: close)", async () => { const finish = channel("http.server.response.finish"); - let writableEndedAtPublish: boolean | undefined; + let stateAtPublish: { finished: boolean; writableEnded: boolean } | undefined; const onFinish = (msg: any) => { - writableEndedAtPublish = msg.socket.writableEnded; + stateAtPublish = { finished: msg.response.finished, writableEnded: msg.socket.writableEnded }; }; finish.subscribe(onFinish); let client: any; @@ -640,7 +640,9 @@ describe("http server channels (#29586)", () => { client.on("error", reject); }); await drain(); - expect(writableEndedAtPublish).toBe(false); + // Node's resOnFinish runs as a 'finish' listener: the response is + // finished but the socket is not yet ended at publish time. + expect(stateAtPublish).toEqual({ finished: true, writableEnded: false }); } finally { client?.destroy?.(); finish.unsubscribe(onFinish); @@ -657,12 +659,12 @@ describe("http server channels (#29586)", () => { const requestStart = channel("http.server.request.start"); const finish = channel("http.server.response.finish"); const seen: string[] = []; - let writableEndedAtPublish: boolean | undefined; + let stateAtPublish: { finished: boolean; writableEnded: boolean } | undefined; const onCreated = () => seen.push("created"); const onStart = () => seen.push("start"); const onFinish = (msg: any) => { seen.push("finish"); - writableEndedAtPublish = msg.socket.writableEnded; + stateAtPublish = { finished: msg.response.finished, writableEnded: msg.socket.writableEnded }; }; created.subscribe(onCreated); requestStart.subscribe(onStart); @@ -694,7 +696,8 @@ describe("http server channels (#29586)", () => { }); await drain(); expect(seen.sort()).toEqual(["created", "finish", "start"]); - expect(writableEndedAtPublish).toBe(false); + // Same publish-time state as the plain http path / Node's resOnFinish. + expect(stateAtPublish).toEqual({ finished: true, writableEnded: false }); } finally { created.unsubscribe(onCreated); requestStart.unsubscribe(onStart);