catalog-backend: fix for not waiting for response buffer to drain as needed

Signed-off-by: Patrik Oldsberg <poldsberg@gmail.com>
This commit is contained in:
Patrik Oldsberg
2024-12-13 16:27:24 +01:00
parent 1d1d7ad658
commit e0cd1894f4
3 changed files with 41 additions and 58 deletions
@@ -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);
@@ -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<boolean>;
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) {
@@ -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<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);
});
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<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;
}