handle event emitter leak when streaming responses
Signed-off-by: Fredrik Adelöw <freben@gmail.com>
This commit is contained in:
@@ -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]);
|
||||
};
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user