diff --git a/.changeset/rich-beers-eat.md b/.changeset/rich-beers-eat.md new file mode 100644 index 0000000000..55115e929e --- /dev/null +++ b/.changeset/rich-beers-eat.md @@ -0,0 +1,5 @@ +--- +'@backstage/plugin-catalog-backend-module-incremental-ingestion': patch +--- + +Wire up the events together in the new backend system diff --git a/plugins/catalog-backend-module-incremental-ingestion/src/module/WrapperProviders.test.ts b/plugins/catalog-backend-module-incremental-ingestion/src/module/WrapperProviders.test.ts index 0f98384893..95b3f203c4 100644 --- a/plugins/catalog-backend-module-incremental-ingestion/src/module/WrapperProviders.test.ts +++ b/plugins/catalog-backend-module-incremental-ingestion/src/module/WrapperProviders.test.ts @@ -68,6 +68,7 @@ describe('WrapperProviders', () => { client, scheduler: scheduler as Partial as SchedulerService, applyDatabaseMigrations, + events: mockServices.events.mock(), }); const wrapped1 = providers.wrap(provider1, { burstInterval: { seconds: 1 }, 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 c322d9adda..738d5733d0 100644 --- a/plugins/catalog-backend-module-incremental-ingestion/src/module/WrapperProviders.ts +++ b/plugins/catalog-backend-module-incremental-ingestion/src/module/WrapperProviders.ts @@ -36,6 +36,7 @@ import { IncrementalEntityProvider, IncrementalEntityProviderOptions, } from '../types'; +import { EventsService } from '@backstage/plugin-events-node'; /** * Helps in the creation of the catalog entity providers that wrap the @@ -53,6 +54,7 @@ export class WrapperProviders { client: Knex; scheduler: SchedulerService; applyDatabaseMigrations?: typeof applyDatabaseMigrations; + events: EventsService; }, ) {} @@ -130,6 +132,20 @@ export class WrapperProviders { frequency, timeout: length, }); + + const topics = engine.supportsEventTopics(); + if (topics.length > 0) { + logger.info( + `Provider ${provider.getProviderName()} subscribing to events for topics: ${topics.join( + ',', + )}`, + ); + await this.options.events.subscribe({ + topics, + id: `catalog-backend-module-incremental-ingestion:${provider.getProviderName()}`, + onEvent: evt => engine.onEvent(evt), + }); + } } catch (error) { logger.warn( `Failed to initialize incremental ingestion provider ${provider.getProviderName()}, ${stringifyError( diff --git a/plugins/catalog-backend-module-incremental-ingestion/src/module/catalogModuleIncrementalIngestionEntityProvider.ts b/plugins/catalog-backend-module-incremental-ingestion/src/module/catalogModuleIncrementalIngestionEntityProvider.ts index 8752bd3137..25ff33d6ba 100644 --- a/plugins/catalog-backend-module-incremental-ingestion/src/module/catalogModuleIncrementalIngestionEntityProvider.ts +++ b/plugins/catalog-backend-module-incremental-ingestion/src/module/catalogModuleIncrementalIngestionEntityProvider.ts @@ -25,6 +25,7 @@ import { IncrementalEntityProviderOptions, } from '@backstage/plugin-catalog-backend-module-incremental-ingestion'; import { WrapperProviders } from './WrapperProviders'; +import { eventsServiceRef } from '@backstage/plugin-events-node'; /** * @public @@ -106,6 +107,7 @@ export const catalogModuleIncrementalIngestionEntityProvider = httpRouter: coreServices.httpRouter, logger: coreServices.logger, scheduler: coreServices.scheduler, + events: eventsServiceRef, }, async init({ catalog, @@ -114,6 +116,7 @@ export const catalogModuleIncrementalIngestionEntityProvider = httpRouter, logger, scheduler, + events, }) { const client = await database.getClient(); @@ -122,6 +125,7 @@ export const catalogModuleIncrementalIngestionEntityProvider = logger, client, scheduler, + events, }); for (const entry of addedProviders) {