diff --git a/plugins/catalog-backend/src/service/createRouter.ts b/plugins/catalog-backend/src/service/createRouter.ts index 601b0b1098..50ccb08b4c 100644 --- a/plugins/catalog-backend/src/service/createRouter.ts +++ b/plugins/catalog-backend/src/service/createRouter.ts @@ -58,7 +58,6 @@ import { } from '@backstage/backend-plugin-api'; import { LocationAnalyzer } from '@backstage/plugin-catalog-node'; import { AuthorizedValidationService } from './AuthorizedValidationService'; -import { DeferredPromise, createDeferred } from '@backstage/types'; import { createEntityArrayJsonStream, processEntitiesResponseItems, @@ -175,22 +174,6 @@ export async function createRouter( return; } - // For other read-the-entire-world cases, use queryEntities and stream - // out results. - - // The write lock is used for back pressure, preventing slow readers - // from forcing our read loop to pile up response data in userspace - // buffers faster than the kernel buffer is emptied. - // https://nodejs.org/api/http.html#http_response_write_chunk_encoding_callback - const locks: { writeLock?: DeferredPromise } = {}; - const controller = new AbortController(); - const signal = controller.signal; - req.on('end', () => { - controller.abort(new Error('Client closed connection')); - locks.writeLock?.resolve(); - delete locks.writeLock; - }); - const responseStream = createEntityArrayJsonStream(res); const limit = 10000; let cursor: Cursor | undefined; @@ -211,31 +194,17 @@ export async function createRouter( ); if (result.items.entities.length) { - await locks?.writeLock; - - signal.throwIfAborted(); - if (!disableRelationsCompatibility) { result.items = processEntitiesResponseItems( result.items, expandLegacyCompoundRelationsInEntity, ); } - if (!responseStream.send(result.items)) { - // The kernel buffer is full. Create the lock but do not await it - // yet - we can better spend our time going to the next round of - // the loop and read from the database while we wait for it to - // drain. - locks.writeLock = createDeferred(); - res.once('drain', () => { - locks.writeLock?.resolve(); - delete locks.writeLock; - }); + if (await responseStream.send(result.items)) { + return; // Client closed connection } } - signal.throwIfAborted(); - cursor = result.pageInfo?.nextCursor; } while (cursor); diff --git a/plugins/catalog-backend/src/service/response/createEntityArrayJsonStream.ts b/plugins/catalog-backend/src/service/response/createEntityArrayJsonStream.ts index 4762afd2d9..5bdc9d8994 100644 --- a/plugins/catalog-backend/src/service/response/createEntityArrayJsonStream.ts +++ b/plugins/catalog-backend/src/service/response/createEntityArrayJsonStream.ts @@ -16,9 +16,10 @@ import { EntitiesResponseItems } from '../../catalog/types'; import { Response } from 'express'; +import { writeResponseData } from './write'; export interface EntityArrayJsonStream { - send(entities: EntitiesResponseItems): boolean; + send(entities: EntitiesResponseItems): Promise; complete(): void; close(): void; } @@ -33,7 +34,7 @@ export function createEntityArrayJsonStream( let completed = false; return { - send(response) { + async send(response) { if (firstSend) { res.setHeader('Content-Type', 'application/json; charset=utf-8'); res.status(200); @@ -41,13 +42,15 @@ export function createEntityArrayJsonStream( } if (response.type === 'raw') { - let needsDrain = false; for (const item of response.entities) { const prefix = firstSend ? '[' : ','; firstSend = false; - needsDrain ||= !res.write(prefix + item, 'utf8'); + + if (await writeResponseData(res, prefix + item)) { + return true; + } } - return !needsDrain; + return false; } let data: string; @@ -60,7 +63,7 @@ export function createEntityArrayJsonStream( } firstSend = false; - return res.write(data, 'utf8'); + return writeResponseData(res, 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 c2bca0953a..e76369fbc6 100644 --- a/plugins/catalog-backend/src/service/response/write.ts +++ b/plugins/catalog-backend/src/service/response/write.ts @@ -78,26 +78,37 @@ export async function writeEntitiesResponse( const prefix = first ? '[' : ','; first = false; - const needsDrain = !res.write(prefix + entity, 'utf8'); - if (needsDrain) { - 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); - }); - if (closed) { - return; - } + if (await writeResponseData(res, prefix + entity)) { + return; } } res.end(`${first ? '[' : ''}]${trailing}`); } + +/** + * Writes a data 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 + */ +export async function writeResponseData(res: Response, data: string | Buffer) { + const ok = res.write(data, 'utf8'); + if (!ok) { + 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; +}