diff --git a/plugins/catalog-backend/src/database/DefaultProcessingDatabase.test.ts b/plugins/catalog-backend/src/database/DefaultProcessingDatabase.test.ts index a09e938575..790762a428 100644 --- a/plugins/catalog-backend/src/database/DefaultProcessingDatabase.test.ts +++ b/plugins/catalog-backend/src/database/DefaultProcessingDatabase.test.ts @@ -58,6 +58,7 @@ describe('DefaultProcessingDatabase', () => { minSeconds: 100, maxSeconds: 150, }), + events: mockServices.events.mock(), }), }; } diff --git a/plugins/catalog-backend/src/database/DefaultProcessingDatabase.ts b/plugins/catalog-backend/src/database/DefaultProcessingDatabase.ts index 81c26c66ab..50c3cbb1c8 100644 --- a/plugins/catalog-backend/src/database/DefaultProcessingDatabase.ts +++ b/plugins/catalog-backend/src/database/DefaultProcessingDatabase.ts @@ -42,11 +42,7 @@ import { checkLocationKeyConflict } from './operations/refreshState/checkLocatio import { insertUnprocessedEntity } from './operations/refreshState/insertUnprocessedEntity'; import { updateUnprocessedEntity } from './operations/refreshState/updateUnprocessedEntity'; import { generateStableHash, generateTargetKey } from './util'; -import { - EventBroker, - EventParams, - EventsService, -} from '@backstage/plugin-events-node'; +import { EventParams, EventsService } from '@backstage/plugin-events-node'; import { DateTime } from 'luxon'; import { CATALOG_CONFLICTS_TOPIC } from '../constants'; import { CatalogConflictEventPayload } from '../catalog/types'; @@ -64,7 +60,7 @@ export class DefaultProcessingDatabase implements ProcessingDatabase { database: Knex; logger: LoggerService; refreshInterval: ProcessingIntervalFunction; - eventBroker?: EventBroker | EventsService; + events: EventsService; }, ) { initDatabaseMetrics(options.database); @@ -367,7 +363,7 @@ export class DefaultProcessingDatabase implements ProcessingDatabase { this.options.logger.warn( `Detected conflicting entityRef ${entityRef} already referenced by ${conflictingKey} and now also ${locationKey}`, ); - if (this.options.eventBroker && locationKey) { + if (locationKey) { const eventParams: EventParams = { topic: CATALOG_CONFLICTS_TOPIC, eventPayload: { @@ -378,7 +374,7 @@ export class DefaultProcessingDatabase implements ProcessingDatabase { lastConflictAt: DateTime.now().toISO()!, }, }; - await this.options.eventBroker?.publish(eventParams); + await this.options.events.publish(eventParams); } } } diff --git a/plugins/catalog-backend/src/processing/DefaultCatalogProcessingEngine.test.ts b/plugins/catalog-backend/src/processing/DefaultCatalogProcessingEngine.test.ts index 93c22aab1e..5e6312ec4a 100644 --- a/plugins/catalog-backend/src/processing/DefaultCatalogProcessingEngine.test.ts +++ b/plugins/catalog-backend/src/processing/DefaultCatalogProcessingEngine.test.ts @@ -71,6 +71,7 @@ describe('DefaultCatalogProcessingEngine', () => { stitcher: stitcher, createHash: () => hash, scheduler: mockServices.scheduler(), + events: mockServices.events.mock(), }); db.transaction.mockImplementation(cb => cb((() => {}) as any)); @@ -139,6 +140,7 @@ describe('DefaultCatalogProcessingEngine', () => { stitcher: stitcher, scheduler: mockServices.scheduler(), createHash: () => hash, + events: mockServices.events.mock(), }); db.transaction.mockImplementation(cb => cb((() => {}) as any)); @@ -223,6 +225,7 @@ describe('DefaultCatalogProcessingEngine', () => { stitcher: stitcher, scheduler: mockServices.scheduler(), createHash: () => hash, + events: mockServices.events.mock(), }); db.transaction.mockImplementation(cb => cb((() => {}) as any)); @@ -301,6 +304,7 @@ describe('DefaultCatalogProcessingEngine', () => { stitcher: stitcher, scheduler: mockServices.scheduler(), createHash: () => hash, + events: mockServices.events.mock(), }); db.transaction.mockImplementation(cb => cb((() => {}) as any)); @@ -362,6 +366,7 @@ describe('DefaultCatalogProcessingEngine', () => { scheduler: mockServices.scheduler(), createHash: () => hash, pollingIntervalMs: 100, + events: mockServices.events.mock(), }); db.transaction.mockImplementation(cb => cb((() => {}) as any)); @@ -478,6 +483,7 @@ describe('DefaultCatalogProcessingEngine', () => { scheduler: mockServices.scheduler(), createHash: () => hash, pollingIntervalMs: 100, + events: mockServices.events.mock(), }); db.transaction.mockImplementation(cb => cb((() => {}) as any)); @@ -584,6 +590,7 @@ describe('DefaultCatalogProcessingEngine', () => { scheduler: mockServices.scheduler(), createHash: () => hash, pollingIntervalMs: 100, + events: mockServices.events.mock(), }); db.transaction.mockImplementation(cb => cb((() => {}) as any)); @@ -668,6 +675,7 @@ describe('DefaultCatalogProcessingEngine', () => { scheduler: mockServices.scheduler(), createHash: () => hash, pollingIntervalMs: 100, + events: mockServices.events.mock(), }); db.transaction.mockImplementation(cb => cb((() => {}) as any)); @@ -757,6 +765,7 @@ describe('DefaultCatalogProcessingEngine', () => { scheduler: mockServices.scheduler(), createHash: () => hash, pollingIntervalMs: 100, + events: mockServices.events.mock(), }); db.transaction.mockImplementation(cb => cb((() => {}) as any)); diff --git a/plugins/catalog-backend/src/processing/DefaultCatalogProcessingEngine.ts b/plugins/catalog-backend/src/processing/DefaultCatalogProcessingEngine.ts index 7061309175..b2a562ea6a 100644 --- a/plugins/catalog-backend/src/processing/DefaultCatalogProcessingEngine.ts +++ b/plugins/catalog-backend/src/processing/DefaultCatalogProcessingEngine.ts @@ -36,7 +36,7 @@ import { withActiveSpan, } from '../util/opentelemetry'; import { deleteOrphanedEntities } from '../database/operations/util/deleteOrphanedEntities'; -import { EventBroker, EventsService } from '@backstage/plugin-events-node'; +import { EventsService } from '@backstage/plugin-events-node'; import { CATALOG_ERRORS_TOPIC } from '../constants'; import { LoggerService, SchedulerService } from '@backstage/backend-plugin-api'; @@ -73,7 +73,7 @@ export class DefaultCatalogProcessingEngine { errors: Error[]; }) => Promise | void; private readonly tracker: ProgressTracker; - private readonly eventBroker?: EventBroker | EventsService; + private readonly events: EventsService; private stopFunc?: () => void; @@ -93,7 +93,7 @@ export class DefaultCatalogProcessingEngine { errors: Error[]; }) => Promise | void; tracker?: ProgressTracker; - eventBroker?: EventBroker | EventsService; + events: EventsService; }) { this.config = options.config; this.scheduler = options.scheduler; @@ -107,7 +107,7 @@ export class DefaultCatalogProcessingEngine { this.orphanCleanupIntervalMs = options.orphanCleanupIntervalMs ?? 30_000; this.onProcessingError = options.onProcessingError; this.tracker = options.tracker ?? progressTracker(); - this.eventBroker = options.eventBroker; + this.events = options.events; this.stopFunc = undefined; } @@ -201,7 +201,7 @@ export class DefaultCatalogProcessingEngine { const location = unprocessedEntity?.metadata?.annotations?.[ANNOTATION_LOCATION]; if (result.errors.length) { - this.eventBroker?.publish({ + this.events.publish({ topic: CATALOG_ERRORS_TOPIC, eventPayload: { entity: entityRef, diff --git a/plugins/catalog-backend/src/processing/DefaultCatalogProcessingOrchestrator.test.ts b/plugins/catalog-backend/src/processing/DefaultCatalogProcessingOrchestrator.test.ts index 31d7c34b94..beb3611206 100644 --- a/plugins/catalog-backend/src/processing/DefaultCatalogProcessingOrchestrator.test.ts +++ b/plugins/catalog-backend/src/processing/DefaultCatalogProcessingOrchestrator.test.ts @@ -97,7 +97,6 @@ describe('DefaultCatalogProcessingOrchestrator', () => { parser: defaultEntityDataParser, policy: EntityPolicies.allOf([]), rulesEnforcer: { isAllowed: () => true }, - legacySingleProcessorValidation: false, }); it('runs a minimal processing', async () => { @@ -192,7 +191,7 @@ describe('DefaultCatalogProcessingOrchestrator', () => { }); }); - it('runs all processor validations when asked to', async () => { + it('runs all processor validations', async () => { const validate = jest.fn(async () => true); const processor1: CatalogProcessor = { getProcessorName: () => 'processor1', @@ -213,7 +212,6 @@ describe('DefaultCatalogProcessingOrchestrator', () => { parser: defaultEntityDataParser, policy: EntityPolicies.allOf([]), rulesEnforcer: { isAllowed: () => true }, - legacySingleProcessorValidation: true, }); const modern = new DefaultCatalogProcessingOrchestrator({ @@ -226,13 +224,12 @@ describe('DefaultCatalogProcessingOrchestrator', () => { parser: defaultEntityDataParser, policy: EntityPolicies.allOf([]), rulesEnforcer: { isAllowed: () => true }, - legacySingleProcessorValidation: false, }); await expect(legacy.process({ entity })).resolves.toMatchObject({ ok: true, }); - expect(validate).toHaveBeenCalledTimes(1); + expect(validate).toHaveBeenCalledTimes(2); validate.mockClear(); @@ -291,7 +288,6 @@ describe('DefaultCatalogProcessingOrchestrator', () => { parser, policy: EntityPolicies.allOf([]), rulesEnforcer, - legacySingleProcessorValidation: false, }); rulesEnforcer.isAllowed.mockReturnValueOnce(true); @@ -333,7 +329,6 @@ describe('DefaultCatalogProcessingOrchestrator', () => { parser, policy: EntityPolicies.allOf([new FailingEntityPolicy()]), rulesEnforcer, - legacySingleProcessorValidation: false, }); await expect( diff --git a/plugins/catalog-backend/src/processing/DefaultCatalogProcessingOrchestrator.ts b/plugins/catalog-backend/src/processing/DefaultCatalogProcessingOrchestrator.ts index b1fcba5d40..79873264a3 100644 --- a/plugins/catalog-backend/src/processing/DefaultCatalogProcessingOrchestrator.ts +++ b/plugins/catalog-backend/src/processing/DefaultCatalogProcessingOrchestrator.ts @@ -95,7 +95,6 @@ export class DefaultCatalogProcessingOrchestrator parser: CatalogProcessorParser; policy: EntityPolicy; rulesEnforcer: CatalogRulesEnforcer; - legacySingleProcessorValidation: boolean; }, ) {} @@ -309,9 +308,6 @@ export class DefaultCatalogProcessingOrchestrator ); if (thisValid) { valid = true; - if (this.options.legacySingleProcessorValidation) { - break; - } } } catch (e) { throw new InputError( diff --git a/plugins/catalog-backend/src/service/CatalogBuilder.ts b/plugins/catalog-backend/src/service/CatalogBuilder.ts index 40fcc3a624..0b35e84102 100644 --- a/plugins/catalog-backend/src/service/CatalogBuilder.ts +++ b/plugins/catalog-backend/src/service/CatalogBuilder.ts @@ -55,7 +55,7 @@ import { PlaceholderResolver, ScmLocationAnalyzer, } from '@backstage/plugin-catalog-node'; -import { EventBroker, EventsService } from '@backstage/plugin-events-node'; +import { EventsService } from '@backstage/plugin-events-node'; import { Permission, PermissionAuthorizer, @@ -123,6 +123,7 @@ export type CatalogEnvironment = { auth: AuthService; httpAuth: HttpAuthService; auditor: AuditorService; + events: EventsService; }; /** @@ -168,8 +169,6 @@ export class CatalogBuilder { private readonly permissions: Permission[]; private readonly permissionRules: CatalogPermissionRuleInput[]; private allowedLocationType: string[]; - private legacySingleProcessorValidation = false; - private eventBroker?: EventBroker | EventsService; /** * Creates a catalog builder. @@ -216,31 +215,6 @@ export class CatalogBuilder { return this; } - /** - * Processing interval determines how often entities should be processed. - * Seconds provided will be multiplied by 1.5 - * The default processing interval is 100-150 seconds. - * setting this too low will potentially deplete request quotas to upstream services. - */ - setProcessingIntervalSeconds(seconds: number): CatalogBuilder { - this.processingInterval = createRandomProcessingInterval({ - minSeconds: seconds, - maxSeconds: seconds * 1.5, - }); - return this; - } - - /** - * Overwrites the default processing interval function used to spread - * entity updates in the catalog. - */ - setProcessingInterval( - processingInterval: ProcessingIntervalFunction, - ): CatalogBuilder { - this.processingInterval = processingInterval; - return this; - } - /** * Overwrites the default location analyzer. */ @@ -427,23 +401,6 @@ export class CatalogBuilder { return this; } - /** - * Enables the legacy behaviour of canceling validation early whenever only a - * single processor declares an entity kind to be valid. - */ - useLegacySingleProcessorValidation(): this { - this.legacySingleProcessorValidation = true; - return this; - } - - /** - * Enables the publishing of events for conflicts in the DefaultProcessingDatabase - */ - setEventBroker(broker: EventBroker | EventsService): CatalogBuilder { - this.eventBroker = broker; - return this; - } - /** * Wires up and returns all of the component parts of the catalog */ @@ -461,6 +418,7 @@ export class CatalogBuilder { auditor, auth, httpAuth, + events, } = this.env; const enableRelationsCompatibility = Boolean( @@ -485,8 +443,8 @@ export class CatalogBuilder { const processingDatabase = new DefaultProcessingDatabase({ database: dbClient, logger, + events, refreshInterval: this.processingInterval, - eventBroker: this.eventBroker, }); const providerDatabase = new DefaultProviderDatabase({ database: dbClient, @@ -523,7 +481,6 @@ export class CatalogBuilder { logger, parser, policy, - legacySingleProcessorValidation: this.legacySingleProcessorValidation, }); const entitiesCatalog = new AuthorizedEntitiesCatalog( @@ -588,7 +545,7 @@ export class CatalogBuilder { onProcessingError: event => { this.onProcessingError?.(event); }, - eventBroker: this.eventBroker, + events, }); const locationAnalyzer = diff --git a/plugins/catalog-backend/src/service/CatalogPlugin.ts b/plugins/catalog-backend/src/service/CatalogPlugin.ts index a9bb2c9285..24147e3cc5 100644 --- a/plugins/catalog-backend/src/service/CatalogPlugin.ts +++ b/plugins/catalog-backend/src/service/CatalogPlugin.ts @@ -271,10 +271,9 @@ export const catalogPlugin = createBackendPlugin({ auth, httpAuth, auditor, + events, }); - builder.setEventBroker(events); - if (processingExtensions.onProcessingErrorHandler) { builder.subscribe({ onProcessingError: processingExtensions.onProcessingErrorHandler, diff --git a/plugins/catalog-backend/src/service/DefaultRefreshService.test.ts b/plugins/catalog-backend/src/service/DefaultRefreshService.test.ts index f999400b8a..29519e9e6b 100644 --- a/plugins/catalog-backend/src/service/DefaultRefreshService.test.ts +++ b/plugins/catalog-backend/src/service/DefaultRefreshService.test.ts @@ -57,6 +57,7 @@ describe('DefaultRefreshService', () => { database: knex, logger, refreshInterval: () => 100, + events: mockServices.events.mock(), }), catalogDb: new DefaultCatalogDatabase({ database: knex, @@ -162,6 +163,7 @@ describe('DefaultRefreshService', () => { }, createHash: () => createHash('sha1'), pollingIntervalMs: 50, + events: mockServices.events.mock(), }); return engine; diff --git a/plugins/catalog-backend/src/tests/integration.test.ts b/plugins/catalog-backend/src/tests/integration.test.ts index 486d1bfb99..39871583af 100644 --- a/plugins/catalog-backend/src/tests/integration.test.ts +++ b/plugins/catalog-backend/src/tests/integration.test.ts @@ -244,6 +244,7 @@ class TestHarness { const processingDatabase = new DefaultProcessingDatabase({ database: options.db, logger, + events: mockServices.events.mock(), refreshInterval: () => 0.05, }); @@ -273,7 +274,6 @@ class TestHarness { logger, parser: defaultEntityDataParser, policy: EntityPolicies.allOf([]), - legacySingleProcessorValidation: false, }); const stitcher = DefaultStitcher.fromConfig(config, { knex: options.db, @@ -303,6 +303,7 @@ class TestHarness { proxyProgressTracker.reportError(event.unprocessedEntity, event.errors); }, tracker: proxyProgressTracker, + events: mockServices.events.mock(), }); const refresh = new DefaultRefreshService({ database: catalogDatabase }); diff --git a/plugins/catalog-backend/src/tests/performance/getProcessableEntitiesPerformance.test.ts b/plugins/catalog-backend/src/tests/performance/getProcessableEntitiesPerformance.test.ts index e2ed160b42..a1cc1d3238 100644 --- a/plugins/catalog-backend/src/tests/performance/getProcessableEntitiesPerformance.test.ts +++ b/plugins/catalog-backend/src/tests/performance/getProcessableEntitiesPerformance.test.ts @@ -79,6 +79,7 @@ describePerformanceTest('getProcessableEntities', () => { const sut = new DefaultProcessingDatabase({ database: knex, logger: mockServices.logger.mock(), + events: mockServices.events.mock(), refreshInterval: () => 0, });