Merge pull request #15922 from thefrontside/tm/incremental-ingestion-return-emittor
Return event subscriber from addIncrementalEntityProvider
This commit is contained in:
@@ -0,0 +1,5 @@
|
||||
---
|
||||
'@backstage/plugin-catalog-backend-module-incremental-ingestion': minor
|
||||
---
|
||||
|
||||
Return EventSubscriber from addIncrementalEntityProvider to hook up to EventsBackend
|
||||
@@ -11,6 +11,7 @@ import type { Config } from '@backstage/config';
|
||||
import type { DeferredEntity } from '@backstage/plugin-catalog-backend';
|
||||
import type { DurationObjectUnits } from 'luxon';
|
||||
import { EventParams } from '@backstage/plugin-events-node';
|
||||
import { EventSubscriber } from '@backstage/plugin-events-node';
|
||||
import type { Logger } from 'winston';
|
||||
import type { PermissionEvaluator } from '@backstage/plugin-permission-common';
|
||||
import type { PluginDatabaseManager } from '@backstage/backend-common';
|
||||
@@ -37,7 +38,7 @@ export class IncrementalCatalogBuilder {
|
||||
addIncrementalEntityProvider<TCursor, TContext>(
|
||||
provider: IncrementalEntityProvider<TCursor, TContext>,
|
||||
options: IncrementalEntityProviderOptions,
|
||||
): void;
|
||||
): EventSubscriber;
|
||||
// (undocumented)
|
||||
build(): Promise<{
|
||||
incrementalAdminRouter: Router;
|
||||
|
||||
+12
-2
@@ -26,6 +26,7 @@ import { applyDatabaseMigrations } from '../database/migrations';
|
||||
import { IncrementalIngestionDatabaseManager } from '../database/IncrementalIngestionDatabaseManager';
|
||||
import { IncrementalProviderRouter } from '../router/routes';
|
||||
import { Deferred } from '../util';
|
||||
import { EventParams, EventSubscriber } from '@backstage/plugin-events-node';
|
||||
|
||||
/** @public */
|
||||
export class IncrementalCatalogBuilder {
|
||||
@@ -71,13 +72,15 @@ export class IncrementalCatalogBuilder {
|
||||
addIncrementalEntityProvider<TCursor, TContext>(
|
||||
provider: IncrementalEntityProvider<TCursor, TContext>,
|
||||
options: IncrementalEntityProviderOptions,
|
||||
) {
|
||||
): EventSubscriber {
|
||||
const { burstInterval, burstLength, restLength } = options;
|
||||
const { logger: catalogLogger, scheduler } = this.env;
|
||||
const ready = this.ready;
|
||||
|
||||
const manager = this.manager;
|
||||
|
||||
let engine: IncrementalIngestionEngine;
|
||||
|
||||
this.builder.addEntityProvider({
|
||||
getProviderName: provider.getProviderName.bind(provider),
|
||||
async connect(connection) {
|
||||
@@ -87,7 +90,7 @@ export class IncrementalCatalogBuilder {
|
||||
|
||||
logger.info(`Connecting`);
|
||||
|
||||
const engine = new IncrementalIngestionEngine({
|
||||
engine = new IncrementalIngestionEngine({
|
||||
...options,
|
||||
ready,
|
||||
manager,
|
||||
@@ -112,5 +115,12 @@ export class IncrementalCatalogBuilder {
|
||||
});
|
||||
},
|
||||
});
|
||||
|
||||
return {
|
||||
onEvent: (params: EventParams) => engine.onEvent(params),
|
||||
supportsEventTopics() {
|
||||
return engine.supportsEventTopics();
|
||||
},
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user