From b7497a992c6c9453b50c61ff923bc6ddfd918209 Mon Sep 17 00:00:00 2001 From: Johan Haals Date: Mon, 13 Sep 2021 11:34:13 +0200 Subject: [PATCH] catalog-backend: tests for processingEngine refresh Co-authored-by: Patrik Oldsberg Signed-off-by: Johan Haals --- .../DefaultCatalogProcessingEngine.test.ts | 242 +++++++++++++++++- .../next/DefaultCatalogProcessingEngine.ts | 2 + .../DefaultProcessingDatabase.test.ts | 83 ------ .../database/DefaultProcessingDatabase.ts | 2 +- 4 files changed, 243 insertions(+), 86 deletions(-) diff --git a/plugins/catalog-backend/src/next/DefaultCatalogProcessingEngine.test.ts b/plugins/catalog-backend/src/next/DefaultCatalogProcessingEngine.test.ts index f7dede4adb..a0d62c8642 100644 --- a/plugins/catalog-backend/src/next/DefaultCatalogProcessingEngine.test.ts +++ b/plugins/catalog-backend/src/next/DefaultCatalogProcessingEngine.test.ts @@ -15,13 +15,27 @@ */ import { getVoidLogger } from '@backstage/backend-common'; -import { Hash } from 'crypto'; +import { TestDatabaseId, TestDatabases } from '@backstage/backend-test-utils'; +import { createHash, Hash } from 'crypto'; +import { Knex } from 'knex'; +import { Logger } from 'winston'; import { DateTime } from 'luxon'; +import { DatabaseManager } from './database/DatabaseManager'; import waitForExpect from 'wait-for-expect'; import { DefaultProcessingDatabase } from './database/DefaultProcessingDatabase'; +import { + DbRefreshStateReferencesRow, + DbRefreshStateRow, +} from './database/tables'; +import { ProcessingDatabase } from './database/types'; import { DefaultCatalogProcessingEngine } from './DefaultCatalogProcessingEngine'; -import { CatalogProcessingOrchestrator } from './processing/types'; +import { + CatalogProcessingOrchestrator, + EntityProcessingRequest, +} from './processing/types'; import { Stitcher } from './stitching/Stitcher'; +import { Entity, stringifyEntityRef } from '@backstage/catalog-model'; +import { v4 as uuid } from 'uuid'; describe('DefaultCatalogProcessingEngine', () => { const db = { @@ -236,3 +250,227 @@ describe('DefaultCatalogProcessingEngine', () => { await engine.stop(); }); }); + +describe('DefaultCatalogProcessingEngine integration', () => { + const defaultLogger = getVoidLogger(); + const databases = TestDatabases.create({ + ids: ['POSTGRES_13', 'POSTGRES_9', 'SQLITE_3'], + }); + + async function createDatabase( + databaseId: TestDatabaseId, + logger: Logger = defaultLogger, + ) { + const knex = await databases.init(databaseId); + await DatabaseManager.createDatabase(knex); + return { + knex, + db: new DefaultProcessingDatabase({ + database: knex, + logger, + refreshInterval: () => 100, + }), + }; + } + + const createPopulatedEngine = async (options: { + db: ProcessingDatabase; + knex: Knex; + entities: Entity[]; + references: { [source: string]: string[] }; + }) => { + const { db, knex, entities, references } = options; + + const entityMap = new Map( + entities.map(entity => [stringifyEntityRef(entity), entity]), + ); + + for (const entity of entities) { + await knex('refresh_state').insert({ + entity_id: uuid(), + entity_ref: stringifyEntityRef(entity), + unprocessed_entity: JSON.stringify(entity), + errors: '[]', + next_update_at: '2031-01-01 23:00:00', + last_discovery_at: '2021-04-01 13:37:00', + }); + } + + for (const entityRef of entityMap.keys()) { + if (!(entityRef in references)) { + await knex( + 'refresh_state_references', + ).insert({ + source_key: 'ConfigLocationProvider', + target_entity_ref: entityRef, + }); + } + } + for (const [sourceRef, targetRefs] of Object.entries(references)) { + for (const targetRef of targetRefs) { + await knex( + 'refresh_state_references', + ).insert({ + source_entity_ref: sourceRef, + target_entity_ref: targetRef, + }); + } + } + + const engine = new DefaultCatalogProcessingEngine( + defaultLogger, + [], + db, + { + async process(request: EntityProcessingRequest) { + const entityRef = stringifyEntityRef(request.entity); + const entity = entityMap.get(entityRef); + if (!entity) { + throw new Error(`Unexpected entity: ${entityRef}`); + } + const deferredEntities = + references[entityRef]?.map(ref => { + const e = entityMap.get(ref); + if (!e) { + throw new Error(`Target entity not found: ${ref}`); + } + return { entity: e, locationKey: ref }; + }) || []; + + return { + ok: true, + completedEntity: { + ...entity, + metadata: { + ...entity.metadata, + annotations: { + ...entity.metadata.annotations, + 'refresh-completed': 'true', + }, + }, + }, + relations: [], + errors: [], + deferredEntities, + state: new Map(), + }; + }, + }, + new Stitcher(knex, defaultLogger), + () => createHash('sha1'), + 50, + ); + + return engine; + }; + + const waitForRefresh = async (knex: Knex, entityRef: string) => { + for (;;) { + const [result] = await knex('refresh_state') + .where('entity_ref', entityRef) + .select(); + + const entity = result.processed_entity + ? (JSON.parse(result.processed_entity) as Entity) + : undefined; + if (entity?.metadata?.annotations?.['refresh-completed']) { + return true; + } + await new Promise(resolve => setTimeout(resolve, 500)); + } + }; + + it.each(databases.eachSupportedId())( + 'should refresh the parent location, %p', + async databaseId => { + const { knex, db } = await createDatabase(databaseId); + + const engine = await createPopulatedEngine({ + db, + knex, + entities: [ + { + kind: 'Location', + apiVersion: '1.0.0', + metadata: { + name: 'myloc', + }, + }, + { + kind: 'Component', + apiVersion: '1.0.0', + metadata: { + name: 'mycomp', + }, + }, + ], + references: { + 'location:default/myloc': ['component:default/mycomp'], + }, + }); + + await engine.start(); + + await engine.refresh({ + entityRef: 'component:default/mycomp', + }); + + await expect( + waitForRefresh(knex, 'component:default/mycomp'), + ).resolves.toBe(true); + + await engine.stop(); + }, + ); + + it.each(databases.eachSupportedId())( + 'should refresh the location further up the tree, %p', + async databaseId => { + const { knex, db } = await createDatabase(databaseId); + + const engine = await createPopulatedEngine({ + db, + knex, + entities: [ + { + kind: 'Location', + apiVersion: '1.0.0', + metadata: { + name: 'myloc', + }, + }, + { + kind: 'Component', + apiVersion: '1.0.0', + metadata: { + name: 'mycomp', + }, + }, + { + kind: 'Api', + apiVersion: '1.0.0', + metadata: { + name: 'myapi', + }, + }, + ], + references: { + 'location:default/myloc': ['component:default/mycomp'], + 'component:default/mycomp': ['api:default/myapi'], + }, + }); + + await engine.start(); + + await engine.refresh({ + entityRef: 'api:default/myapi', + }); + + await expect(waitForRefresh(knex, 'api:default/myapi')).resolves.toBe( + true, + ); + + await engine.stop(); + }, + ); +}); diff --git a/plugins/catalog-backend/src/next/DefaultCatalogProcessingEngine.ts b/plugins/catalog-backend/src/next/DefaultCatalogProcessingEngine.ts index 0e77e6e9b4..b1cfc33885 100644 --- a/plugins/catalog-backend/src/next/DefaultCatalogProcessingEngine.ts +++ b/plugins/catalog-backend/src/next/DefaultCatalogProcessingEngine.ts @@ -104,6 +104,7 @@ export class DefaultCatalogProcessingEngine implements CatalogProcessingEngine { private readonly orchestrator: CatalogProcessingOrchestrator, private readonly stitcher: Stitcher, private readonly createHash: () => Hash, + private readonly pollingIntervalMs: number = 1000, ) {} async start() { @@ -123,6 +124,7 @@ export class DefaultCatalogProcessingEngine implements CatalogProcessingEngine { this.stopFunc = startTaskPipeline({ lowWatermark: 5, highWatermark: 10, + pollingIntervalMs: this.pollingIntervalMs, loadTasks: async count => { try { const { items } = await this.processingDatabase.transaction( diff --git a/plugins/catalog-backend/src/next/database/DefaultProcessingDatabase.test.ts b/plugins/catalog-backend/src/next/database/DefaultProcessingDatabase.test.ts index e027830601..23829bd33f 100644 --- a/plugins/catalog-backend/src/next/database/DefaultProcessingDatabase.test.ts +++ b/plugins/catalog-backend/src/next/database/DefaultProcessingDatabase.test.ts @@ -1017,87 +1017,4 @@ describe('Default Processing Database', () => { 60_000, ); }); - - describe('refreshEntities', () => { - it.each(databases.eachSupportedId())( - 'should refresh the location further up the tree, %p', - async databaseId => { - const { knex, db } = await createDatabase(databaseId); - - await knex('refresh_state').insert({ - entity_id: '7', - entity_ref: 'location:default/myloc', - unprocessed_entity: JSON.stringify({ - kind: 'Location', - apiVersion: '1.0.0', - metadata: { - name: 'xyz', - }, - } as Entity), - errors: '[]', - next_update_at: '2031-01-01 23:00:00', - last_discovery_at: '2021-04-01 13:37:00', - }); - - await knex('refresh_state').insert({ - entity_id: '8', - entity_ref: 'component:default/mycomp', - unprocessed_entity: JSON.stringify({ - kind: 'Component', - apiVersion: '1.0.0', - metadata: { - name: 'xyz', - }, - } as Entity), - errors: '[]', - next_update_at: '2031-01-01 23:00:00', - last_discovery_at: '2021-04-01 13:37:00', - }); - - await knex('refresh_state').insert({ - entity_id: '9', - entity_ref: 'api:default/myapi', - unprocessed_entity: JSON.stringify({ - kind: 'Api', - apiVersion: '1.0.0', - metadata: { - name: 'xyz', - }, - } as Entity), - errors: '[]', - next_update_at: '2031-01-01 23:00:00', - last_discovery_at: '2021-04-01 13:37:00', - }); - - await insertRefRow(knex, { - source_entity_ref: 'component:default/mycomp', - target_entity_ref: 'api:default/myapi', - }); - await insertRefRow(knex, { - source_entity_ref: 'location:default/myloc', - target_entity_ref: 'component:default/mycomp', - }); - - await insertRefRow(knex, { - source_key: 'ConfigLocationProvider', - target_entity_ref: 'location:default/myloc', - }); - - await db.transaction(async tx => { - await db.refreshUnprocessedEntities(tx, { - entityRef: 'api:default/myapi', - }); - }); - - const [result] = await knex('refresh_state') - .where('entity_ref', 'api:default/myapi') - .select(); - - // TODO: This is going to break after 2031 - expect(parseDate(result.next_update_at).year).toEqual( - DateTime.local().year, - ); - }, - ); - }); }); diff --git a/plugins/catalog-backend/src/next/database/DefaultProcessingDatabase.ts b/plugins/catalog-backend/src/next/database/DefaultProcessingDatabase.ts index a19f88188d..369deb36a4 100644 --- a/plugins/catalog-backend/src/next/database/DefaultProcessingDatabase.ts +++ b/plugins/catalog-backend/src/next/database/DefaultProcessingDatabase.ts @@ -514,7 +514,7 @@ export class DefaultProcessingDatabase implements ProcessingDatabase { }; } - // TODO(jhaals): Rename this to refreshEntity? + // TODO(jhaals): Rename this to refreshEntity/refreshEntities? async refreshUnprocessedEntities( txOpaque: Transaction, options: RefreshUnprocessedEntitiesOptions,