catalog-backend: tests for processingEngine refresh

Co-authored-by: Patrik Oldsberg <poldsberg@gmail.com>
Signed-off-by: Johan Haals <johan.haals@gmail.com>
This commit is contained in:
Johan Haals
2021-09-13 11:34:13 +02:00
parent 9c9353e2d3
commit b7497a992c
4 changed files with 243 additions and 86 deletions
@@ -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<DbRefreshStateRow>('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<DbRefreshStateReferencesRow>(
'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<DbRefreshStateReferencesRow>(
'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<DbRefreshStateRow>('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();
},
);
});
@@ -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<RefreshStateItem>({
lowWatermark: 5,
highWatermark: 10,
pollingIntervalMs: this.pollingIntervalMs,
loadTasks: async count => {
try {
const { items } = await this.processingDatabase.transaction(
@@ -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<DbRefreshStateRow>('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<DbRefreshStateRow>('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<DbRefreshStateRow>('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<DbRefreshStateRow>('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,
);
},
);
});
});
@@ -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,