diff --git a/plugins/catalog-backend/src/next/CatalogProcessingEngineImpl.ts b/plugins/catalog-backend/src/next/CatalogProcessingEngineImpl.ts index e94cdf2f31..8a2a0349d6 100644 --- a/plugins/catalog-backend/src/next/CatalogProcessingEngineImpl.ts +++ b/plugins/catalog-backend/src/next/CatalogProcessingEngineImpl.ts @@ -23,7 +23,6 @@ import { CatalogProcessingOrchestrator, } from './types'; -import { JsonObject } from '@backstage/config'; import { EntitiesCatalog } from '../catalog/types'; import { Logger } from 'winston'; @@ -50,42 +49,33 @@ export class CatalogProcessingEngineImpl implements CatalogProcessingEngine { this.running = true; while (this.running) { - const { entity, state, eager } = await this.stateManager.pop(); + const { + id, + entity, + state: intialState, + } = await this.stateManager.getNextProccessingItem(); const { - completeEntities, + completedEntity, deferredEntites, errors, + state, } = await this.orchestrator.process({ entity, - state, - eager, + state: intialState, }); for (const error of errors) { this.logger.warn(error.message); } - for (const deferred of deferredEntites) { - this.stateManager.setResult({ - request: { entity: deferred, state: new Map() }, - nextRefresh: 'now()', - }); - } - - if (completeEntities.length) { - await this.entitiesCatalog.batchAddOrUpdateEntities( - completeEntities.map(e => ({ - entity: e, - relations: [], - })), - { - locationId: 'xyz', - dryRun: false, - outputEntities: true, - }, - ); - } + await this.stateManager.setProcessingItemResult({ + id, + entity: completedEntity, + state, + errors, + }); + await this.stateManager.addProcessingItems({ entities: deferredEntites }); } } @@ -96,28 +86,19 @@ export class CatalogProcessingEngineImpl implements CatalogProcessingEngine { subscription.unsubscribe(); } } + private async onNext(message: EntityMessage) { if ('all' in message) { // TODO unhandled rejection - await Promise.all( - message.all.map(entity => - this.stateManager.setResult({ - request: { entity, state: new Map() }, - nextRefresh: 'now()', - }), - ), - ); + await this.stateManager.addProcessingItems({ + entities: message.all, + }); } if ('added' in message) { - await Promise.all( - message.added.map(entity => - this.stateManager.setResult({ - request: { entity, state: new Map() }, - nextRefresh: 'now()', - }), - ), - ); + await this.stateManager.addProcessingItems({ + entities: message.added, + }); // TODO deletions of message.removed } diff --git a/plugins/catalog-backend/src/next/CatalogProcessingOrchestratorImpl.ts b/plugins/catalog-backend/src/next/CatalogProcessingOrchestratorImpl.ts index f78ebfbef2..add7db53b3 100644 --- a/plugins/catalog-backend/src/next/CatalogProcessingOrchestratorImpl.ts +++ b/plugins/catalog-backend/src/next/CatalogProcessingOrchestratorImpl.ts @@ -19,14 +19,18 @@ import { LocationSpec, stringifyLocationReference, } from '@backstage/catalog-model'; -import { JsonObject } from '@backstage/config'; import { CatalogProcessor, CatalogProcessorEmit, CatalogProcessorParser, CatalogProcessorResult, } from '../ingestion/processors'; -import { CatalogProcessingOrchestrator, EntityProcessingError } from './types'; +import { + CatalogProcessingOrchestrator, + EntityProcessingError, + EntityProcessingRequest, + EntityProcessingResult, +} from './types'; import { Logger } from 'winston'; import * as result from '../ingestion/processors/results'; import { locationToEntity } from './LocationToEntity'; @@ -41,18 +45,10 @@ export class CatalogProcessingOrchestratorImpl }, ) {} - async process(request: { - entity: Entity; - eager?: boolean | undefined; - state: Map; - }): Promise<{ - state: Map; - completeEntities: Entity[]; - deferredEntites: Entity[]; - errors: EntityProcessingError[]; - }> { + async process( + request: EntityProcessingRequest, + ): Promise { const { entity, eager, state } = request; - const completedEntities: Entity[] = []; const deferredEntites: Entity[] = []; const errors: EntityProcessingError[] = []; @@ -89,7 +85,7 @@ export class CatalogProcessingOrchestratorImpl return { deferredEntites, - completeEntities: completedEntities, + completedEntity, errors, state, }; diff --git a/plugins/catalog-backend/src/next/ProcessingStateManagerImpl.ts b/plugins/catalog-backend/src/next/ProcessingStateManagerImpl.ts index d8736cfa77..872ee78fdb 100644 --- a/plugins/catalog-backend/src/next/ProcessingStateManagerImpl.ts +++ b/plugins/catalog-backend/src/next/ProcessingStateManagerImpl.ts @@ -14,52 +14,61 @@ * limitations under the License. */ -import { Entity } from '@backstage/catalog-model'; -import { JsonObject } from '@backstage/config'; -import { DbRefreshStateRow, ProcessingDatabase } from './database/types'; -import { ProcessingResult, ProcessingStateManager } from './types'; +import { ProcessingDatabase, RefreshStateItem } from './database/types'; +import { + AddProcessingItemRequest, + ProccessingItem, + ProcessingItemResult, + ProcessingStateManager, +} from './types'; export class ProcessingStateManagerImpl implements ProcessingStateManager { constructor(private readonly db: ProcessingDatabase) {} - async setResult(result: ProcessingResult) { - await this.db.transaction(async tx => { - this.db.addEntityRefreshState(tx, [ - { - entity: result.request.entity, - nextRefresh: result.nextRefresh, - }, - ]); + async setProcessingItemResult(result: ProcessingItemResult) { + return this.db.transaction(async tx => { + await this.db.updateProcessedEntity(tx, { + id: result.id, + processedEntity: result.entity, + errors: JSON.stringify(result.errors), + state: result.state, + }); }); } - async pop(): Promise<{ - entity: Entity; - eager?: boolean | undefined; - state: Map; - }> { - const entities = await new Promise(resolve => + async addProcessingItems(request: AddProcessingItemRequest) { + return this.db.transaction(async tx => { + await this.db.addUnprocessedEntities(tx, { + unprocessedEntities: request.entities, + }); + }); + } + + async getNextProccessingItem(): Promise { + const entities = await new Promise(resolve => this.popFromQueue(resolve), ); - const result = entities[0]; + const { id, state, unprocessedEntity } = entities[0]; return { - entity: JSON.parse(result.unproccessed_entity) as Entity, - state: new Map(JSON.parse(result.cache)), + id, + entity: unprocessedEntity, + state, }; } - async popFromQueue(resolve: (rows: DbRefreshStateRow[]) => void) { + async popFromQueue(resolve: (rows: RefreshStateItem[]) => void) { const entities = await this.db.transaction(async tx => { return this.db.getProcessableEntities(tx, { processBatchSize: 1, }); }); - if (!entities.length) { + // No entities require refresh, wait and try again. + if (!entities.items.length) { setTimeout(() => this.popFromQueue(resolve), 1000); return; } - resolve(entities); + resolve(entities.items); } } diff --git a/plugins/catalog-backend/src/next/database/ProcessingDatabaseImpl.ts b/plugins/catalog-backend/src/next/database/ProcessingDatabaseImpl.ts index af0d331199..3fe25c1eae 100644 --- a/plugins/catalog-backend/src/next/database/ProcessingDatabaseImpl.ts +++ b/plugins/catalog-backend/src/next/database/ProcessingDatabaseImpl.ts @@ -21,10 +21,10 @@ import { ProcessingDatabase, AddUnprocessedEntitiesOptions, UpdateProcessedEntityOptions, - GetProcessedEntitiesResult, + GetProcessableEntitiesResult, } from './types'; import type { Logger } from 'winston'; -import { Entity } from '@backstage/catalog-model'; +import { Entity, stringifyEntityRef } from '@backstage/catalog-model'; import { v4 as uuid } from 'uuid'; export type DbRefreshStateRequest = { @@ -54,11 +54,11 @@ class ProcessingDatabaseImpl implements ProcessingDatabase { options: UpdateProcessedEntityOptions, ): Promise { const tx = txOpaque as Knex.Transaction; - const { id, processedEntity, cache, errors } = options; + const { id, processedEntity, state, errors } = options; const result = await tx('refresh_state') .update({ - processed_entity: processedEntity, - cache, + processed_entity: JSON.stringify(processedEntity), + cache: JSON.stringify(state), errors, }) .where('id', id); @@ -93,7 +93,7 @@ class ProcessingDatabaseImpl implements ProcessingDatabase { async getProcessableEntities( txOpaque: Transaction, request: { processBatchSize: number }, - ): Promise { + ): Promise { const tx = txOpaque as Knex.Transaction; const items = await tx('refresh_state') @@ -114,7 +114,18 @@ class ProcessingDatabaseImpl implements ProcessingDatabase { : tx.raw(`now() + interval '30 seconds'`), }); - return items; + return { + items: items.map(i => ({ + id: i.id, + entityRef: i.entity_ref, + unprocessedEntity: JSON.parse(i.unprocessed_entity) as Entity, + processedEntity: JSON.parse(i.processed_entity) as Entity, + nextUpdateAt: i.next_update_at, + lastDiscoveryAt: i.last_discovery_at, + state: JSON.parse(i.cache), + errors: i.errors, + })), + }; } async transaction(fn: (tx: Transaction) => Promise): Promise { @@ -148,9 +159,3 @@ class ProcessingDatabaseImpl implements ProcessingDatabase { } } } - -function stringifyEntityRef( - entity: Entity, -): Knex.MaybeRawColumn | undefined { - throw new Error('Function not implemented.'); -} diff --git a/plugins/catalog-backend/src/next/database/types.ts b/plugins/catalog-backend/src/next/database/types.ts index 310b6b4b7e..6cd42e690f 100644 --- a/plugins/catalog-backend/src/next/database/types.ts +++ b/plugins/catalog-backend/src/next/database/types.ts @@ -26,24 +26,24 @@ export type AddUnprocessedEntitiesResult = {}; export type UpdateProcessedEntityOptions = { id: string; - processedEntity?: string; - cache?: string; + processedEntity?: Entity; + state?: Map; errors?: string; }; export type RefreshStateItem = { id: string; entityRef: string; - unprocessedEntity: string; - processedEntity: string; + unprocessedEntity: Entity; + processedEntity: Entity; nextUpdateAt: string; lastDiscoveryAt: string; // remove? - cache: JsonObject; + state: Map; errors: string; }; -export type GetProcessedEntitiesResult = { - items: RefreshStateItem; +export type GetProcessableEntitiesResult = { + items: RefreshStateItem[]; }; export interface ProcessingDatabase { @@ -56,7 +56,7 @@ export interface ProcessingDatabase { getProcessableEntities( txOpaque: Transaction, request: { processBatchSize: number }, - ): Promise; + ): Promise; /** * Updates the diff --git a/plugins/catalog-backend/src/next/types.ts b/plugins/catalog-backend/src/next/types.ts index c7f1e370e0..58d88636f8 100644 --- a/plugins/catalog-backend/src/next/types.ts +++ b/plugins/catalog-backend/src/next/types.ts @@ -69,8 +69,7 @@ export interface EntityProvider { entityChange$(): Observable; } -// interface CatalogProcessor {} -type EntityProcessingRequest = { +export type EntityProcessingRequest = { entity: Entity; eager?: boolean; state: Map; // Versions for multiple deployments etc @@ -81,9 +80,9 @@ export type EntityProcessingError = { message: String; }; -type EntityProcessingResult = { +export type EntityProcessingResult = { state: Map; - completeEntities: Entity[]; + completedEntity: Entity; deferredEntites: Entity[]; errors: EntityProcessingError[]; }; @@ -92,12 +91,24 @@ export interface CatalogProcessingOrchestrator { process(request: EntityProcessingRequest): Promise; } -export type ProcessingResult = { - nextRefresh: string; - request: EntityProcessingRequest; +export type ProcessingItemResult = { + id: string; + entity: Entity; + state: Map; + errors: EntityProcessingError[]; }; +export type AddProcessingItemRequest = { + entities: Entity[]; +}; + +export type ProccessingItem = { + id: string; + entity: Entity; + state: Map; +}; export interface ProcessingStateManager { - setResult(result: ProcessingResult): Promise; - pop(): Promise; + setProcessingItemResult(result: ProcessingItemResult): Promise; + getNextProccessingItem(): Promise; + addProcessingItems(request: AddProcessingItemRequest): Promise; }