handle event emitter leak when streaming responses

Signed-off-by: Fredrik Adelöw <freben@gmail.com>
This commit is contained in:
Fredrik Adelöw
2025-08-06 10:55:32 +02:00
parent cdc0fa00d0
commit 1752be6b4b
4 changed files with 62 additions and 34 deletions
+5
View File
@@ -0,0 +1,5 @@
---
'@backstage/plugin-catalog-backend': patch
---
Attempt to circumvent event listener memory leak in compression middleware
@@ -214,7 +214,7 @@ export async function createRouter(
let cursor: Cursor | undefined;
try {
let currentWrite: Promise<boolean> | 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
}
@@ -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<boolean>;
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) {
@@ -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<boolean>(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]);
};
}