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(), })); }