diff --git a/.changeset/shaggy-spiders-notice.md b/.changeset/shaggy-spiders-notice.md new file mode 100644 index 0000000000..513f6cd521 --- /dev/null +++ b/.changeset/shaggy-spiders-notice.md @@ -0,0 +1,14 @@ +--- +'@backstage/plugin-catalog-backend': patch +--- + +CatalogBuilder supports now subscription to processing engine errors. + +```ts +subscribe(options: { + onProcessingError: (event: { unprocessedEntity: Entity, error: Error }) => Promise | void; +}); +``` + +If you want to get notified on errors while processing the entities, you call CatalogBuilder.subscribe +to get notifications with the parameters defined as above. diff --git a/plugins/catalog-backend/api-report.md b/plugins/catalog-backend/api-report.md index 09ca4d8dfe..d9c2842b24 100644 --- a/plugins/catalog-backend/api-report.md +++ b/plugins/catalog-backend/api-report.md @@ -144,6 +144,13 @@ export class CatalogBuilder { processingInterval: ProcessingIntervalFunction, ): CatalogBuilder; setProcessingIntervalSeconds(seconds: number): CatalogBuilder; + // (undocumented) + subscribe(options: { + onProcessingError: (event: { + unprocessedEntity: Entity; + errors: Error[]; + }) => Promise | void; + }): void; } // @alpha diff --git a/plugins/catalog-backend/src/processing/DefaultCatalogProcessingEngine.ts b/plugins/catalog-backend/src/processing/DefaultCatalogProcessingEngine.ts index 82062a8e2c..fb9518c783 100644 --- a/plugins/catalog-backend/src/processing/DefaultCatalogProcessingEngine.ts +++ b/plugins/catalog-backend/src/processing/DefaultCatalogProcessingEngine.ts @@ -14,8 +14,8 @@ * limitations under the License. */ -import { stringifyEntityRef } from '@backstage/catalog-model'; -import { assertError, serializeError } from '@backstage/errors'; +import { Entity, stringifyEntityRef } from '@backstage/catalog-model'; +import { assertError, serializeError, stringifyError } from '@backstage/errors'; import { Hash } from 'crypto'; import stableStringify from 'fast-json-stable-stringify'; import { Logger } from 'winston'; @@ -25,7 +25,7 @@ import { CatalogProcessingEngine, CatalogProcessingOrchestrator, EntityProcessingResult, -} from '../processing/types'; +} from './types'; import { Stitcher } from '../stitching/Stitcher'; import { startTaskPipeline } from './TaskPipeline'; @@ -42,6 +42,10 @@ export class DefaultCatalogProcessingEngine implements CatalogProcessingEngine { private readonly stitcher: Stitcher, private readonly createHash: () => Hash, private readonly pollingIntervalMs: number = 1000, + private readonly onProcessingError?: (event: { + unprocessedEntity: Entity; + errors: Error[]; + }) => Promise | void, ) {} async start() { @@ -122,6 +126,7 @@ export class DefaultCatalogProcessingEngine implements CatalogProcessingEngine { ); let hashBuilder = this.createHash().update(errorsString); + if (result.ok) { const { entityRefs: parents } = await this.processingDatabase.transaction(tx => @@ -155,6 +160,22 @@ 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 listener if the entity can not be processed. + Promise.resolve(undefined) + .then(() => + this.onProcessingError?.({ + unprocessedEntity, + errors: result.errors, + }), + ) + .catch(error => { + this.logger.debug( + `Processing error listener threw an exception, ${stringifyError( + error, + )}`, + ); + }); + await this.processingDatabase.transaction(async tx => { await this.processingDatabase.updateProcessedEntityErrors(tx, { id, diff --git a/plugins/catalog-backend/src/processing/types.ts b/plugins/catalog-backend/src/processing/types.ts index e94125b6f5..b6790db812 100644 --- a/plugins/catalog-backend/src/processing/types.ts +++ b/plugins/catalog-backend/src/processing/types.ts @@ -28,7 +28,7 @@ export type EntityProcessingRequest = { }; /** * The result of processing an entity. - * @public + * @internal */ export type EntityProcessingResult = | { diff --git a/plugins/catalog-backend/src/service/CatalogBuilder.ts b/plugins/catalog-backend/src/service/CatalogBuilder.ts index 0d4ddc8e71..390dc817d1 100644 --- a/plugins/catalog-backend/src/service/CatalogBuilder.ts +++ b/plugins/catalog-backend/src/service/CatalogBuilder.ts @@ -17,6 +17,7 @@ import { PluginDatabaseManager, UrlReader } from '@backstage/backend-common'; import { DefaultNamespaceEntityPolicy, + Entity, EntityPolicies, EntityPolicy, FieldFormatEntityPolicy, @@ -56,7 +57,7 @@ import { } from '../modules/core/PlaceholderProcessor'; import { defaultEntityDataParser } from '../modules/util/parse'; import { LocationAnalyzer } from '../ingestion/types'; -import { CatalogProcessingEngine } from '../processing/types'; +import { CatalogProcessingEngine } from '../processing'; import { DefaultProcessingDatabase } from '../database/DefaultProcessingDatabase'; import { applyDatabaseMigrations } from '../database/migrations'; import { DefaultCatalogProcessingEngine } from '../processing/DefaultCatalogProcessingEngine'; @@ -133,6 +134,10 @@ export class CatalogBuilder { private processors: CatalogProcessor[]; private processorsReplace: boolean; private parser: CatalogProcessorParser | undefined; + private onProcessingError?: (event: { + unprocessedEntity: Entity; + errors: Error[]; + }) => Promise | void; private processingInterval: ProcessingIntervalFunction = createRandomProcessingInterval({ minSeconds: 100, @@ -447,6 +452,10 @@ export class CatalogBuilder { orchestrator, stitcher, () => createHash('sha1'), + 1000, + event => { + this.onProcessingError?.(event); + }, ); const locationAnalyzer = @@ -478,6 +487,15 @@ export class CatalogBuilder { }; } + subscribe(options: { + onProcessingError: (event: { + unprocessedEntity: Entity; + errors: Error[]; + }) => Promise | void; + }) { + this.onProcessingError = options.onProcessingError; + } + private buildEntityPolicy(): EntityPolicy { const entityPolicies: EntityPolicy[] = this.entityPoliciesReplace ? [new SchemaValidEntityPolicy(), ...this.entityPolicies]