From 1752be6b4b07c6fe9b243bcb9b2c997d6d19f4f1 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Fredrik=20Adel=C3=B6w?= Date: Wed, 6 Aug 2025 10:55:32 +0200 Subject: [PATCH] handle event emitter leak when streaming responses MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Fredrik Adelöw --- .changeset/deep-mangos-dig.md | 5 ++ .../src/service/createRouter.ts | 4 +- .../response/createEntityArrayJsonStream.ts | 14 ++-- .../src/service/response/write.ts | 73 ++++++++++++------- 4 files changed, 62 insertions(+), 34 deletions(-) create mode 100644 .changeset/deep-mangos-dig.md diff --git a/.changeset/deep-mangos-dig.md b/.changeset/deep-mangos-dig.md new file mode 100644 index 0000000000..237f625e27 --- /dev/null +++ b/.changeset/deep-mangos-dig.md @@ -0,0 +1,5 @@ +--- +'@backstage/plugin-catalog-backend': patch +--- + +Attempt to circumvent event listener memory leak in compression middleware diff --git a/plugins/catalog-backend/src/service/createRouter.ts b/plugins/catalog-backend/src/service/createRouter.ts index 2c9e63e8c4..9c9b97eb4c 100644 --- a/plugins/catalog-backend/src/service/createRouter.ts +++ b/plugins/catalog-backend/src/service/createRouter.ts @@ -214,7 +214,7 @@ export async function createRouter( let cursor: Cursor | undefined; try { - let currentWrite: Promise | undefined = undefined; + let currentWrite: Promise<'ok' | 'closed'> | undefined = undefined; do { const result = await entitiesCatalog.queryEntities( !cursor @@ -230,7 +230,7 @@ export async function createRouter( ); // Wait for previous write to complete - if (await currentWrite) { + if ((await currentWrite) === 'closed') { return; // Client closed connection } diff --git a/plugins/catalog-backend/src/service/response/createEntityArrayJsonStream.ts b/plugins/catalog-backend/src/service/response/createEntityArrayJsonStream.ts index 5bdc9d8994..ac557f5a2f 100644 --- a/plugins/catalog-backend/src/service/response/createEntityArrayJsonStream.ts +++ b/plugins/catalog-backend/src/service/response/createEntityArrayJsonStream.ts @@ -16,10 +16,10 @@ import { EntitiesResponseItems } from '../../catalog/types'; import { Response } from 'express'; -import { writeResponseData } from './write'; +import { createResponseDataWriter } from './write'; export interface EntityArrayJsonStream { - send(entities: EntitiesResponseItems): Promise; + send(entities: EntitiesResponseItems): Promise<'ok' | 'closed'>; complete(): void; close(): void; } @@ -28,6 +28,8 @@ export interface EntityArrayJsonStream { export function createEntityArrayJsonStream( res: Response, ): EntityArrayJsonStream { + const writer = createResponseDataWriter(res); + // Imitate the httpRouter behavior of pretty-printing in development const prettyPrint = process.env.NODE_ENV === 'development'; let firstSend = true; @@ -46,11 +48,11 @@ export function createEntityArrayJsonStream( const prefix = firstSend ? '[' : ','; firstSend = false; - if (await writeResponseData(res, prefix + item)) { - return true; + if ((await writer(prefix + item)) === 'closed') { + return 'closed'; } } - return false; + return 'ok'; } let data: string; @@ -63,7 +65,7 @@ export function createEntityArrayJsonStream( } firstSend = false; - return writeResponseData(res, data); + return writer(data); }, complete() { if (firstSend) { diff --git a/plugins/catalog-backend/src/service/response/write.ts b/plugins/catalog-backend/src/service/response/write.ts index b43064a1e4..8fbd54d01c 100644 --- a/plugins/catalog-backend/src/service/response/write.ts +++ b/plugins/catalog-backend/src/service/response/write.ts @@ -16,7 +16,7 @@ import { Response } from 'express'; import { EntitiesResponseItems } from '../../catalog/types'; -import { JsonValue } from '@backstage/types'; +import { createDeferred, DeferredPromise, JsonValue } from '@backstage/types'; import { NotFoundError } from '@backstage/errors'; import { processEntitiesResponseItems } from './process'; @@ -50,6 +50,8 @@ export async function writeEntitiesResponse(options: { alwaysUseObjectMode?: boolean; }) { const { res, responseWrapper, alwaysUseObjectMode } = options; + const writer = createResponseDataWriter(res); + const items = alwaysUseObjectMode ? processEntitiesResponseItems(options.items, e => e) : options.items; @@ -83,7 +85,7 @@ export async function writeEntitiesResponse(options: { const prefix = first ? '[' : ','; first = false; - if (await writeResponseData(res, prefix + entity)) { + if ((await writer(prefix + entity)) === 'closed') { return; } } @@ -91,32 +93,51 @@ export async function writeEntitiesResponse(options: { } /** - * Writes a data to the response and waits if the response buffer needs draining. + * Creates a data writer that writes to the response and waits if the response + * buffer needs draining. * * @internal - * @returns true if the response was closed while waiting for the buffer to drain + * @returns A writer function. If a write attempt returns 'closed', the + * connection has become closed prematurely and the caller should stop trying to + * write. */ -export async function writeResponseData(res: Response, data: string | Buffer) { - const ok = res.write(data, 'utf8'); - if (!ok) { - if (res.closed) { - return true; +export function createResponseDataWriter( + res: Response, +): (data: string | Buffer) => Promise<'ok' | 'closed'> { + // See https://github.com/backstage/backstage/issues/30659 + // + // This code goes to some lengths to just add listeners once at the top, + // instead of on every need to drain. Hence it is more complex that seems to + // be necessary, just to avoid listener leaks. + + let drainPromise: DeferredPromise<'ok'> | undefined; + + const closePromise = new Promise<'closed'>(resolve => { + function onClose() { + res.off('drain', onDrain); + res.off('close', onClose); + res.off('finish', onClose); + resolve('closed'); } - const closed = await new Promise(resolve => { - function onContinue() { - res.off('drain', onContinue); - res.off('close', onClose); - resolve(false); - } - function onClose() { - res.off('drain', onContinue); - res.off('close', onClose); - resolve(true); - } - res.on('drain', onContinue); - res.on('close', onClose); - }); - return closed; - } - return false; + function onDrain() { + drainPromise?.resolve('ok'); + drainPromise = undefined; + } + res.on('drain', onDrain); + res.on('close', onClose); + res.on('finish', onClose); + }); + + return async data => { + if (res.write(data, 'utf8')) { + return 'ok'; + } + + if (res.closed) { + return 'closed'; + } + + drainPromise = createDeferred(); + return Promise.race([drainPromise, closePromise]); + }; }