diff --git a/plugins/catalog-backend/src/next/CatalogProcessingEngineImpl.ts b/plugins/catalog-backend/src/next/CatalogProcessingEngineImpl.ts index 8a2a0349d6..93616b6399 100644 --- a/plugins/catalog-backend/src/next/CatalogProcessingEngineImpl.ts +++ b/plugins/catalog-backend/src/next/CatalogProcessingEngineImpl.ts @@ -55,27 +55,28 @@ export class CatalogProcessingEngineImpl implements CatalogProcessingEngine { state: intialState, } = await this.stateManager.getNextProccessingItem(); - const { - completedEntity, - deferredEntites, - errors, - state, - } = await this.orchestrator.process({ + const result = await this.orchestrator.process({ entity, state: intialState, }); - for (const error of errors) { + for (const error of result.errors) { this.logger.warn(error.message); } + if (!result.ok) { + return; + } + await this.stateManager.setProcessingItemResult({ id, - entity: completedEntity, - state, - errors, + entity: result.completedEntity, + state: result.state, + errors: result.errors, + }); + await this.stateManager.addProcessingItems({ + entities: result.deferredEntites, }); - await this.stateManager.addProcessingItems({ entities: deferredEntites }); } } diff --git a/plugins/catalog-backend/src/next/CatalogProcessingOrchestratorImpl.ts b/plugins/catalog-backend/src/next/CatalogProcessingOrchestratorImpl.ts index add7db53b3..1288561b07 100644 --- a/plugins/catalog-backend/src/next/CatalogProcessingOrchestratorImpl.ts +++ b/plugins/catalog-backend/src/next/CatalogProcessingOrchestratorImpl.ts @@ -16,24 +16,23 @@ import { Entity, - LocationSpec, - stringifyLocationReference, + EntityRelationSpec, + stringifyEntityRef, + LOCATION_ANNOTATION, + EntityPolicy, } from '@backstage/catalog-model'; import { CatalogProcessor, - CatalogProcessorEmit, CatalogProcessorParser, CatalogProcessorResult, } from '../ingestion/processors'; import { CatalogProcessingOrchestrator, - EntityProcessingError, EntityProcessingRequest, EntityProcessingResult, } from './types'; import { Logger } from 'winston'; -import * as result from '../ingestion/processors/results'; -import { locationToEntity } from './LocationToEntity'; +import { InputError } from '@backstage/errors'; export class CatalogProcessingOrchestratorImpl implements CatalogProcessingOrchestrator { @@ -42,6 +41,7 @@ export class CatalogProcessingOrchestratorImpl processors: CatalogProcessor[]; logger: Logger; parser: CatalogProcessorParser; + policy: EntityPolicy; }, ) {} @@ -49,101 +49,144 @@ export class CatalogProcessingOrchestratorImpl request: EntityProcessingRequest, ): Promise { const { entity, eager, state } = request; - const deferredEntites: Entity[] = []; - const errors: EntityProcessingError[] = []; - if (eager) { - const stack = [entity]; - const emit = (i: CatalogProcessorResult) => stack.push(i); - while (stack.length) { - const item = stack.pop(); - stack.push(); - } - } else { - const emit = (i: CatalogProcessorResult) => { - if (i.type === 'entity') { - deferredEntites.push(i.entity); - } - if (i.type === 'location') { - deferredEntites.push( - locationToEntity(i.location.type, i.location.type), - ); - } - }; - if (entity.spec) { - if (entity.spec.kind === 'Location') { - this.handleLocation( - { - type: (entity.spec.type as unknown) as string, - target: (entity.spec.target as unknown) as string, - }, - emit, - ); - } - } - } + const result = await this.processSingleEntity(entity); - return { - deferredEntites, - completedEntity, - errors, - state, - }; + console.log('result', JSON.stringify(result)); + return result; } - private async handleLocation( - location: LocationSpec, - emit: CatalogProcessorEmit, - ) { - const { processors, logger } = this.options; + private async processSingleEntity( + unprocessedEntity: Entity, + ): Promise { + // TODO: validate that this doesn't change during processing + const entityRef = stringifyEntityRef(unprocessedEntity); + // TODO: which one do we actually use here? source-location? - maybe probably doesn't exist yet? + const locationRef = + unprocessedEntity.metadata?.annotations?.[LOCATION_ANNOTATION]; - const validatedEmit: CatalogProcessorEmit = emitResult => { - if (emitResult.type === 'relation') { - throw new Error('readLocation may not emit entity relations'); - } - if ( - emitResult.type === 'location' && - emitResult.location.type === location.type && - emitResult.location.target === location.target - ) { - // Ignore self-referential locations silently (this can happen for - // example if you use a glob target like "**/*.yaml" in a Location - // entity) - return; - } - emit(emitResult); - }; + const emitter = createEmitter(this.options.logger); - for (const processor of processors) { - if (processor.readLocation) { - try { - if ( - await processor.readLocation( - location, - // TODO: figure out how this will work in the new processing system - locations from PRs might be optional - location.presence !== 'required', - validatedEmit, - this.options.parser, - ) - ) { - return; + try { + // Pre-process phase, used to populate entities with data that is required during main processing step + let entity = unprocessedEntity; + for (const processor of this.options.processors) { + if (processor.preProcessEntity) { + try { + entity = await processor.preProcessEntity({ + entity, + emit: emitter.emit, + }); + } catch (e) { + throw new Error( + `Processor ${processor.constructor.name} threw an error while preprocessing entity ${entityRef} at ${locationRef}, ${e}`, + ); } - } catch (e) { - const message = `Processor ${ - processor.constructor.name - } threw an error while reading location ${stringifyLocationReference( - location, - )}, ${e}`; - emit(result.generalError(location, message)); - logger.warn(message); } } - } - const message = `No processor was able to read location ${stringifyLocationReference( - location, - )}`; - emit(result.inputError(location, message)); - logger.warn(message); + // Enforce entity policies making sure that entities conform to a general schema + let policyEnforcedEntity; + try { + policyEnforcedEntity = await this.options.policy.enforce(entity); + } catch (e) { + throw new InputError( + `Policy check failed while analyzing entity ${entityRef} at ${locationRef}, ${e}`, + ); + } + if (!policyEnforcedEntity) { + throw new Error( + `Policy unexpectedly returned no data while analyzing entity ${entityRef} at ${locationRef}`, + ); + } + entity = policyEnforcedEntity; + + // Validate the given entity kind against its schema + let handled = false; + for (const processor of this.options.processors) { + if (processor.validateEntityKind) { + try { + handled = await processor.validateEntityKind(entity); + if (handled) { + break; + } + } catch (e) { + throw new InputError( + `Processor ${processor.constructor.name} threw an error while validating the entity ${entityRef} at ${locationRef}, ${e}`, + ); + } + } + } + if (!handled) { + throw new InputError( + `No processor recognized the entity ${entityRef} at ${locationRef}`, + ); + } + + // Main processing step of the entity + for (const processor of this.options.processors) { + if (processor.processEntity) { + try { + entity = await processor.processEntity({ + entity, + emit: emitter.emit, + }); + } catch (e) { + throw new Error( + `Processor ${processor.constructor.name} threw an error while postprocessing entity ${entityRef} at ${locationRef}, ${e}`, + ); + } + } + } + + return { + ...emitter.results(), + completedEntity: entity, + state: new Map(), + ok: true, + }; + } catch (error) { + this.options.logger.warn(error.message); + return { ok: false, errors: emitter.results().errors.concat(error) }; + } } } + +function createEmitter(logger: Logger) { + let done = false; + + const errors = new Array(); + const relations = new Array(); + const deferredEntites = new Array(); + + const emit = (i: CatalogProcessorResult) => { + console.log('CatalogProcessorResult', i); + if (done) { + logger.warn( + `Item if type ${i.type} was emitted after processing had completed at ${ + new Error().stack + }`, + ); + return; + } + if (i.type === 'entity') { + deferredEntites.push(i.entity); + } else if (i.type === 'relation') { + relations.push(i.relation); + } else if (i.type === 'error') { + errors.push(i.error); + } + }; + + return { + emit, + results() { + done = true; + return { + errors, + relations, + deferredEntites, + }; + }, + }; +} diff --git a/plugins/catalog-backend/src/next/LocationServiceImpl.ts b/plugins/catalog-backend/src/next/LocationServiceImpl.ts index 29b41e0a09..e8e8745f01 100644 --- a/plugins/catalog-backend/src/next/LocationServiceImpl.ts +++ b/plugins/catalog-backend/src/next/LocationServiceImpl.ts @@ -50,7 +50,7 @@ export class LocationServiceImpl implements LocationService { return { location: { ...spec, id: `${spec.type}:${spec.target}` }, - entities: processed.completeEntities, + entities: [processed.completedEntity], }; } diff --git a/plugins/catalog-backend/src/next/database/ProcessingDatabaseImpl.ts b/plugins/catalog-backend/src/next/database/ProcessingDatabaseImpl.ts index 3fe25c1eae..25181335c4 100644 --- a/plugins/catalog-backend/src/next/database/ProcessingDatabaseImpl.ts +++ b/plugins/catalog-backend/src/next/database/ProcessingDatabaseImpl.ts @@ -43,7 +43,7 @@ export type DbRefreshStateRow = { errors: string; }; -class ProcessingDatabaseImpl implements ProcessingDatabase { +export class ProcessingDatabaseImpl implements ProcessingDatabase { constructor( private readonly database: Knex, private readonly logger: Logger, @@ -82,6 +82,7 @@ class ProcessingDatabaseImpl implements ProcessingDatabase { id: uuid(), entity_ref: stringifyEntityRef(entity), unprocessed_entity: JSON.stringify(entity), + errors: '', next_update_at: tx.fn.now(), last_discovery_at: tx.fn.now(), }) diff --git a/plugins/catalog-backend/src/next/types.ts b/plugins/catalog-backend/src/next/types.ts index 58d88636f8..091ddc8bcd 100644 --- a/plugins/catalog-backend/src/next/types.ts +++ b/plugins/catalog-backend/src/next/types.ts @@ -75,17 +75,18 @@ export type EntityProcessingRequest = { state: Map; // Versions for multiple deployments etc }; -export type EntityProcessingError = { - // some error stuff here - message: String; -}; - -export type EntityProcessingResult = { - state: Map; - completedEntity: Entity; - deferredEntites: Entity[]; - errors: EntityProcessingError[]; -}; +export type EntityProcessingResult = + | { + ok: true; + state: Map; + completedEntity: Entity; + deferredEntites: Entity[]; + errors: Error[]; + } + | { + ok: false; + errors: Error[]; + }; export interface CatalogProcessingOrchestrator { process(request: EntityProcessingRequest): Promise; @@ -95,7 +96,7 @@ export type ProcessingItemResult = { id: string; entity: Entity; state: Map; - errors: EntityProcessingError[]; + errors: Error[]; }; export type AddProcessingItemRequest = {