diff --git a/plugins/catalog-backend/api-report.md b/plugins/catalog-backend/api-report.md index 3ac5094b31..ad898d2e22 100644 --- a/plugins/catalog-backend/api-report.md +++ b/plugins/catalog-backend/api-report.md @@ -145,6 +145,10 @@ export class CatalogBuilder { processingInterval: ProcessingIntervalFunction, ): CatalogBuilder; setProcessingIntervalSeconds(seconds: number): CatalogBuilder; + // @alpha (undocumented) + subscribe( + catalogProcessingErrorListeners: CatalogProcessingErrorListener[], + ): void; } // @alpha @@ -202,15 +206,13 @@ export type CatalogPermissionRule = // @public (undocumented) export interface CatalogProcessingEngine { - // (undocumented) - addErrorListener?(errorListener: CatalogProcessingErrorListener): void; // (undocumented) start(): Promise; // (undocumented) stop(): Promise; } -// @public +// @alpha export interface CatalogProcessingErrorListener { // (undocumented) onError( diff --git a/plugins/catalog-backend/src/processing/DefaultCatalogProcessingEngine.ts b/plugins/catalog-backend/src/processing/DefaultCatalogProcessingEngine.ts index 7ebb55273b..69295a0662 100644 --- a/plugins/catalog-backend/src/processing/DefaultCatalogProcessingEngine.ts +++ b/plugins/catalog-backend/src/processing/DefaultCatalogProcessingEngine.ts @@ -34,7 +34,6 @@ const CACHE_TTL = 5; export class DefaultCatalogProcessingEngine implements CatalogProcessingEngine { private readonly tracker = progressTracker(); - private readonly errorListeners: CatalogProcessingErrorListener[] = []; private stopFunc?: () => void; constructor( @@ -44,6 +43,7 @@ export class DefaultCatalogProcessingEngine implements CatalogProcessingEngine { private readonly stitcher: Stitcher, private readonly createHash: () => Hash, private readonly pollingIntervalMs: number = 1000, + private readonly catalogProcessingErrorListener?: CatalogProcessingErrorListener, ) {} async start() { @@ -157,9 +157,11 @@ export class DefaultCatalogProcessingEngine implements CatalogProcessingEngine { // just store the errors and trigger a stich so that they become visible to // the outside. if (!result.ok) { - // notify the error listeners if the entity can not be processed. - this.errorListeners.forEach(listener => - listener.onError(unprocessedEntity, result, resultHash), + // notify the error listener if the entity can not be processed. + this.catalogProcessingErrorListener?.onError( + unprocessedEntity, + result, + resultHash, ); await this.processingDatabase.transaction(async tx => { @@ -231,10 +233,6 @@ export class DefaultCatalogProcessingEngine implements CatalogProcessingEngine { this.stopFunc = undefined; } } - - addErrorListener(errorListener: CatalogProcessingErrorListener) { - this.errorListeners.push(errorListener); - } } // Helps wrap the timing and logging behaviors diff --git a/plugins/catalog-backend/src/processing/types.ts b/plugins/catalog-backend/src/processing/types.ts index 54118f8422..e6c58cf36f 100644 --- a/plugins/catalog-backend/src/processing/types.ts +++ b/plugins/catalog-backend/src/processing/types.ts @@ -66,13 +66,12 @@ export type DeferredEntity = { export interface CatalogProcessingEngine { start(): Promise; stop(): Promise; - addErrorListener?(errorListener: CatalogProcessingErrorListener): void; } /** * An error listener for catalog processing engine. It can be used to listen and track entity errors. * - * @public + * @alpha */ export interface CatalogProcessingErrorListener { onError( diff --git a/plugins/catalog-backend/src/service/CatalogBuilder.ts b/plugins/catalog-backend/src/service/CatalogBuilder.ts index 0d4ddc8e71..be427998fc 100644 --- a/plugins/catalog-backend/src/service/CatalogBuilder.ts +++ b/plugins/catalog-backend/src/service/CatalogBuilder.ts @@ -56,7 +56,10 @@ import { } from '../modules/core/PlaceholderProcessor'; import { defaultEntityDataParser } from '../modules/util/parse'; import { LocationAnalyzer } from '../ingestion/types'; -import { CatalogProcessingEngine } from '../processing/types'; +import { + CatalogProcessingEngine, + CatalogProcessingErrorListener, +} from '../processing'; import { DefaultProcessingDatabase } from '../database/DefaultProcessingDatabase'; import { applyDatabaseMigrations } from '../database/migrations'; import { DefaultCatalogProcessingEngine } from '../processing/DefaultCatalogProcessingEngine'; @@ -133,6 +136,7 @@ export class CatalogBuilder { private processors: CatalogProcessor[]; private processorsReplace: boolean; private parser: CatalogProcessorParser | undefined; + private catalogProcessingErrorListeners?: CatalogProcessingErrorListener[]; private processingInterval: ProcessingIntervalFunction = createRandomProcessingInterval({ minSeconds: 100, @@ -447,6 +451,14 @@ export class CatalogBuilder { orchestrator, stitcher, () => createHash('sha1'), + 1000, + { + onError: async (unprocessedEntity, result, resultHash) => { + this.catalogProcessingErrorListeners?.forEach(listener => + listener.onError(unprocessedEntity, result, resultHash), + ); + }, + }, ); const locationAnalyzer = @@ -478,6 +490,14 @@ export class CatalogBuilder { }; } + /** + * @alpha + * @param catalogProcessingErrorListeners - a list of listeners to get notified if an error occurs while processing an entity + */ + subscribe(catalogProcessingErrorListeners: CatalogProcessingErrorListener[]) { + this.catalogProcessingErrorListeners = catalogProcessingErrorListeners; + } + private buildEntityPolicy(): EntityPolicy { const entityPolicies: EntityPolicy[] = this.entityPoliciesReplace ? [new SchemaValidEntityPolicy(), ...this.entityPolicies]