From 810cff1bef9d8f1145c9d1cea8585108ac95afa4 Mon Sep 17 00:00:00 2001 From: Alan Agius <17563226+alan-agius4@users.noreply.github.com> Date: Fri, 31 Jul 2026 18:44:44 +0000 Subject: [PATCH] fix(@angular/ssr): settle writeResponseToNodeResponse when client disconnects When a client disconnects while a response is backpressured or streaming, the Node response (`ServerResponse` or `Http2ServerResponse`) is closed or destroyed without emitting a `drain` event. Previously, this caused `writeResponseToNodeResponse()` to park indefinitely waiting for `drain` and never settle. This change monitors whether the Node response is closed or destroyed and removes event listeners, cancels the reader, and resolves the returned Promise when a client disconnect occurs. Closes #33719 --- packages/angular/ssr/node/src/response.ts | 104 +++++++++++++++--- .../ssr/node/test/response_http1_spec.ts | 44 ++++++++ .../ssr/node/test/response_http2_spec.ts | 44 ++++++++ 3 files changed, 178 insertions(+), 14 deletions(-) diff --git a/packages/angular/ssr/node/src/response.ts b/packages/angular/ssr/node/src/response.ts index 56f091deed5f..90ed5bb31a67 100644 --- a/packages/angular/ssr/node/src/response.ts +++ b/packages/angular/ssr/node/src/response.ts @@ -9,6 +9,24 @@ import type { ServerResponse } from 'node:http'; import type { Http2ServerResponse } from 'node:http2'; +/** + * Checks whether a Node.js `ServerResponse` or `Http2ServerResponse` is destroyed, closed, or ended. + * + * @param destination - The HTTP/1.1 or HTTP/2 server response to check. + * @returns `true` if the response or its underlying stream is destroyed, closed, or ended; otherwise `false`. + */ +function isResponseDestroyedOrClosed(destination: ServerResponse | Http2ServerResponse): boolean { + return ( + Boolean(destination.destroyed) || + Boolean(destination.closed) || + Boolean(destination.writableEnded) || + ('stream' in destination && + (!destination.stream || + Boolean(destination.stream.destroyed) || + Boolean(destination.stream.closed))) + ); +} + /** * Streams a web-standard `Response` into a Node.js `ServerResponse` * or `Http2ServerResponse`. @@ -24,6 +42,10 @@ export async function writeResponseToNodeResponse( source: Response, destination: ServerResponse | Http2ServerResponse, ): Promise { + if (isResponseDestroyedOrClosed(destination)) { + return; + } + const { status, headers, body } = source; destination.statusCode = status; @@ -48,27 +70,52 @@ export async function writeResponseToNodeResponse( } if (!body) { - destination.end(); + if (!isResponseDestroyedOrClosed(destination)) { + destination.end(); + } return; } - try { - const reader = body.getReader(); - - destination.on('close', () => { - reader.cancel().catch((error) => { - // eslint-disable-next-line no-console - console.error( - `An error occurred while writing the response body for: ${destination.req.url}.`, - error, - ); - }); + let isClosed = isResponseDestroyedOrClosed(destination); + const isDestroyedOrClosed = () => isClosed || isResponseDestroyedOrClosed(destination); + + let readerCancelled = false; + const reader = body.getReader(); + const cancelReader = (error?: unknown) => { + if (readerCancelled) { + return; + } + readerCancelled = true; + isClosed = true; + destination.off('close', cancelReader); + destination.off('error', cancelReader); + reader.cancel(error).catch((err) => { + // eslint-disable-next-line no-console + console.error( + `An error occurred while writing the response body for: ${destination.req.url}.`, + err, + ); }); + }; + + destination.once('close', cancelReader); + destination.once('error', cancelReader); + try { // eslint-disable-next-line no-constant-condition while (true) { + if (isDestroyedOrClosed()) { + cancelReader(); + break; + } + const { done, value } = await reader.read(); + if (isDestroyedOrClosed()) { + cancelReader(); + break; + } + if (done) { destination.end(); break; @@ -78,10 +125,39 @@ export async function writeResponseToNodeResponse( if (canContinue === false) { // Explicitly check for `false`, as AWS may return `undefined` even though this is not valid. // See: https://github.com/CodeGenieApp/serverless-express/issues/683 - await new Promise((resolve) => destination.once('drain', resolve)); + await new Promise((resolve) => { + if (isDestroyedOrClosed()) { + resolve(); + + return; + } + + const onDrain = () => { + destination.off('close', onClose); + destination.off('error', onClose); + resolve(); + }; + + const onClose = () => { + destination.off('drain', onDrain); + destination.off('close', onClose); + destination.off('error', onClose); + cancelReader(); + resolve(); + }; + + destination.once('drain', onDrain); + destination.once('close', onClose); + destination.once('error', onClose); + }); } } } catch { - destination.end('Internal server error.'); + if (!isDestroyedOrClosed()) { + destination.end('Internal server error.'); + } + } finally { + destination.off('close', cancelReader); + destination.off('error', cancelReader); } } diff --git a/packages/angular/ssr/node/test/response_http1_spec.ts b/packages/angular/ssr/node/test/response_http1_spec.ts index 8deae2b7e3b4..c7648f03a627 100644 --- a/packages/angular/ssr/node/test/response_http1_spec.ts +++ b/packages/angular/ssr/node/test/response_http1_spec.ts @@ -107,4 +107,48 @@ describe('writeResponseToNodeResponse (HTTP/1.1)', () => { expect(response.headers['set-cookie']).toEqual(cookieValue); }); + + it('should resolve and cancel reader when client disconnects while response is backpressured', async () => { + let writePromise!: Promise; + let readerCancelled = false; + const largeChunk = 'x'.repeat(1024 * 1024 * 4); // 4MB to ensure backpressure + const stream = new ReadableStream({ + start(controller) { + controller.enqueue(largeChunk); + controller.enqueue(largeChunk); + }, + cancel() { + readerCancelled = true; + }, + }); + + server.once('request', (_, nodeResponse) => { + writePromise = writeResponseToNodeResponse(new Response(stream), nodeResponse); + }); + + await new Promise((resolve) => { + const { port } = server.address() as AddressInfo; + const clientRequest = requestCb( + { + host: 'localhost', + port, + }, + (response) => { + response.once('data', () => { + clientRequest.destroy(); + resolve(); + }); + }, + ); + + clientRequest.on('error', () => { + // Expected when destroying the socket + }); + + clientRequest.end(); + }); + + await writePromise; + expect(readerCancelled).toBeTrue(); + }); }); diff --git a/packages/angular/ssr/node/test/response_http2_spec.ts b/packages/angular/ssr/node/test/response_http2_spec.ts index c2de8e977561..431bc9de0c38 100644 --- a/packages/angular/ssr/node/test/response_http2_spec.ts +++ b/packages/angular/ssr/node/test/response_http2_spec.ts @@ -116,4 +116,48 @@ describe('writeResponseToNodeResponse (HTTP/2)', () => { expect(resHeaders['set-cookie']).toEqual(cookieValue); }); + + it('should resolve and cancel reader when client disconnects while response is backpressured', async () => { + let writePromise!: Promise; + let readerCancelled = false; + const largeChunk = 'x'.repeat(1024 * 1024); // 1MB to exceed HTTP/2 flow control window + const stream = new ReadableStream({ + start(controller) { + controller.enqueue(largeChunk); + controller.enqueue(largeChunk); + }, + cancel() { + readerCancelled = true; + }, + }); + + server.once('request', (_, nodeResponse) => { + writePromise = writeResponseToNodeResponse(new Response(stream), nodeResponse); + }); + + await new Promise((resolve) => { + const { port } = server.address() as AddressInfo; + const client = connect(`http://localhost:${port}`); + const req = client.request({ + ':path': '/', + }); + + req.once('response', () => { + req.once('data', () => { + req.destroy(); + client.destroy(); + resolve(); + }); + }); + + req.on('error', () => { + // Expected when destroying/closing the stream + }); + + req.end(); + }); + + await writePromise; + expect(readerCancelled).toBeTrue(); + }); });