Skip to content
Closed
Show file tree
Hide file tree
Changes from 4 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
38 changes: 38 additions & 0 deletions src/js/node/_http_server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -91,6 +91,26 @@

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");
Expand Down Expand Up @@ -660,6 +680,13 @@
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 }),
);

Check warning on line 688 in src/js/node/_http_server.ts

View check run for this annotation

Claude / Claude Code Review

response.finish publishes after socket.end() runs (listener-ordering Node divergence)

The `emitResponseFinishChannel` listener is registered (line 685) *after* `endSocketOnFinishIfNeeded` (line 666), so on the `kMustCloseConnection` path (HTTP/1.0, `Connection: close`, `shouldKeepAlive=false`, close-delimited body, etc.) `socket.end()` runs before the channel publishes — a subscriber sees `payload.socket.writableEnded === true` where Node shows `false` (Node's `resOnFinish` publishes as its first action, before `socket.destroySoon()`). Practical impact is near-zero, but moving th
Comment thread
robobun marked this conversation as resolved.
Outdated

setIsNextIncomingMessageHTTPS(prevIsNextIncomingMessageHTTPS);
handle.onabort = onServerRequestEvent.bind(socket);
// start buffering data if any, the user will need to resume() or .on("data") to read it
Expand Down Expand Up @@ -737,6 +764,11 @@
http_req._dumpAndCloseReadable();
}

// 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 });
}
Comment thread
claude[bot] marked this conversation as resolved.
if (reachedRequestsLimit) {
server.emit("dropRequest", http_req, socket);
http_res.writeHead(503);
Expand Down Expand Up @@ -1482,6 +1514,12 @@
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 });
}

Check warning on line 1522 in src/js/node/_http_server.ts

View check run for this annotation

Claude / Claude Code Review

http2 allowHTTP1 fallback emits response.created but not request.start/response.finish

Moving the `response.created` publish into the `ServerResponse` constructor is correct per Node, but it now also fires from `connectionListenerHTTP1` in `src/js/node/http2.ts:5526` (the `Http2SecureServer({allowHTTP1: true})` HTTP/1 fallback). That path doesn't go through `onNodeHTTPRequest`, so it never publishes `http.server.request.start` or attaches the `response.finish` listener — an APM subscriber pairing created↔finish will now see orphaned `response.created` events on this path (Node rou
Comment thread
robobun marked this conversation as resolved.
}
$toClass(ServerResponse, "ServerResponse", OutgoingMessage);

Expand Down
277 changes: 277 additions & 0 deletions test/js/node/diagnostics_channel/diagnostics_channel.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -345,6 +347,281 @@ describe("TracingChannel", () => {
test.todo("TODO");
});

describe("http server channels (#29586)", () => {
const listen = (s: http.Server) => new Promise<number>(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<void>(r => setImmediate(r));
await new Promise<void>(r => setImmediate(r));
};

test("publishes http.server.request.start, response.created, response.finish", async () => {
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"));
const port = await listen(server);
await (await fetch(`http://127.0.0.1:${port}/`)).text();
await drain();

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 {
channel("http.server.response.created").unsubscribe(onCreated);
channel("http.server.request.start").unsubscribe(onStart);
channel("http.server.response.finish").unsubscribe(onFinish);
}
Comment thread
claude[bot] marked this conversation as resolved.
});
Comment thread
coderabbitai[bot] marked this conversation as resolved.

// 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");
});
const port = await listen(server);
await (await fetch(`http://127.0.0.1:${port}/`)).text();
await drain();

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<void>();
const { promise: mayFinish, resolve: continueFinish } = Promise.withResolvers<void>();

try {
await using server = http.createServer(async (_req, res) => {
onHandlerEntered();
// Hold the response open until we've subscribed.
await mayFinish;
res.end("done");
});
const port = await listen(server);
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 drain();

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");
});
const port = await listen(server);
await new Promise<void>((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 drain();
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"));
const port = await listen(server);
// Raw socket — we want `Expect: weird` to reach the 417 branch.
await new Promise<void>((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 drain();
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();
});
const port = await listen(server);
await new Promise<void>((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 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
// 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);
Comment thread
claude[bot] marked this conversation as resolved.
} finally {
requestStart.unsubscribe(onStart);
}
});

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"));
const port = await listen(server);
const body = await (await fetch(`http://127.0.0.1:${port}/`)).text();
expect(body).toBe("ok");
});
});

const mocks = new Map();

function mustCall<T>(fn: (...args: any[]) => T, expected?: number) {
Expand Down
Loading