feat(catalog): batch size + mode config for entity streaming

adds parameters to control the batch size and mode for the
streamEntities method. with array mode, the entities are returned as
batches while the default single mode returns entities one by one.

Signed-off-by: Hellgren Heikki <heikki.hellgren@op.fi>
This commit is contained in:
Hellgren Heikki
2025-08-29 08:30:24 +03:00
parent 0e9ec444b7
commit 2f5162c4b7
11 changed files with 86 additions and 32 deletions
@@ -71,7 +71,14 @@ export class InMemoryCatalogClient implements CatalogApi {
// (undocumented)
removeLocationById(_id: string): Promise<void>;
// (undocumented)
streamEntities(request?: StreamEntitiesRequest): AsyncIterable<Entity>;
streamEntities<
T extends StreamEntitiesRequest,
R = T extends {
mode: 'array';
} & StreamEntitiesRequest
? Entity[]
: Entity,
>(request?: T): AsyncIterable<R>;
// (undocumented)
validateEntity(
_entity: Entity,
+13 -6
View File
@@ -91,7 +91,7 @@ export interface CatalogApi {
streamEntities(
request?: StreamEntitiesRequest,
options?: CatalogRequestOptions,
): AsyncIterable<Entity>;
): AsyncIterable<Entity | Entity[]>;
validateEntity(
entity: Entity,
locationRef: string,
@@ -174,10 +174,14 @@ export class CatalogClient implements CatalogApi {
id: string,
options?: CatalogRequestOptions,
): Promise<void>;
streamEntities(
request?: StreamEntitiesRequest,
options?: CatalogRequestOptions,
): AsyncIterable<Entity>;
streamEntities<
T extends StreamEntitiesRequest,
R = T extends {
mode: 'array';
} & StreamEntitiesRequest
? Entity[]
: Entity,
>(request?: T, options?: CatalogRequestOptions): AsyncIterable<R>;
validateEntity(
entity: Entity,
locationRef: string,
@@ -329,7 +333,10 @@ export type QueryEntitiesResponse = {
export type StreamEntitiesRequest = Omit<
QueryEntitiesRequest,
'limit' | 'offset'
>;
> & {
batchSize?: number;
mode?: 'single' | 'array';
};
// @public
export type ValidateEntityResponse =
@@ -579,7 +579,7 @@ describe('CatalogClient', () => {
);
});
it('should stream entities', async () => {
it('should stream entities one by one', async () => {
const stream = client.streamEntities({}, { token });
const results: Entity[] = [];
for await (const entity of stream) {
@@ -588,6 +588,15 @@ describe('CatalogClient', () => {
expect(results).toEqual(defaultResponse.items);
});
it('should stream entities in batches', async () => {
const stream = client.streamEntities({ mode: 'array' }, { token });
const results: Entity[] = [];
for await (const entityBatch of stream) {
results.push(...entityBatch);
}
expect(results).toEqual(defaultResponse.items);
});
it('should handle errors', async () => {
const mockedEndpoint = jest
.fn()
@@ -597,7 +606,7 @@ describe('CatalogClient', () => {
const stream = client.streamEntities({}, { token });
await expect(async () => {
const results: Entity[] = [];
const results = [];
for await (const entity of stream) {
results.push(entity);
}
+15 -10
View File
@@ -51,7 +51,7 @@ import type {
} from '@backstage/plugin-catalog-common';
// Number of entities to return in a single streamEntities request
const STREAM_ENTITIES_LIMIT = 500;
const DEFAULT_STREAM_ENTITIES_LIMIT = 500;
/**
* A frontend and backend compatible client for communicating with the Backstage
@@ -463,20 +463,25 @@ export class CatalogClient implements CatalogApi {
/**
* {@inheritdoc CatalogApi.streamEntities}
*/
async *streamEntities(
request?: StreamEntitiesRequest,
options?: CatalogRequestOptions,
): AsyncIterable<Entity> {
async *streamEntities<
T extends StreamEntitiesRequest,
R = T extends { mode: 'array' } & StreamEntitiesRequest ? Entity[] : Entity,
>(request?: T, options?: CatalogRequestOptions): AsyncIterable<R> {
let cursor: string | undefined = undefined;
const limit = request?.batchSize ?? DEFAULT_STREAM_ENTITIES_LIMIT;
const mode = request?.mode ?? 'single';
do {
const res = await this.queryEntities(
cursor
? { ...request, cursor, limit: STREAM_ENTITIES_LIMIT }
: { ...request, limit: STREAM_ENTITIES_LIMIT },
cursor ? { ...request, cursor, limit } : { ...request, limit },
options,
);
for (const entity of res.items) {
yield entity;
if (mode === 'single') {
for (const entity of res.items) {
yield entity as R;
}
} else {
yield res.items as R;
}
cursor = res.pageInfo.nextCursor;
@@ -123,7 +123,7 @@ describe('InMemoryCatalogClient', () => {
});
});
it('streamEntities', async () => {
it('streamEntities single', async () => {
const client = new InMemoryCatalogClient({ entities });
const stream = client.streamEntities();
const results: Entity[] = [];
@@ -133,6 +133,16 @@ describe('InMemoryCatalogClient', () => {
expect(results).toEqual(entities);
});
it('streamEntities batch', async () => {
const client = new InMemoryCatalogClient({ entities });
const stream = client.streamEntities({ mode: 'array' });
const results: Entity[] = [];
for await (const entity of stream) {
results.push(...entity);
}
expect(results).toEqual(entities);
});
it('getEntityAncestors', async () => {
const client = new InMemoryCatalogClient({ entities });
await expect(
@@ -280,16 +280,23 @@ export class InMemoryCatalogClient implements CatalogApi {
throw new NotImplementedError('Method not implemented.');
}
async *streamEntities(
request?: StreamEntitiesRequest,
): AsyncIterable<Entity> {
async *streamEntities<
T extends StreamEntitiesRequest,
R = T extends { mode: 'array' } & StreamEntitiesRequest ? Entity[] : Entity,
>(request?: T): AsyncIterable<R> {
let cursor: string | undefined = undefined;
const mode = request?.mode ?? 'single';
do {
const res = await this.queryEntities(
cursor ? { ...request, cursor } : request,
);
for (const entity of res.items) {
yield entity;
if (mode === 'single') {
for (const entity of res.items) {
yield entity as R;
}
} else {
yield res.items as R;
}
cursor = res.pageInfo.nextCursor;
+11 -2
View File
@@ -473,7 +473,16 @@ export type QueryEntitiesResponse = {
export type StreamEntitiesRequest = Omit<
QueryEntitiesRequest,
'limit' | 'offset'
>;
> & {
/**
* The number of entities to fetch in each batch. Defaults to 500.
*/
batchSize?: number;
/**
* The mode in which entities should be yielded. Defaults to 'single'.
*/
mode?: 'single' | 'array';
};
/**
* A client for interacting with the Backstage software catalog through its API.
@@ -715,5 +724,5 @@ export interface CatalogApi {
streamEntities(
request?: StreamEntitiesRequest,
options?: CatalogRequestOptions,
): AsyncIterable<Entity>;
): AsyncIterable<Entity | Entity[]>;
}
+1 -1
View File
@@ -111,7 +111,7 @@ export interface CatalogServiceMock extends CatalogService, CatalogApi {
streamEntities(
request?: StreamEntitiesRequest,
options?: CatalogServiceRequestOptions | CatalogRequestOptions,
): AsyncIterable<Entity>;
): AsyncIterable<Entity | Entity[]>;
// (undocumented)
validateEntity(
entity: Entity,
+1 -1
View File
@@ -199,7 +199,7 @@ export interface CatalogService {
streamEntities(
request: StreamEntitiesRequest | undefined,
options: CatalogServiceRequestOptions,
): AsyncIterable<Entity>;
): AsyncIterable<Entity | Entity[]>;
// (undocumented)
validateEntity(
entity: Entity,
+2 -2
View File
@@ -146,7 +146,7 @@ export interface CatalogService {
streamEntities(
request: StreamEntitiesRequest | undefined,
options: CatalogServiceRequestOptions,
): AsyncIterable<Entity>;
): AsyncIterable<Entity | Entity[]>;
}
class DefaultCatalogService implements CatalogService {
@@ -329,7 +329,7 @@ class DefaultCatalogService implements CatalogService {
async *streamEntities(
request: StreamEntitiesRequest | undefined,
options: CatalogServiceRequestOptions,
): AsyncIterable<Entity> {
): AsyncIterable<Entity | Entity[]> {
yield* this.#catalogApi.streamEntities(
request,
await this.#getOptions(options),
+1 -1
View File
@@ -140,5 +140,5 @@ export interface CatalogServiceMock extends CatalogService, CatalogApi {
streamEntities(
request?: StreamEntitiesRequest,
options?: CatalogServiceRequestOptions | CatalogRequestOptions,
): AsyncIterable<Entity>;
): AsyncIterable<Entity | Entity[]>;
}