From d0e85ef9aff735322068a24e16d36d10610f731e Mon Sep 17 00:00:00 2001 From: Damon Kaswell Date: Wed, 9 Nov 2022 16:26:20 -0800 Subject: [PATCH] Code hygiene and cleanup Signed-off-by: Damon Kaswell --- .../incremental-ingestion-backend/README.md | 2 +- .../api-report.md | 30 +-------- .../package.json | 4 +- .../src/engine/IncrementalIngestionEngine.ts | 62 +++++++++++-------- .../src/index.ts | 17 ++++- .../src/routes.ts | 2 +- .../src/service/IncrementalCatalogBuilder.ts | 23 +++---- .../src/types.ts | 14 ++--- yarn.lock | 12 +--- 9 files changed, 75 insertions(+), 91 deletions(-) diff --git a/plugins/incremental-ingestion-backend/README.md b/plugins/incremental-ingestion-backend/README.md index 4e9d4fe65b..0e0a4e56e6 100644 --- a/plugins/incremental-ingestion-backend/README.md +++ b/plugins/incremental-ingestion-backend/README.md @@ -42,7 +42,7 @@ The Incremental Entity Provider backend is designed for data sources that provid ## Installation -1. Install `@backstage/plugin-incremental-ingestion-backend` with `yarn add @backstage/plugin-incremental-ingestion-backend` +1. Install `@backstage/plugin-incremental-ingestion-backend` with `yarn workspace backend add @backstage/plugin-incremental-ingestion-backend` 2. Import `IncrementalCatalogBuilder` from `@backstage/plugin-incremental-ingestion-backend` and instantiate it with `await IncrementalCatalogBuilder.create(env, builder)`. You have to pass `builder` into `IncrementalCatalogBuilder.create` function because `IncrementalCatalogBuilder` will convert an `IncrementalEntityProvider` into an `EntityProvider` and call `builder.addEntityProvider`. ```ts diff --git a/plugins/incremental-ingestion-backend/api-report.md b/plugins/incremental-ingestion-backend/api-report.md index 375eedeae0..d843aad1a3 100644 --- a/plugins/incremental-ingestion-backend/api-report.md +++ b/plugins/incremental-ingestion-backend/api-report.md @@ -10,21 +10,19 @@ import type { Config } from '@backstage/config'; import type { DeferredEntity } from '@backstage/plugin-catalog-backend'; import { Duration } from 'luxon'; import type { DurationObjectUnits } from 'luxon'; -import type { EntityProviderConnection } from '@backstage/plugin-catalog-backend'; import { Knex } from 'knex'; import type { Logger } from 'winston'; import type { PermissionAuthorizer } from '@backstage/plugin-permission-common'; import type { PluginDatabaseManager } from '@backstage/backend-common'; import type { PluginTaskScheduler } from '@backstage/backend-tasks'; import { Router } from 'express'; -import type { TaskFunction } from '@backstage/backend-tasks'; import type { UrlReader } from '@backstage/backend-common'; // @public export interface EntityIteratorResult { - cursor: T; + cursor?: T; done: boolean; - entities: DeferredEntity[]; + entities?: DeferredEntity[]; } // @public @@ -198,30 +196,6 @@ export interface IngestionUpsertIFace { | 'backing off'; } -// @public (undocumented) -export interface IterationEngine { - // (undocumented) - taskFn: TaskFunction; -} - -// @public (undocumented) -export interface IterationEngineOptions { - // (undocumented) - backoff?: IncrementalEntityProviderOptions['backoff']; - // (undocumented) - connection: EntityProviderConnection; - // (undocumented) - logger: Logger; - // (undocumented) - manager: IncrementalIngestionDatabaseManager; - // (undocumented) - provider: IncrementalEntityProvider; - // (undocumented) - ready: Promise; - // (undocumented) - restLength: DurationObjectUnits; -} - // @public export interface MarkRecord { // (undocumented) diff --git a/plugins/incremental-ingestion-backend/package.json b/plugins/incremental-ingestion-backend/package.json index 6fdcdbb138..195adab7ea 100644 --- a/plugins/incremental-ingestion-backend/package.json +++ b/plugins/incremental-ingestion-backend/package.json @@ -1,10 +1,10 @@ { "name": "@backstage/plugin-incremental-ingestion-backend", + "description": "An entity provider for streaming large asset sources into the catalog", "version": "0.1.0", "main": "src/index.ts", "types": "src/index.ts", "license": "Apache-2.0", - "private": true, "publishConfig": { "access": "public", "main": "dist/index.cjs.js", @@ -30,7 +30,7 @@ "@backstage/plugin-catalog-backend": "workspace:^", "@backstage/plugin-permission-common": "workspace:^", "@types/express": "^4.17.6", - "@types/luxon": "3.0.0", + "@types/luxon": "^3.0.0", "express": "^4.17.1", "express-promise-router": "^4.1.0", "knex": "^2.0.0", diff --git a/plugins/incremental-ingestion-backend/src/engine/IncrementalIngestionEngine.ts b/plugins/incremental-ingestion-backend/src/engine/IncrementalIngestionEngine.ts index 4b8588cdfe..6d35f2db4b 100644 --- a/plugins/incremental-ingestion-backend/src/engine/IncrementalIngestionEngine.ts +++ b/plugins/incremental-ingestion-backend/src/engine/IncrementalIngestionEngine.ts @@ -27,6 +27,7 @@ import { performance } from 'perf_hooks'; import { Duration, DurationObjectUnits } from 'luxon'; import { v4 } from 'uuid'; +/** @public */ export class IncrementalIngestionEngine implements IterationEngine { restLength: Duration; backoff: DurationObjectUnits[]; @@ -193,7 +194,7 @@ export class IncrementalIngestionEngine implements IterationEngine { async ingestOneBurst(id: string, signal: AbortSignal) { const lastMark = await this.manager.getLastMark(id); - const cursor = lastMark ? lastMark.cursor : void 0; + const cursor = lastMark ? lastMark.cursor : undefined; let sequence = lastMark ? lastMark.sequence + 1 : 0; const start = performance.now(); @@ -206,10 +207,15 @@ export class IncrementalIngestionEngine implements IterationEngine { await this.options.provider.around(async (context: unknown) => { let next = await this.options.provider.next(context, cursor); count++; - // eslint-disable-next-line no-constant-condition - while (true) { + for (;;) { done = next.done; - await this.mark(id, sequence, next.entities, next.done, next.cursor); + await this.mark({ + id, + sequence, + entities: next?.entities, + done: next.done, + cursor: next?.cursor, + }); if (signal.aborted || next.done) { break; } else { @@ -228,17 +234,20 @@ export class IncrementalIngestionEngine implements IterationEngine { return done; } - async mark( - id: string, - sequence: number, - entities: DeferredEntity[], - done: boolean, - cursor?: unknown, - ) { + async mark(options: { + id: string; + sequence: number; + entities?: DeferredEntity[]; + done: boolean; + cursor?: unknown; + }) { + const { id, sequence, entities, done, cursor } = options; this.options.logger.debug( `incremental-engine: Ingestion '${id}': MARK ${ - entities.length - } entities, cursor: ${JSON.stringify(cursor)}, done: ${done}`, + entities ? entities.length : 0 + } entities, cursor: ${ + cursor ? JSON.stringify(cursor) : 'none' + }, done: ${done}`, ); const markId = v4(); @@ -251,24 +260,25 @@ export class IncrementalIngestionEngine implements IterationEngine { }, }); - if (entities.length > 0) { + if (entities && entities.length > 0) { await this.manager.createMarkEntities(markId, entities); } - const added = entities.map(deferred => ({ - ...deferred, - entity: { - ...deferred.entity, - metadata: { - ...deferred.entity.metadata, - annotations: { - ...deferred.entity.metadata.annotations, - [INCREMENTAL_ENTITY_PROVIDER_ANNOTATION]: - this.options.provider.getProviderName(), + const added = + entities?.map(deferred => ({ + ...deferred, + entity: { + ...deferred.entity, + metadata: { + ...deferred.entity.metadata, + annotations: { + ...deferred.entity.metadata.annotations, + [INCREMENTAL_ENTITY_PROVIDER_ANNOTATION]: + this.options.provider.getProviderName(), + }, }, }, - }, - })); + })) ?? []; const removed: DeferredEntity[] = done ? [] diff --git a/plugins/incremental-ingestion-backend/src/index.ts b/plugins/incremental-ingestion-backend/src/index.ts index 46afb24a25..da50c6b37f 100644 --- a/plugins/incremental-ingestion-backend/src/index.ts +++ b/plugins/incremental-ingestion-backend/src/index.ts @@ -13,6 +13,17 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -export * from './service/IncrementalCatalogBuilder'; -export * from './types'; -export * from './database/IncrementalIngestionDatabaseManager'; +export { IncrementalCatalogBuilder } from './service/IncrementalCatalogBuilder'; +export type { IncrementalIngestionDatabaseManager } from './database/IncrementalIngestionDatabaseManager'; +export type { + EntityIteratorResult, + IncrementalEntityProvider, + IncrementalEntityProviderOptions, + IngestionRecord, + IngestionRecordUpdate, + IngestionUpsertIFace, + MarkRecord, + MarkRecordInsert, + INCREMENTAL_ENTITY_PROVIDER_ANNOTATION, + PluginEnvironment, +} from './types'; diff --git a/plugins/incremental-ingestion-backend/src/routes.ts b/plugins/incremental-ingestion-backend/src/routes.ts index 81ebc1d7b6..2f51ef6499 100644 --- a/plugins/incremental-ingestion-backend/src/routes.ts +++ b/plugins/incremental-ingestion-backend/src/routes.ts @@ -17,7 +17,7 @@ import { errorHandler } from '@backstage/backend-common'; import express from 'express'; import Router from 'express-promise-router'; import { Logger } from 'winston'; -import { IncrementalIngestionDatabaseManager } from './database/IncrementalIngestionDatabaseManager'; +import { IncrementalIngestionDatabaseManager } from './'; /** @public */ export const createIncrementalProviderRouter = async ( diff --git a/plugins/incremental-ingestion-backend/src/service/IncrementalCatalogBuilder.ts b/plugins/incremental-ingestion-backend/src/service/IncrementalCatalogBuilder.ts index e6e1e05091..7bf39e160f 100644 --- a/plugins/incremental-ingestion-backend/src/service/IncrementalCatalogBuilder.ts +++ b/plugins/incremental-ingestion-backend/src/service/IncrementalCatalogBuilder.ts @@ -27,12 +27,11 @@ import { IncrementalIngestionDatabaseManager } from '../database/IncrementalInge import { createIncrementalProviderRouter } from '../routes'; class Deferred implements Promise { - // eslint-disable-next-line @typescript-eslint/ban-ts-comment - // @ts-ignore - resolve: (value: T) => void; - // eslint-disable-next-line @typescript-eslint/ban-ts-comment - // @ts-ignore - reject: (error: Error) => void; + #resolve?: (value: T) => void; + #reject?: (error: Error) => void; + + get resolve() { return this.#resolve!; } + get reject() { return this.#reject!; } then: Promise['then']; catch: Promise['catch']; @@ -40,8 +39,8 @@ class Deferred implements Promise { constructor() { const promise = new Promise((resolve, reject) => { - this.resolve = resolve; - this.reject = reject; + this.#resolve = resolve; + this.#reject = reject; }); this.then = promise.then.bind(promise); @@ -98,9 +97,11 @@ export class IncrementalCatalogBuilder { options: IncrementalEntityProviderOptions, ) { const { burstInterval, burstLength, restLength } = options; - const { logger: catalogLogger, database, scheduler } = this.env; + const { logger: catalogLogger, scheduler } = this.env; const ready = this.ready; + const manager = this.manager; + this.builder.addEntityProvider({ getProviderName: provider.getProviderName.bind(provider), async connect(connection) { @@ -110,10 +111,6 @@ export class IncrementalCatalogBuilder { logger.info(`Connecting`); - const client = await database.getClient(); - - const manager = new IncrementalIngestionDatabaseManager({ client }); - const engine = new IncrementalIngestionEngine({ ...options, ready, diff --git a/plugins/incremental-ingestion-backend/src/types.ts b/plugins/incremental-ingestion-backend/src/types.ts index 8a3233e3b7..2671752ad6 100644 --- a/plugins/incremental-ingestion-backend/src/types.ts +++ b/plugins/incremental-ingestion-backend/src/types.ts @@ -29,7 +29,7 @@ import type { import type { PermissionAuthorizer } from '@backstage/plugin-permission-common'; import type { DurationObjectUnits } from 'luxon'; import type { Logger } from 'winston'; -import { IncrementalIngestionDatabaseManager } from './database/IncrementalIngestionDatabaseManager'; +import { IncrementalIngestionDatabaseManager } from './'; /** * Entity annotation containing the incremental entity provider. * @@ -100,12 +100,12 @@ export interface EntityIteratorResult { * A value that marks the page of entities after this one. It will * be used to pass into the following invocation of `next()` */ - cursor: T; + cursor?: T; /** * The entities to ingest. */ - entities: DeferredEntity[]; + entities?: DeferredEntity[]; } /** @public */ @@ -225,7 +225,7 @@ export interface IngestionUpsertIFace { /** * This interface is for updating an existing ingestion record. - * + * * @public */ export interface IngestionRecordUpdate { @@ -235,7 +235,7 @@ export interface IngestionRecordUpdate { /** * The expected response from the `ingestion.ingestion_marks` table. - * + * * @public */ export interface MarkRecord { @@ -248,7 +248,7 @@ export interface MarkRecord { /** * The expected response from the `ingestion.ingestions` table. - * + * * @public */ export interface IngestionRecord extends IngestionUpsertIFace { @@ -262,7 +262,7 @@ export interface IngestionRecord extends IngestionUpsertIFace { /** * This interface supplies all the values for adding an ingestion mark. - * + * * @public */ export interface MarkRecordInsert { diff --git a/yarn.lock b/yarn.lock index 4cc281f7dd..17e002800a 100644 --- a/yarn.lock +++ b/yarn.lock @@ -6023,7 +6023,7 @@ __metadata: languageName: unknown linkType: soft -"@backstage/plugin-incremental-ingestion-backend@workspace:^, @backstage/plugin-incremental-ingestion-backend@workspace:plugins/incremental-ingestion-backend": +"@backstage/plugin-incremental-ingestion-backend@workspace:plugins/incremental-ingestion-backend": version: 0.0.0-use.local resolution: "@backstage/plugin-incremental-ingestion-backend@workspace:plugins/incremental-ingestion-backend" dependencies: @@ -6035,7 +6035,7 @@ __metadata: "@backstage/plugin-catalog-backend": "workspace:^" "@backstage/plugin-permission-common": "workspace:^" "@types/express": ^4.17.6 - "@types/luxon": 3.0.0 + "@types/luxon": ^3.0.0 express: ^4.17.1 express-promise-router: ^4.1.0 knex: ^2.0.0 @@ -13805,13 +13805,6 @@ __metadata: languageName: node linkType: hard -"@types/luxon@npm:3.0.0": - version: 3.0.0 - resolution: "@types/luxon@npm:3.0.0" - checksum: 7738d3f4b91097a8139eca966ab53c6b40bb417b2cc2014c3bbd0aee033d6ff54af14c62046c2eb042434923b6bf796874f24b2af90e07ef422e6d841463d4b2 - languageName: node - linkType: hard - "@types/markdown-it@npm:^12.2.3": version: 12.2.3 resolution: "@types/markdown-it@npm:12.2.3" @@ -21379,7 +21372,6 @@ __metadata: "@backstage/plugin-events-backend": "workspace:^" "@backstage/plugin-events-node": "workspace:^" "@backstage/plugin-graphql-backend": "workspace:^" - "@backstage/plugin-incremental-ingestion-backend": "workspace:^" "@backstage/plugin-jenkins-backend": "workspace:^" "@backstage/plugin-kafka-backend": "workspace:^" "@backstage/plugin-kubernetes-backend": "workspace:^"