From acb2fdb38fb2ed80364b0bb1ef837c89365ed38c Mon Sep 17 00:00:00 2001 From: Hellgren Heikki Date: Fri, 29 Aug 2025 14:11:05 +0300 Subject: [PATCH] feat(catalog): add possibility to stream entity pages adds a new method `streamEntityPages` that can be used to stream entity pages from catalog instead streaming entities one by one. this can be more efficient load wise but not usable for every use case. Signed-off-by: Hellgren Heikki --- .changeset/olive-moons-burn.md | 16 +++++++++- .../catalog-client/report-testUtils.api.md | 11 ++----- packages/catalog-client/report.api.md | 23 +++++++------ .../catalog-client/src/CatalogClient.test.ts | 14 ++++---- packages/catalog-client/src/CatalogClient.ts | 32 ++++++++++++------- .../testUtils/InMemoryCatalogClient.test.ts | 14 ++++---- .../src/testUtils/InMemoryCatalogClient.ts | 27 +++++++++------- packages/catalog-client/src/types/api.ts | 20 +++++++++--- plugins/catalog-node/report-testUtils.api.md | 7 +++- plugins/catalog-node/report.api.md | 7 +++- plugins/catalog-node/src/catalogService.ts | 19 +++++++++-- .../src/testUtils/catalogServiceMock.ts | 1 + plugins/catalog-node/src/testUtils/types.ts | 7 +++- .../src/testUtils/catalogApiMock.ts | 1 + 14 files changed, 132 insertions(+), 67 deletions(-) diff --git a/.changeset/olive-moons-burn.md b/.changeset/olive-moons-burn.md index feeff93d3d..652c8fff14 100644 --- a/.changeset/olive-moons-burn.md +++ b/.changeset/olive-moons-burn.md @@ -4,7 +4,7 @@ '@backstage/plugin-catalog-node': minor --- -Introduced new `streamEntities` async generator method for the catalog. +Introduced new `streamEntities` and `streamEntityPages` async generator methods for the catalog. Catalog API and Catalog Service now includes a `streamEntities` method that allows for streaming entities from the catalog. This method is designed to handle large datasets efficiently by processing entities in a stream rather than loading them @@ -19,3 +19,17 @@ for await (const entity of stream) { // Handle entity } ``` + +Additionally, a `streamEntityPages` method is available that streams entities in pages, allowing for batch processing of entities. +This method is more efficient than `streamEntities` when you can process entities in chunks. +Example usage: + +```ts +const pageStream = catalogClient.streamEntityPages( + { batchSize: 100 }, + { token }, +); +for await (const page of pageStream) { + // Handle page of entities +} +``` diff --git a/packages/catalog-client/report-testUtils.api.md b/packages/catalog-client/report-testUtils.api.md index cadf58aeaf..01bb94214b 100644 --- a/packages/catalog-client/report-testUtils.api.md +++ b/packages/catalog-client/report-testUtils.api.md @@ -71,14 +71,9 @@ export class InMemoryCatalogClient implements CatalogApi { // (undocumented) removeLocationById(_id: string): Promise; // (undocumented) - streamEntities< - T extends StreamEntitiesRequest, - R = T extends { - mode: 'array'; - } & StreamEntitiesRequest - ? Entity[] - : Entity, - >(request?: T): AsyncIterable; + streamEntities(request?: StreamEntitiesRequest): AsyncIterable; + // (undocumented) + streamEntityPages(request?: StreamEntitiesRequest): AsyncIterable; // (undocumented) validateEntity( _entity: Entity, diff --git a/packages/catalog-client/report.api.md b/packages/catalog-client/report.api.md index 6000686fbd..5bb36e16fc 100644 --- a/packages/catalog-client/report.api.md +++ b/packages/catalog-client/report.api.md @@ -91,7 +91,11 @@ export interface CatalogApi { streamEntities( request?: StreamEntitiesRequest, options?: CatalogRequestOptions, - ): AsyncIterable; + ): AsyncIterable; + streamEntityPages( + request?: StreamEntitiesRequest, + options?: CatalogRequestOptions, + ): AsyncIterable; validateEntity( entity: Entity, locationRef: string, @@ -174,14 +178,14 @@ export class CatalogClient implements CatalogApi { id: string, options?: CatalogRequestOptions, ): Promise; - streamEntities< - T extends StreamEntitiesRequest, - R = T extends { - mode: 'array'; - } & StreamEntitiesRequest - ? Entity[] - : Entity, - >(request?: T, options?: CatalogRequestOptions): AsyncIterable; + streamEntities( + request?: StreamEntitiesRequest, + options?: CatalogRequestOptions, + ): AsyncIterable; + streamEntityPages( + request?: StreamEntitiesRequest, + options?: CatalogRequestOptions, + ): AsyncIterable; validateEntity( entity: Entity, locationRef: string, @@ -335,7 +339,6 @@ export type StreamEntitiesRequest = Omit< 'limit' | 'offset' > & { batchSize?: number; - mode?: 'single' | 'array'; }; // @public diff --git a/packages/catalog-client/src/CatalogClient.test.ts b/packages/catalog-client/src/CatalogClient.test.ts index c5e2c8cc21..e330eefd33 100644 --- a/packages/catalog-client/src/CatalogClient.test.ts +++ b/packages/catalog-client/src/CatalogClient.test.ts @@ -579,7 +579,7 @@ describe('CatalogClient', () => { ); }); - it('should stream entities one by one', async () => { + it('should stream entities', async () => { const stream = client.streamEntities({}, { token }); const results: Entity[] = []; for await (const entity of stream) { @@ -588,13 +588,13 @@ 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); + it('should stream entity pages', async () => { + const stream = client.streamEntityPages({}, { token }); + const results: Entity[][] = []; + for await (const entityPage of stream) { + results.push(entityPage); } - expect(results).toEqual(defaultResponse.items); + expect(results).toEqual([defaultResponse.items, []]); }); it('should handle errors', async () => { diff --git a/packages/catalog-client/src/CatalogClient.ts b/packages/catalog-client/src/CatalogClient.ts index b090966807..e83fd362da 100644 --- a/packages/catalog-client/src/CatalogClient.ts +++ b/packages/catalog-client/src/CatalogClient.ts @@ -463,26 +463,34 @@ export class CatalogClient implements CatalogApi { /** * {@inheritdoc CatalogApi.streamEntities} */ - async *streamEntities< - T extends StreamEntitiesRequest, - R = T extends { mode: 'array' } & StreamEntitiesRequest ? Entity[] : Entity, - >(request?: T, options?: CatalogRequestOptions): AsyncIterable { + async *streamEntities( + request?: StreamEntitiesRequest, + options?: CatalogRequestOptions, + ): AsyncIterable { + const pages = this.streamEntityPages(request, options); + for await (const page of pages) { + for (const entity of page) { + yield entity; + } + } + } + + /** + * {@inheritdoc CatalogApi.streamEntityPages} + */ + async *streamEntityPages( + request?: StreamEntitiesRequest, + options?: CatalogRequestOptions, + ): AsyncIterable { 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 } : { ...request, limit }, options, ); - if (mode === 'single') { - for (const entity of res.items) { - yield entity as R; - } - } else { - yield res.items as R; - } + yield res.items; cursor = res.pageInfo.nextCursor; } while (cursor); diff --git a/packages/catalog-client/src/testUtils/InMemoryCatalogClient.test.ts b/packages/catalog-client/src/testUtils/InMemoryCatalogClient.test.ts index afc26346a7..35f0a89248 100644 --- a/packages/catalog-client/src/testUtils/InMemoryCatalogClient.test.ts +++ b/packages/catalog-client/src/testUtils/InMemoryCatalogClient.test.ts @@ -123,7 +123,7 @@ describe('InMemoryCatalogClient', () => { }); }); - it('streamEntities single', async () => { + it('streamEntities', async () => { const client = new InMemoryCatalogClient({ entities }); const stream = client.streamEntities(); const results: Entity[] = []; @@ -133,14 +133,14 @@ describe('InMemoryCatalogClient', () => { expect(results).toEqual(entities); }); - it('streamEntities batch', async () => { + it('streamEntityPages', async () => { const client = new InMemoryCatalogClient({ entities }); - const stream = client.streamEntities({ mode: 'array' }); - const results: Entity[] = []; - for await (const entity of stream) { - results.push(...entity); + const stream = client.streamEntityPages(); + const results: Entity[][] = []; + for await (const page of stream) { + results.push(page); } - expect(results).toEqual(entities); + expect(results).toEqual([entities]); }); it('getEntityAncestors', async () => { diff --git a/packages/catalog-client/src/testUtils/InMemoryCatalogClient.ts b/packages/catalog-client/src/testUtils/InMemoryCatalogClient.ts index fe23e7ceb1..c033ced387 100644 --- a/packages/catalog-client/src/testUtils/InMemoryCatalogClient.ts +++ b/packages/catalog-client/src/testUtils/InMemoryCatalogClient.ts @@ -280,24 +280,27 @@ export class InMemoryCatalogClient implements CatalogApi { throw new NotImplementedError('Method not implemented.'); } - async *streamEntities< - T extends StreamEntitiesRequest, - R = T extends { mode: 'array' } & StreamEntitiesRequest ? Entity[] : Entity, - >(request?: T): AsyncIterable { + async *streamEntities( + request?: StreamEntitiesRequest, + ): AsyncIterable { + const pages = this.streamEntityPages(request); + for await (const page of pages) { + for (const entity of page) { + yield entity; + } + } + } + + async *streamEntityPages( + request?: StreamEntitiesRequest, + ): AsyncIterable { let cursor: string | undefined = undefined; - const mode = request?.mode ?? 'single'; do { const res = await this.queryEntities( cursor ? { ...request, cursor } : request, ); - if (mode === 'single') { - for (const entity of res.items) { - yield entity as R; - } - } else { - yield res.items as R; - } + yield res.items; cursor = res.pageInfo.nextCursor; } while (cursor); diff --git a/packages/catalog-client/src/types/api.ts b/packages/catalog-client/src/types/api.ts index fbf3d82a83..56872c60c8 100644 --- a/packages/catalog-client/src/types/api.ts +++ b/packages/catalog-client/src/types/api.ts @@ -478,10 +478,6 @@ export type StreamEntitiesRequest = Omit< * 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'; }; /** @@ -724,5 +720,19 @@ export interface CatalogApi { streamEntities( request?: StreamEntitiesRequest, options?: CatalogRequestOptions, - ): AsyncIterable; + ): AsyncIterable; + + /** + * Asynchronously streams entity pages from the catalog. Uses `queryEntities` + * to fetch entities in batches, and yields them one page at a time. + * + * @public + * + * @param request - Request parameters + * @param options - Additional options + */ + streamEntityPages( + request?: StreamEntitiesRequest, + options?: CatalogRequestOptions, + ): AsyncIterable; } diff --git a/plugins/catalog-node/report-testUtils.api.md b/plugins/catalog-node/report-testUtils.api.md index 4d0f7fdd48..e06fb741d1 100644 --- a/plugins/catalog-node/report-testUtils.api.md +++ b/plugins/catalog-node/report-testUtils.api.md @@ -111,7 +111,12 @@ export interface CatalogServiceMock extends CatalogService, CatalogApi { streamEntities( request?: StreamEntitiesRequest, options?: CatalogServiceRequestOptions | CatalogRequestOptions, - ): AsyncIterable; + ): AsyncIterable; + // (undocumented) + streamEntityPages( + request?: StreamEntitiesRequest, + options?: CatalogServiceRequestOptions | CatalogRequestOptions, + ): AsyncIterable; // (undocumented) validateEntity( entity: Entity, diff --git a/plugins/catalog-node/report.api.md b/plugins/catalog-node/report.api.md index 6e476f88f6..9173f89e29 100644 --- a/plugins/catalog-node/report.api.md +++ b/plugins/catalog-node/report.api.md @@ -199,7 +199,12 @@ export interface CatalogService { streamEntities( request: StreamEntitiesRequest | undefined, options: CatalogServiceRequestOptions, - ): AsyncIterable; + ): AsyncIterable; + // (undocumented) + streamEntityPages( + request: StreamEntitiesRequest | undefined, + options: CatalogServiceRequestOptions, + ): AsyncIterable; // (undocumented) validateEntity( entity: Entity, diff --git a/plugins/catalog-node/src/catalogService.ts b/plugins/catalog-node/src/catalogService.ts index d8509ddb20..aa6a69a69b 100644 --- a/plugins/catalog-node/src/catalogService.ts +++ b/plugins/catalog-node/src/catalogService.ts @@ -146,7 +146,12 @@ export interface CatalogService { streamEntities( request: StreamEntitiesRequest | undefined, options: CatalogServiceRequestOptions, - ): AsyncIterable; + ): AsyncIterable; + + streamEntityPages( + request: StreamEntitiesRequest | undefined, + options: CatalogServiceRequestOptions, + ): AsyncIterable; } class DefaultCatalogService implements CatalogService { @@ -329,13 +334,23 @@ class DefaultCatalogService implements CatalogService { async *streamEntities( request: StreamEntitiesRequest | undefined, options: CatalogServiceRequestOptions, - ): AsyncIterable { + ): AsyncIterable { yield* this.#catalogApi.streamEntities( request, await this.#getOptions(options), ); } + async *streamEntityPages( + request: StreamEntitiesRequest | undefined, + options: CatalogServiceRequestOptions, + ): AsyncIterable { + yield* this.#catalogApi.streamEntityPages( + request, + await this.#getOptions(options), + ); + } + async #getOptions( options: CatalogServiceRequestOptions, ): Promise { diff --git a/plugins/catalog-node/src/testUtils/catalogServiceMock.ts b/plugins/catalog-node/src/testUtils/catalogServiceMock.ts index 59c66c0b88..78962dc484 100644 --- a/plugins/catalog-node/src/testUtils/catalogServiceMock.ts +++ b/plugins/catalog-node/src/testUtils/catalogServiceMock.ts @@ -104,5 +104,6 @@ export namespace catalogServiceMock { validateEntity: jest.fn(), analyzeLocation: jest.fn(), streamEntities: jest.fn(), + streamEntityPages: jest.fn(), })); } diff --git a/plugins/catalog-node/src/testUtils/types.ts b/plugins/catalog-node/src/testUtils/types.ts index fafdcb82fa..c8ada67dc5 100644 --- a/plugins/catalog-node/src/testUtils/types.ts +++ b/plugins/catalog-node/src/testUtils/types.ts @@ -140,5 +140,10 @@ export interface CatalogServiceMock extends CatalogService, CatalogApi { streamEntities( request?: StreamEntitiesRequest, options?: CatalogServiceRequestOptions | CatalogRequestOptions, - ): AsyncIterable; + ): AsyncIterable; + + streamEntityPages( + request?: StreamEntitiesRequest, + options?: CatalogServiceRequestOptions | CatalogRequestOptions, + ): AsyncIterable; } diff --git a/plugins/catalog-react/src/testUtils/catalogApiMock.ts b/plugins/catalog-react/src/testUtils/catalogApiMock.ts index 594b3a76be..a37d2669ad 100644 --- a/plugins/catalog-react/src/testUtils/catalogApiMock.ts +++ b/plugins/catalog-react/src/testUtils/catalogApiMock.ts @@ -103,5 +103,6 @@ export namespace catalogApiMock { validateEntity: jest.fn(), analyzeLocation: jest.fn(), streamEntities: jest.fn(), + streamEntityPages: jest.fn(), })); }