From 7e751c5b2bc4b7596e09adcf016bc10ebcfec09d Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Fri, 31 Jul 2026 05:46:22 -0700 Subject: [PATCH] fix(gateway): retain HTTP disconnect and stream error ownership (#116835) Co-authored-by: Peter Steinberger --- src/gateway/http-common.test.ts | 99 ++++++++++++++++++++++++++++++- src/gateway/http-common.ts | 12 ++++ src/gateway/test-http-response.ts | 4 +- 3 files changed, 112 insertions(+), 3 deletions(-) diff --git a/src/gateway/http-common.test.ts b/src/gateway/http-common.test.ts index 94ad10ff9021..5df82e7e926c 100644 --- a/src/gateway/http-common.test.ts +++ b/src/gateway/http-common.test.ts @@ -1,7 +1,7 @@ // HTTP common tests cover JSON/text response helpers, auth failures, security // headers, SSE headers, body parsing, and disconnect diagnostics. import { EventEmitter } from "node:events"; -import type { IncomingMessage, ServerResponse } from "node:http"; +import { ServerResponse, type IncomingMessage } from "node:http"; import { beforeEach, describe, expect, it, vi } from "vitest"; import { onDiagnosticEvent, @@ -337,6 +337,103 @@ describe("watchClientDisconnect", () => { expect(controller.signal.aborted).toBe(true); }); + it("immediately aborts when the request socket was already destroyed", () => { + const socket = Object.assign(new EventEmitter(), { destroyed: true }); + const { req, res } = makeMockHttpReqRes(socket, socket); + const controller = new AbortController(); + const onDisconnect = vi.fn(); + + const cleanup = watchClientDisconnect(req, res, controller, onDisconnect); + + expect(controller.signal.aborted).toBe(true); + expect(onDisconnect).toHaveBeenCalledTimes(1); + expect(socket.listenerCount("close")).toBe(0); + expect(res.listenerCount("error")).toBe(1); + cleanup(); + res.emit("close"); + expect(res.listenerCount("error")).toBe(0); + }); + + it("immediately aborts when the response was already destroyed", () => { + const socket = new EventEmitter(); + const { req, res } = makeMockHttpReqRes(socket, socket); + Object.assign(res, { destroyed: true }); + const controller = new AbortController(); + const onDisconnect = vi.fn(); + + const cleanup = watchClientDisconnect(req, res, controller, onDisconnect); + + expect(controller.signal.aborted).toBe(true); + expect(onDisconnect).toHaveBeenCalledTimes(1); + expect(socket.listenerCount("close")).toBe(0); + expect(res.listenerCount("error")).toBe(1); + cleanup(); + res.emit("close"); + expect(res.listenerCount("error")).toBe(0); + }); + + it("handles response stream errors as client disconnects", () => { + const socket = new EventEmitter(); + const { req, res } = makeMockHttpReqRes(socket, socket); + const controller = new AbortController(); + const onDisconnect = vi.fn(); + const cleanup = watchClientDisconnect(req, res, controller, onDisconnect); + + expect(() => res.emit("error", new Error("response stream failed"))).not.toThrow(); + expect(controller.signal.aborted).toBe(true); + expect(onDisconnect).toHaveBeenCalledTimes(1); + + cleanup(); + res.emit("close"); + expect(res.listenerCount("error")).toBe(0); + }); + + it("keeps real response errors handled after cleanup until the response closes", async () => { + const socket = new EventEmitter(); + const req = { socket } as IncomingMessage; + const res = new ServerResponse({ method: "POST" } as IncomingMessage); + const controller = new AbortController(); + const onDisconnect = vi.fn(); + const cleanup = watchClientDisconnect(req, res, controller, onDisconnect); + + expect(res.listenerCount("error")).toBeGreaterThan(0); + cleanup(); + res.end(); + res.write("late SSE frame"); + await new Promise((resolve) => { + setImmediate(resolve); + }); + + expect(controller.signal.aborted).toBe(true); + expect(onDisconnect).toHaveBeenCalledTimes(1); + + res.emit("close"); + expect(res.listenerCount("error")).toBe(0); + }); + + it("keeps deferred errors handled when a real response was already destroyed", async () => { + const socket = new EventEmitter(); + const req = { socket } as IncomingMessage; + const res = new ServerResponse({ method: "POST" } as IncomingMessage); + const controller = new AbortController(); + const deferredError = new Promise((resolve) => { + process.nextTick(() => { + res.emit("error", new Error("destroyed response failed during cleanup")); + resolve(); + }); + }); + res.destroy(); + + const cleanup = watchClientDisconnect(req, res, controller); + expect(controller.signal.aborted).toBe(true); + expect(res.listenerCount("error")).toBe(1); + + await deferredError; + cleanup(); + res.emit("close"); + expect(res.listenerCount("error")).toBe(0); + }); + it("does not double-abort when the controller is already aborted", () => { const socket = new EventEmitter(); const { req, res } = makeMockHttpReqRes(socket, null); diff --git a/src/gateway/http-common.ts b/src/gateway/http-common.ts index 265bb53a4c51..37cfbe776a58 100644 --- a/src/gateway/http-common.ts +++ b/src/gateway/http-common.ts @@ -182,6 +182,18 @@ export function watchClientDisconnect( abortController.abort(new ClientDisconnectError()); } }; + const stopWatchingResponseErrors = () => { + res.off("error", handleClose); + res.off("close", stopWatchingResponseErrors); + }; + // Finalizers release socket watchers before res.end(); keep its error + // listener until close so a failed flush cannot become process-fatal. + res.on("error", handleClose); + res.once("close", stopWatchingResponseErrors); + if (res.destroyed || sockets.some((socket) => socket.destroyed)) { + handleClose(); + return () => {}; + } for (const socket of sockets) { socket.on("close", handleClose); } diff --git a/src/gateway/test-http-response.ts b/src/gateway/test-http-response.ts index 85585257a0c5..4a25a858d98d 100644 --- a/src/gateway/test-http-response.ts +++ b/src/gateway/test-http-response.ts @@ -1,6 +1,6 @@ // Gateway HTTP test helpers build minimal request/response doubles and collect // client response bodies. -import type { EventEmitter } from "node:events"; +import { EventEmitter } from "node:events"; import type { IncomingMessage, ServerResponse } from "node:http"; import { PassThrough } from "node:stream"; import { vi } from "vitest"; @@ -37,7 +37,7 @@ export function makeMockHttpReqRes( ): { req: IncomingMessage; res: ServerResponse } { return { req: { socket: reqSocket } as unknown as IncomingMessage, - res: { socket: resSocket } as unknown as ServerResponse, + res: Object.assign(new EventEmitter(), { socket: resSocket }) as unknown as ServerResponse, }; }