From 3c1eae932492b2adb5a11e0556a958d912a4f7e9 Mon Sep 17 00:00:00 2001 From: Damon Kaswell Date: Thu, 19 Jan 2023 08:44:30 -0800 Subject: [PATCH] Do not create new event route, rely on events backend instead Signed-off-by: Damon Kaswell --- .../api-report.md | 23 +++-- .../src/engine/IncrementalIngestionEngine.ts | 88 +++++++++---------- .../src/module/WrapperProviders.ts | 4 +- .../src/router/routes.ts | 46 +--------- .../src/types.ts | 42 +++++---- 5 files changed, 84 insertions(+), 119 deletions(-) diff --git a/plugins/catalog-backend-module-incremental-ingestion/api-report.md b/plugins/catalog-backend-module-incremental-ingestion/api-report.md index 3fc932ad2d..4bf6bb9456 100644 --- a/plugins/catalog-backend-module-incremental-ingestion/api-report.md +++ b/plugins/catalog-backend-module-incremental-ingestion/api-report.md @@ -51,23 +51,22 @@ export class IncrementalCatalogBuilder { // @public export interface IncrementalEntityProvider { around(burst: (context: TContext) => Promise): Promise; + eventHandler?: { + onEvent: (params: EventParams) => + | undefined + | { + added: DeferredEntity[]; + removed: { + entityRef: string; + }[]; + }; + supportsEventTopics: () => string[]; + }; getProviderName(): string; next( context: TContext, cursor?: TCursor, ): Promise>; - onEvent?: (params: EventParams) => - | { - delta: - | { - added: DeferredEntity[]; - removed: { - entityRef: string; - }[]; - } - | undefined; - } - | undefined; } // @public (undocumented) diff --git a/plugins/catalog-backend-module-incremental-ingestion/src/engine/IncrementalIngestionEngine.ts b/plugins/catalog-backend-module-incremental-ingestion/src/engine/IncrementalIngestionEngine.ts index c1a62b7a73..228c7a2505 100644 --- a/plugins/catalog-backend-module-incremental-ingestion/src/engine/IncrementalIngestionEngine.ts +++ b/plugins/catalog-backend-module-incremental-ingestion/src/engine/IncrementalIngestionEngine.ts @@ -28,7 +28,6 @@ export class IncrementalIngestionEngine { private readonly restLength: Duration; private readonly backoff: DurationObjectUnits[]; - private readonly providerEventTopic: string; private manager: IncrementalIngestionDatabaseManager; @@ -41,7 +40,6 @@ export class IncrementalIngestionEngine { minutes: 30 }, { hours: 3 }, ]; - this.providerEventTopic = `${options.provider.getProviderName()}-push`; } async taskFn(signal: AbortSignal) { @@ -337,68 +335,68 @@ export class IncrementalIngestionEngine async onEvent(params: EventParams): Promise { const { topic } = params; - if (topic !== this.providerEventTopic) { + if (!this.supportsEventTopics().includes(topic)) { return; } const { logger, provider, connection } = this.options; const providerName = provider.getProviderName(); - logger.debug( - `incremental-engine: Received ${this.providerEventTopic} event`, - ); + logger.debug(`incremental-engine: ${providerName} received ${topic} event`); - if (!provider.onEvent) { + if (!provider.eventHandler) { return; } - const update = provider.onEvent(params); + const delta = provider.eventHandler.onEvent(params); - if (update) { - if (update.delta) { - if (update.delta.added.length > 0) { - const ingestionRecord = await this.manager.getCurrentIngestionRecord( - providerName, + if (delta) { + if (delta.added.length > 0) { + const ingestionRecord = await this.manager.getCurrentIngestionRecord( + providerName, + ); + + if (!ingestionRecord) { + logger.debug( + `incremental-engine: ${providerName} skipping delta addition because incremental ingestion is restarting.`, ); + } else { + const mark = + ingestionRecord.status === 'resting' + ? await this.manager.getLastMark(ingestionRecord.id) + : await this.manager.getFirstMark(ingestionRecord.id); - if (!ingestionRecord) { - logger.debug( - `incremental-engine: Skipping delta addition because incremental ingestion is restarting.`, + if (!mark) { + throw new Error( + `Cannot apply delta, page records are missing! Please re-run incremental ingestion for ${providerName}.`, ); - } else { - const mark = - ingestionRecord.status === 'resting' - ? await this.manager.getLastMark(ingestionRecord.id) - : await this.manager.getFirstMark(ingestionRecord.id); - - if (!mark) { - throw new Error( - `Cannot apply delta, page records are missing! Please re-run incremental ingestion for ${providerName}.`, - ); - } - await this.manager.createMarkEntities(mark.id, update.delta.added); } + await this.manager.createMarkEntities(mark.id, delta.added); } - - if (update.delta.removed.length > 0) { - await this.manager.deleteEntityRecordsByRef(update.delta.removed); - } - - await connection.applyMutation({ - type: 'delta', - ...update.delta, - }); - logger.debug( - `incremental-engine: Processed ${this.providerEventTopic} event`, - ); - } else { - logger.warn( - `incremental-engine: Rejected ${this.providerEventTopic} event - empty or invalid`, - ); } + + if (delta.removed.length > 0) { + await this.manager.deleteEntityRecordsByRef(delta.removed); + } + + await connection.applyMutation({ + type: 'delta', + ...delta, + }); + logger.debug( + `incremental-engine: ${providerName} processed delta from '${topic}' event`, + ); + } else { + logger.warn( + `incremental-engine: Rejected delta from '${topic}' event - empty or invalid`, + ); } } supportsEventTopics(): string[] { - return [this.providerEventTopic]; + const { provider } = this.options; + const topics = provider.eventHandler + ? provider.eventHandler.supportsEventTopics() + : []; + return topics; } } diff --git a/plugins/catalog-backend-module-incremental-ingestion/src/module/WrapperProviders.ts b/plugins/catalog-backend-module-incremental-ingestion/src/module/WrapperProviders.ts index e714568c4b..4738893c95 100644 --- a/plugins/catalog-backend-module-incremental-ingestion/src/module/WrapperProviders.ts +++ b/plugins/catalog-backend-module-incremental-ingestion/src/module/WrapperProviders.ts @@ -81,8 +81,8 @@ export class WrapperProviders { ).createRouter(); } - private async startProvider( - provider: IncrementalEntityProvider, + private async startProvider( + provider: IncrementalEntityProvider, providerOptions: IncrementalEntityProviderOptions, connection: EntityProviderConnection, ) { diff --git a/plugins/catalog-backend-module-incremental-ingestion/src/router/routes.ts b/plugins/catalog-backend-module-incremental-ingestion/src/router/routes.ts index b3e3ae5792..fc26abdf32 100644 --- a/plugins/catalog-backend-module-incremental-ingestion/src/router/routes.ts +++ b/plugins/catalog-backend-module-incremental-ingestion/src/router/routes.ts @@ -15,28 +15,21 @@ */ import { errorHandler } from '@backstage/backend-common'; -import { stringifyError } from '@backstage/errors'; -import { EventBroker, EventPublisher } from '@backstage/plugin-events-node'; import express from 'express'; import Router from 'express-promise-router'; import { Logger } from 'winston'; import { IncrementalIngestionDatabaseManager } from '../database/IncrementalIngestionDatabaseManager'; import { PROVIDER_BASE_PATH, PROVIDER_CLEANUP, PROVIDER_HEALTH } from './paths'; -export class IncrementalProviderRouter implements EventPublisher { +export class IncrementalProviderRouter { private manager: IncrementalIngestionDatabaseManager; private logger: Logger; - private eventBroker: EventBroker | undefined; constructor(manager: IncrementalIngestionDatabaseManager, logger: Logger) { this.manager = manager; this.logger = logger; } - async setEventBroker(eventBroker: EventBroker): Promise { - this.eventBroker = eventBroker; - } - async createRouter() { const router = Router(); router.use(express.json()); @@ -239,43 +232,6 @@ export class IncrementalProviderRouter implements EventPublisher { }); }); - router.post(`${PROVIDER_BASE_PATH}/event`, async (req, res) => { - const { provider } = req.params; - - const topic = `${provider}-push`; - - const eventPayload = req.body; - - if (!this.eventBroker) { - res.status(500).json({ - success: false, - provider, - message: `The payload could not be processed!`, - }); - throw new Error('Event broker not initialized!'); - } - - try { - await this.eventBroker.publish({ - topic, - eventPayload, - }); - res.json({ - success: true, - provider, - message: 'Payload submitted.', - }); - } catch (e) { - res.status(500).json({ - success: false, - provider, - message: `There was an error submitting the payload: ${stringifyError( - e, - )}`, - }); - } - }); - router.use(errorHandler()); return router; diff --git a/plugins/catalog-backend-module-incremental-ingestion/src/types.ts b/plugins/catalog-backend-module-incremental-ingestion/src/types.ts index efadee3a76..cd1e63b044 100644 --- a/plugins/catalog-backend-module-incremental-ingestion/src/types.ts +++ b/plugins/catalog-backend-module-incremental-ingestion/src/types.ts @@ -78,23 +78,35 @@ export interface IncrementalEntityProvider { around(burst: (context: TContext) => Promise): Promise; /** - * This method accepts an incoming event for the provider, and - * (optionally) maps the payload to an object containing a delta - * mutation. + * If set, the IncrementalEntityProvider will receive and respond to + * events. * - * If a valid delta is returned by this method, it will be ingested - * automatically by the provider. + * This system acts as a wrapper for the Backstage events bus, and + * requires the events backend to function. It does not provide its + * own events backend. See {@link https://github.com/backstage/backstage/tree/master/plugins/events-backend}. */ - onEvent?: (params: EventParams) => - | { - delta: - | { - added: DeferredEntity[]; - removed: { entityRef: string }[]; - } - | undefined; - } - | undefined; + eventHandler?: { + /** + * This method accepts an incoming event for the provider, and + * optionally maps the payload to an object containing a delta + * mutation. + * + * If a valid delta is returned by this method, it will be ingested + * automatically by the provider. + */ + onEvent: (params: EventParams) => + | undefined + | { + added: DeferredEntity[]; + removed: { entityRef: string }[]; + }; + + /** + * This method returns an array of topics for the IncrementalEntityProvider + * to respond to. + */ + supportsEventTopics: () => string[]; + }; } /**