diff --git a/packages/catalog-model/src/entity/Entity.ts b/packages/catalog-model/src/entity/Entity.ts index 9e7ece4122..ac43137546 100644 --- a/packages/catalog-model/src/entity/Entity.ts +++ b/packages/catalog-model/src/entity/Entity.ts @@ -143,13 +143,14 @@ export type EntityRelation = { }; /** - * Holds the relationship data for entities + * Holds the relation data for entities. */ export type EntityRelationSpec = { /** * The source entity of this relation. */ source: EntityName; + /** * The type of the relation. */ diff --git a/packages/catalog-model/src/kinds/index.ts b/packages/catalog-model/src/kinds/index.ts index 6df4f1212a..4a3faf977b 100644 --- a/packages/catalog-model/src/kinds/index.ts +++ b/packages/catalog-model/src/kinds/index.ts @@ -44,3 +44,4 @@ export type { UserEntityV1alpha1 as UserEntity, UserEntityV1alpha1, } from './UserEntityV1alpha1'; +export * from './relations'; diff --git a/packages/catalog-model/src/kinds/relations.ts b/packages/catalog-model/src/kinds/relations.ts new file mode 100644 index 0000000000..846313ec8c --- /dev/null +++ b/packages/catalog-model/src/kinds/relations.ts @@ -0,0 +1,55 @@ +/* + * Copyright 2020 Spotify AB + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +/* +Naming rules for relations in priority order: + +1. Use at most two words. One main verb and a specifier, e.g. "ownerOf" +2. Reading out " " should make sense in English. +3. Maintain symmetry between pairs, e.g. "ownedBy" and "ownerOf" rather than "owns". +*/ + +/** + * An ownership relation where the owner is usually an organizational + * entity (user or group), and the other entity can be anything. + */ +export const RELATION_OWNED_BY = 'ownedBy'; +export const RELATION_OWNER_OF = 'ownerOf'; + +/** + * A relation with an API entity, typically from a component or system + */ +export const RELATION_CONSUMES_API = 'consumesApi'; +export const RELATION_PROVIDES_API = 'providesApi'; + +/** + * A relation denoting a dependency on another entity. + */ +export const RELATION_DEPENDS_ON = 'dependsOn'; +export const RELATION_DEPENDENCY_OF = 'dependencyOf'; + +/** + * A parent/child relation to build up a tree, used for example to describe + * the organizational structure between groups. + */ +export const RELATION_PARENT_OF = 'parentOf'; +export const RELATION_CHILD_OF = 'childOf'; + +/** + * A membership relation, typically for users in a group. + */ +export const RELATION_MEMBER_OF = 'memberOf'; +export const RELATION_HAS_MEMBER = 'hasMember'; diff --git a/plugins/catalog-backend/migrations/20201019130742_add_relations_table.js b/plugins/catalog-backend/migrations/20201019130742_add_relations_table.js index accf9035e7..70150ff173 100644 --- a/plugins/catalog-backend/migrations/20201019130742_add_relations_table.js +++ b/plugins/catalog-backend/migrations/20201019130742_add_relations_table.js @@ -28,7 +28,7 @@ exports.up = async function up(knex) { .inTable('entities') .onDelete('CASCADE') .notNullable() - .comment('The originating entity of the relation'); + .comment('The entity that provided the relation'); table .string('source_full_name') .notNullable() diff --git a/plugins/catalog-backend/src/catalog/DatabaseEntitiesCatalog.test.ts b/plugins/catalog-backend/src/catalog/DatabaseEntitiesCatalog.test.ts index d84a8e071c..a801015915 100644 --- a/plugins/catalog-backend/src/catalog/DatabaseEntitiesCatalog.test.ts +++ b/plugins/catalog-backend/src/catalog/DatabaseEntitiesCatalog.test.ts @@ -18,6 +18,7 @@ import { getVoidLogger } from '@backstage/backend-common'; import type { Entity } from '@backstage/catalog-model'; import { Database, DatabaseManager } from '../database'; import { DatabaseEntitiesCatalog } from './DatabaseEntitiesCatalog'; +import { EntityUpsertRequest } from './types'; describe('DatabaseEntitiesCatalog', () => { let db: jest.Mocked; @@ -46,7 +47,7 @@ describe('DatabaseEntitiesCatalog', () => { db.transaction.mockImplementation(async f => f('tx')); }); - describe('addOrUpdateEntity', () => { + describe('batchAddOrUpdateEntities', () => { it('adds when no given uid and no matching by name', async () => { const entity: Entity = { apiVersion: 'a', @@ -58,19 +59,25 @@ describe('DatabaseEntitiesCatalog', () => { }; db.entities.mockResolvedValue([]); - db.addEntities.mockResolvedValue([{ entity }]); + db.addEntities.mockResolvedValue([ + { entity: { ...entity, metadata: { ...entity.metadata, uid: 'u' } } }, + ]); const catalog = new DatabaseEntitiesCatalog(db, getVoidLogger()); - const result = await catalog.addOrUpdateEntity(entity); + const result = await catalog.batchAddOrUpdateEntities([ + { entity, relations: [] }, + ]); - expect(db.entityByName).toHaveBeenCalledTimes(1); - expect(db.entityByName).toHaveBeenCalledWith(expect.anything(), { + expect(db.entities).toHaveBeenCalledTimes(1); + expect(db.entities).toHaveBeenCalledWith(expect.anything(), { kind: 'b', - namespace: 'd', - name: 'c', + 'metadata.namespace': 'd', + 'metadata.name': ['c'], }); + expect(db.setRelations).toHaveBeenCalledTimes(1); + expect(db.setRelations).toHaveBeenCalledWith(expect.anything(), 'u', []); expect(db.addEntities).toHaveBeenCalledTimes(1); - expect(result).toBe(entity); + expect(result).toEqual([{ entityId: 'u' }]); }); it('updates when given uid', async () => { @@ -82,9 +89,11 @@ describe('DatabaseEntitiesCatalog', () => { name: 'c', namespace: 'd', }, + spec: { + x: 'b', + }, }; - - db.entityByUid.mockResolvedValue({ + const existing = { entity: { apiVersion: 'a', kind: 'b', @@ -95,14 +104,28 @@ describe('DatabaseEntitiesCatalog', () => { name: 'c', namespace: 'd', }, + spec: { + x: 'a', + }, }, - }); + }; + + db.entities.mockResolvedValue([existing]); + db.entityByUid.mockResolvedValue(existing); db.updateEntity.mockResolvedValue({ entity }); const catalog = new DatabaseEntitiesCatalog(db, getVoidLogger()); - const result = await catalog.addOrUpdateEntity(entity); + const result = await catalog.batchAddOrUpdateEntities([ + { entity, relations: [] }, + ]); - expect(db.entities).toHaveBeenCalledTimes(0); + expect(db.entities).toHaveBeenCalledTimes(1); + expect(db.entities).toHaveBeenCalledWith(expect.anything(), { + kind: 'b', + 'metadata.namespace': 'd', + 'metadata.name': ['c'], + }); + expect(db.entityByName).not.toHaveBeenCalled(); expect(db.entityByUid).toHaveBeenCalledTimes(1); expect(db.entityByUid).toHaveBeenCalledWith(expect.anything(), 'u'); expect(db.updateEntity).toHaveBeenCalledTimes(1); @@ -114,17 +137,22 @@ describe('DatabaseEntitiesCatalog', () => { kind: 'b', metadata: { uid: 'u', - etag: 'e', - generation: 1, + etag: expect.any(String), + generation: 2, name: 'c', namespace: 'd', }, + spec: { + x: 'b', + }, }, }, 'e', 1, ); - expect(result).toBe(entity); + expect(db.setRelations).toHaveBeenCalledTimes(1); + expect(db.setRelations).toHaveBeenCalledWith(expect.anything(), 'u', []); + expect(result).toEqual([{ entityId: 'u' }]); }); it('update when no given uid and matching by name', async () => { @@ -135,25 +163,42 @@ describe('DatabaseEntitiesCatalog', () => { name: 'c', namespace: 'd', }, + spec: { + x: 'b', + }, }; - const existing: Entity = { - apiVersion: 'a', - kind: 'b', - metadata: { - uid: 'u', - etag: 'e', - generation: 1, - name: 'c', - namespace: 'd', + const existing = { + entity: { + apiVersion: 'a', + kind: 'b', + metadata: { + uid: 'u', + etag: 'e', + generation: 1, + name: 'c', + namespace: 'd', + }, + spec: { + x: 'a', + }, }, }; - db.entityByName.mockResolvedValue({ entity: existing }); - db.updateEntity.mockResolvedValue({ entity: existing }); + db.entities.mockResolvedValue([existing]); + db.entityByName.mockResolvedValue(existing); + db.updateEntity.mockResolvedValue(existing); const catalog = new DatabaseEntitiesCatalog(db, getVoidLogger()); - const result = await catalog.addOrUpdateEntity(added); + const result = await catalog.batchAddOrUpdateEntities([ + { entity: added, relations: [] }, + ]); + expect(db.entities).toHaveBeenCalledTimes(1); + expect(db.entities).toHaveBeenCalledWith(expect.anything(), { + kind: 'b', + 'metadata.namespace': 'd', + 'metadata.name': ['c'], + }); expect(db.entityByName).toHaveBeenCalledTimes(1); expect(db.entityByName).toHaveBeenCalledWith(expect.anything(), { kind: 'b', @@ -169,51 +214,95 @@ describe('DatabaseEntitiesCatalog', () => { kind: 'b', metadata: { uid: 'u', - etag: 'e', - generation: 1, + etag: expect.any(String), + generation: 2, name: 'c', namespace: 'd', }, + spec: { + x: 'b', + }, }, }, 'e', 1, ); - expect(result).toEqual(existing); + expect(result).toEqual([{ entityId: 'u' }]); + }); + + it('should not update if entity is unchanged', async () => { + const entity: Entity = { + apiVersion: 'a', + kind: 'b', + metadata: { + uid: 'u', + name: 'c', + namespace: 'd', + }, + spec: { + x: 'a', + }, + }; + + db.entities.mockResolvedValue([{ entity }]); + db.entityByUid.mockResolvedValue({ entity }); + db.updateEntity.mockResolvedValue({ entity }); + + const catalog = new DatabaseEntitiesCatalog(db, getVoidLogger()); + const result = await catalog.batchAddOrUpdateEntities([ + { entity, relations: [] }, + ]); + + expect(db.entities).toHaveBeenCalledTimes(1); + expect(db.entities).toHaveBeenCalledWith(expect.anything(), { + kind: 'b', + 'metadata.namespace': 'd', + 'metadata.name': ['c'], + }); + expect(db.entityByName).not.toHaveBeenCalled(); + expect(db.entityByUid).not.toHaveBeenCalled(); + expect(db.updateEntity).not.toHaveBeenCalled(); + expect(db.setRelations).toHaveBeenCalledTimes(1); + expect(db.setRelations).toHaveBeenCalledWith(expect.anything(), 'u', []); + expect(result).toEqual([{ entityId: 'u' }]); }); - }); - describe('batchAddOrUpdateEntities', () => { it('both adds and updates', async () => { const catalog = new DatabaseEntitiesCatalog( await DatabaseManager.createTestDatabase(), getVoidLogger(), ); - const entities: Entity[] = []; - for (let i = 0; i < 500; ++i) { + const entities: EntityUpsertRequest[] = []; + for (let i = 0; i < 300; ++i) { entities.push({ - apiVersion: 'a', - kind: 'k', - metadata: { name: `n${i}` }, + entity: { + apiVersion: 'a', + kind: 'k', + metadata: { name: `n${i}` }, + }, + relations: [], }); } await catalog.batchAddOrUpdateEntities(entities); const afterFirst = await catalog.entities(); - expect(afterFirst.length).toBe(500); + expect(afterFirst.length).toBe(300); - entities[40].metadata.op = 'changed'; + entities[40].entity.metadata.op = 'changed'; entities.push({ - apiVersion: 'a', - kind: 'k', - metadata: { name: `n500`, op: 'added' }, + entity: { + apiVersion: 'a', + kind: 'k', + metadata: { name: `n300`, op: 'added' }, + }, + relations: [], }); await catalog.batchAddOrUpdateEntities(entities); const afterSecond = await catalog.entities(); - expect(afterSecond.length).toBe(501); + expect(afterSecond.length).toBe(301); expect(afterSecond.find(e => e.metadata.op === 'changed')).toBeDefined(); expect(afterSecond.find(e => e.metadata.op === 'added')).toBeDefined(); - }); + }, 10000); }); }); diff --git a/plugins/catalog-backend/src/catalog/DatabaseEntitiesCatalog.ts b/plugins/catalog-backend/src/catalog/DatabaseEntitiesCatalog.ts index efb541727a..6c3065aeaa 100644 --- a/plugins/catalog-backend/src/catalog/DatabaseEntitiesCatalog.ts +++ b/plugins/catalog-backend/src/catalog/DatabaseEntitiesCatalog.ts @@ -29,7 +29,11 @@ import limiterFactory from 'p-limit'; import { Logger } from 'winston'; import type { Database, DbEntityResponse, EntityFilters } from '../database'; import { durationText } from '../util/timing'; -import type { EntitiesCatalog } from './types'; +import type { + EntitiesCatalog, + EntityUpsertRequest, + EntityUpsertResponse, +} from './types'; type BatchContext = { kind: string; @@ -63,7 +67,7 @@ export class DatabaseEntitiesCatalog implements EntitiesCatalog { return items.map(i => i.entity); } - async addOrUpdateEntity( + private async addOrUpdateEntity( entity: Entity, locationId?: string, ): Promise { @@ -96,15 +100,6 @@ export class DatabaseEntitiesCatalog implements EntitiesCatalog { }); } - async addEntities(entities: Entity[], locationId?: string): Promise { - await this.database.transaction(async tx => { - await this.database.addEntities( - tx, - entities.map(entity => ({ locationId, entity })), - ); - }); - } - async removeEntityByUid(uid: string): Promise { return await this.database.transaction(async tx => { const entityResponse = await this.database.entityByUid(tx, uid); @@ -138,27 +133,31 @@ export class DatabaseEntitiesCatalog implements EntitiesCatalog { * @param entities Some entities * @param locationId The location that they all belong to */ - async batchAddOrUpdateEntities(entities: Entity[], locationId?: string) { + async batchAddOrUpdateEntities( + requests: EntityUpsertRequest[], + locationId?: string, + ): Promise { // Group the entities by unique kind+namespace combinations - const entitiesByKindAndNamespace = groupBy(entities, entity => { + const entitiesByKindAndNamespace = groupBy(requests, ({ entity }) => { const name = getEntityName(entity); return `${name.kind}:${name.namespace}`.toLowerCase(); }); const limiter = limiterFactory(BATCH_CONCURRENCY); - const tasks: Promise[] = []; + const tasks: Promise[] = []; - for (const groupEntities of Object.values(entitiesByKindAndNamespace)) { - const { kind, namespace } = getEntityName(groupEntities[0]); + for (const groupRequests of Object.values(entitiesByKindAndNamespace)) { + const { kind, namespace } = getEntityName(groupRequests[0].entity); // Go through the new entities in reasonable chunk sizes (sometimes, // sources produce tens of thousands of entities, and those are too large // batch sizes to reasonably send to the database) - for (const batch of chunk(groupEntities, BATCH_SIZE)) { + for (const batch of chunk(groupRequests, BATCH_SIZE)) { tasks.push( limiter(async () => { - const first = serializeEntityRef(batch[0]); - const last = serializeEntityRef(batch[batch.length - 1]); + const first = serializeEntityRef(batch[0].entity); + const last = serializeEntityRef(batch[batch.length - 1].entity); + const modifiedEntityIds: EntityUpsertResponse[] = []; this.logger.debug( `Considering batch ${first}-${last} (${batch.length} entries)`, ); @@ -167,12 +166,28 @@ export class DatabaseEntitiesCatalog implements EntitiesCatalog { const context = { kind, namespace, locationId }; for (let attempt = 1; attempt <= BATCH_ATTEMPTS; ++attempt) { try { - const { toAdd, toUpdate } = await this.analyzeBatch( + const { toAdd, toUpdate, toIgnore } = await this.analyzeBatch( batch, context, ); - if (toAdd.length) await this.batchAdd(toAdd, context); - if (toUpdate.length) await this.batchUpdate(toUpdate, context); + if (toAdd.length) { + modifiedEntityIds.push( + ...(await this.batchAdd(toAdd, context)), + ); + } + if (toUpdate.length) { + modifiedEntityIds.push( + ...(await this.batchUpdate(toUpdate, context)), + ); + } + // TODO(Rugvip): We currently always update relations, but we + // likely want to figure out a way to avoid that + for (const { entity, relations } of toIgnore) { + const entityId = entity.metadata.uid!; + await this.setRelations(entityId, relations); + modifiedEntityIds.push({ entityId }); + } + break; } catch (e) { if (e instanceof ConflictError && attempt < BATCH_ATTEMPTS) { @@ -184,16 +199,18 @@ export class DatabaseEntitiesCatalog implements EntitiesCatalog { } } } + + return modifiedEntityIds; }), ); } } - await Promise.all(tasks); + return (await Promise.all(tasks)).flat(); } // Set the relations originating from an entity using the DB layer - async setRelations( + private async setRelations( originatingEntityId: string, relations: EntityRelationSpec[], ): Promise { @@ -207,15 +224,16 @@ export class DatabaseEntitiesCatalog implements EntitiesCatalog { // produce the list of entities to be added, and the list of entities to be // updated private async analyzeBatch( - newEntities: Entity[], + requests: EntityUpsertRequest[], { kind, namespace }: BatchContext, ): Promise<{ - toAdd: Entity[]; - toUpdate: Entity[]; + toAdd: EntityUpsertRequest[]; + toUpdate: EntityUpsertRequest[]; + toIgnore: EntityUpsertRequest[]; }> { const markTimestamp = process.hrtime(); - const names = newEntities.map(e => e.metadata.name); + const names = requests.map(({ entity }) => entity.metadata.name); const oldEntities = await this.entities({ kind: kind, 'metadata.namespace': namespace, @@ -226,18 +244,22 @@ export class DatabaseEntitiesCatalog implements EntitiesCatalog { oldEntities.map(e => [e.metadata.name, e]), ); - const toAdd: Entity[] = []; - const toUpdate: Entity[] = []; + const toAdd: EntityUpsertRequest[] = []; + const toUpdate: EntityUpsertRequest[] = []; + const toIgnore: EntityUpsertRequest[] = []; - for (const newEntity of newEntities) { + for (const request of requests) { + const newEntity = request.entity; const oldEntity = oldEntitiesByName.get(newEntity.metadata.name); if (!oldEntity) { - toAdd.push(newEntity); + toAdd.push(request); } else if (entityHasChanges(oldEntity, newEntity)) { // TODO(freben): This currently uses addOrUpdateEntity under the hood, // but should probably calculate the end result entity right here // instead and call a dedicated batch update database method instead - toUpdate.push(newEntity); + toUpdate.push(request); + } else { + toIgnore.push(request); } } @@ -247,33 +269,60 @@ export class DatabaseEntitiesCatalog implements EntitiesCatalog { } entities to update in ${durationText(markTimestamp)}`, ); - return { toAdd, toUpdate }; + return { toAdd, toUpdate, toIgnore }; } // Efficiently adds the given entities to storage, under the assumption that // they do not conflict with any existing entities - private async batchAdd(entities: Entity[], { locationId }: BatchContext) { + private async batchAdd( + requests: EntityUpsertRequest[], + { locationId }: BatchContext, + ): Promise { const markTimestamp = process.hrtime(); - await this.addEntities(entities, locationId); + const res = await this.database.transaction( + async tx => + await this.database.addEntities( + tx, + requests.map(({ entity }) => ({ locationId, entity })), + ), + ); + + const entityIds = res.map(({ entity }) => ({ + entityId: entity.metadata.uid!, + })); + + for (const [index, { entityId }] of entityIds.entries()) { + await this.setRelations(entityId, requests[index].relations); + } this.logger.debug( - `Added ${entities.length} entities in ${durationText(markTimestamp)}`, + `Added ${requests.length} entities in ${durationText(markTimestamp)}`, ); + + return entityIds; } // Efficiently updates the given entities into storage, under the assumption // that there already exist entities with the same names - private async batchUpdate(entities: Entity[], { locationId }: BatchContext) { + private async batchUpdate( + requests: EntityUpsertRequest[], + { locationId }: BatchContext, + ): Promise { const markTimestamp = process.hrtime(); - + const responseIds: EntityUpsertResponse[] = []; // TODO(freben): Still not batched - for (const entity of entities) { - await this.addOrUpdateEntity(entity, locationId); + for (const entity of requests) { + const res = await this.addOrUpdateEntity(entity.entity, locationId); + const entityId = res.metadata.uid!; + responseIds.push({ entityId }); + await this.setRelations(entityId, entity.relations); } this.logger.debug( - `Updated ${entities.length} entities in ${durationText(markTimestamp)}`, + `Updated ${requests.length} entities in ${durationText(markTimestamp)}`, ); + + return responseIds; } } diff --git a/plugins/catalog-backend/src/catalog/types.ts b/plugins/catalog-backend/src/catalog/types.ts index e50b5f15db..4d821c428a 100644 --- a/plugins/catalog-backend/src/catalog/types.ts +++ b/plugins/catalog-backend/src/catalog/types.ts @@ -21,10 +21,17 @@ import type { EntityFilters } from '../database'; // Entities // +export type EntityUpsertRequest = { + entity: Entity; + relations: EntityRelationSpec[]; +}; + +export type EntityUpsertResponse = { + entityId: string; +}; + export type EntitiesCatalog = { entities(filters?: EntityFilters): Promise; - addOrUpdateEntity(entity: Entity, locationId?: string): Promise; - addEntities(entities: Entity[], locationId?: string): Promise; removeEntityByUid(uid: string): Promise; /** @@ -34,15 +41,9 @@ export type EntitiesCatalog = { * @param locationId The location that they all belong to */ batchAddOrUpdateEntities( - entities: Entity[], + entities: EntityUpsertRequest[], locationId?: string, - ): Promise; - - // Same as the DB layer - setRelations( - entityUid: string, - relations: EntityRelationSpec[], - ): Promise; + ): Promise; }; // diff --git a/plugins/catalog-backend/src/ingestion/HigherOrderOperations.test.ts b/plugins/catalog-backend/src/ingestion/HigherOrderOperations.test.ts index 34320fe18f..933433b47e 100644 --- a/plugins/catalog-backend/src/ingestion/HigherOrderOperations.test.ts +++ b/plugins/catalog-backend/src/ingestion/HigherOrderOperations.test.ts @@ -31,10 +31,7 @@ describe('HigherOrderOperations', () => { beforeAll(() => { entitiesCatalog = { entities: jest.fn(), - addOrUpdateEntity: jest.fn(), - addEntities: jest.fn(), removeEntityByUid: jest.fn(), - setRelations: jest.fn(), batchAddOrUpdateEntities: jest.fn(), }; locationsCatalog = { @@ -69,7 +66,10 @@ describe('HigherOrderOperations', () => { }; locationsCatalog.addLocation.mockImplementation(x => Promise.resolve(x)); locationsCatalog.locations.mockResolvedValue([]); - locationReader.read.mockResolvedValue({ entities: [], errors: [] }); + locationReader.read.mockResolvedValue({ + entities: [], + errors: [], + }); const result = await higherOrderOperation.addLocation(spec); @@ -83,7 +83,7 @@ describe('HigherOrderOperations', () => { expect(locationsCatalog.locations).toBeCalledTimes(1); expect(locationReader.read).toBeCalledTimes(1); expect(locationReader.read).toBeCalledWith({ type: 'a', target: 'b' }); - expect(entitiesCatalog.addOrUpdateEntity).not.toBeCalled(); + expect(entitiesCatalog.batchAddOrUpdateEntities).not.toBeCalled(); expect(locationsCatalog.addLocation).toBeCalledTimes(1); expect(locationsCatalog.addLocation).toBeCalledWith( expect.objectContaining({ @@ -109,7 +109,10 @@ describe('HigherOrderOperations', () => { data: location, }, ]); - locationReader.read.mockResolvedValue({ entities: [], errors: [] }); + locationReader.read.mockResolvedValue({ + entities: [], + errors: [], + }); const result = await higherOrderOperation.addLocation(spec); @@ -118,7 +121,7 @@ describe('HigherOrderOperations', () => { expect(locationsCatalog.locations).toBeCalledTimes(1); expect(locationReader.read).toBeCalledTimes(1); expect(locationReader.read).toBeCalledWith({ type: 'a', target: 'b' }); - expect(entitiesCatalog.addOrUpdateEntity).not.toBeCalled(); + expect(entitiesCatalog.batchAddOrUpdateEntities).not.toBeCalled(); expect(locationsCatalog.addLocation).not.toBeCalled(); }); @@ -136,7 +139,7 @@ describe('HigherOrderOperations', () => { locationsCatalog.locations.mockResolvedValue([]); locationReader.read.mockResolvedValue({ - entities: [{ entity, location }], + entities: [{ entity, location, relations: [] }], errors: [{ error: new Error('abcd'), location }], }); @@ -144,7 +147,7 @@ describe('HigherOrderOperations', () => { /abcd/, ); expect(locationsCatalog.locations).toBeCalledTimes(1); - expect(entitiesCatalog.addOrUpdateEntity).not.toBeCalled(); + expect(entitiesCatalog.batchAddOrUpdateEntities).not.toBeCalled(); expect(locationsCatalog.addLocation).not.toBeCalled(); }); }); @@ -159,7 +162,7 @@ describe('HigherOrderOperations', () => { expect(locationsCatalog.locations).toHaveBeenCalledTimes(1); expect(locationReader.read).not.toHaveBeenCalled(); - expect(entitiesCatalog.addOrUpdateEntity).not.toHaveBeenCalled(); + expect(entitiesCatalog.batchAddOrUpdateEntities).not.toHaveBeenCalled(); }); it('can update a single location where a matching entity did not exist', async () => { @@ -179,15 +182,18 @@ describe('HigherOrderOperations', () => { metadata: { name: 'c1' }, spec: { type: 'service' }, }; + const entityId = 'xyz123'; locationsCatalog.locations.mockResolvedValue([ { currentStatus: locationStatus, data: location }, ]); locationReader.read.mockResolvedValue({ - entities: [{ entity: desc, location }], + entities: [{ entity: desc, location, relations: [] }], errors: [], }); - entitiesCatalog.batchAddOrUpdateEntities.mockResolvedValue(undefined); + entitiesCatalog.batchAddOrUpdateEntities.mockResolvedValue([ + { entityId }, + ]); await expect( higherOrderOperation.refreshAllLocations(), @@ -200,9 +206,13 @@ describe('HigherOrderOperations', () => { target: 'thing', }); expect(entitiesCatalog.batchAddOrUpdateEntities).toHaveBeenCalledTimes(1); - expect(entitiesCatalog.batchAddOrUpdateEntities).toHaveBeenNthCalledWith( - 1, - [expect.objectContaining({ metadata: { name: 'c1' } })], + expect(entitiesCatalog.batchAddOrUpdateEntities).toHaveBeenCalledWith( + [ + expect.objectContaining({ + entity: expect.objectContaining({ metadata: { name: 'c1' } }), + relations: [], + }), + ], '123', ); }); @@ -229,11 +239,11 @@ describe('HigherOrderOperations', () => { { currentStatus: locationStatus, data: location }, ]); locationReader.read.mockResolvedValue({ - entities: [{ entity: desc, location }], + entities: [{ entity: desc, location, relations: [] }], errors: [], }); entitiesCatalog.entities.mockResolvedValue([]); - entitiesCatalog.addEntities.mockResolvedValue(undefined); + entitiesCatalog.batchAddOrUpdateEntities.mockResolvedValue([]); await expect( higherOrderOperation.refreshAllLocations(), diff --git a/plugins/catalog-backend/src/ingestion/HigherOrderOperations.ts b/plugins/catalog-backend/src/ingestion/HigherOrderOperations.ts index b1f4e40524..23f35a4945 100644 --- a/plugins/catalog-backend/src/ingestion/HigherOrderOperations.ts +++ b/plugins/catalog-backend/src/ingestion/HigherOrderOperations.ts @@ -15,7 +15,7 @@ */ import { InputError } from '@backstage/backend-common'; -import { Entity, Location, LocationSpec } from '@backstage/catalog-model'; +import { Location, LocationSpec } from '@backstage/catalog-model'; import { v4 as uuidv4 } from 'uuid'; import { Logger } from 'winston'; import { EntitiesCatalog, LocationsCatalog } from '../catalog'; @@ -94,16 +94,20 @@ export class HigherOrderOperations implements HigherOrderOperation { if (!previousLocation) { await this.locationsCatalog.addLocation(location); } - const outputEntities: Entity[] = []; - for (const entity of readerOutput.entities) { - const out = await this.entitiesCatalog.addOrUpdateEntity( - entity.entity, - location.id, - ); - outputEntities.push(out); + if (readerOutput.entities.length === 0) { + return { location, entities: [] }; } - return { location, entities: outputEntities }; + const writtenEntities = await this.entitiesCatalog.batchAddOrUpdateEntities( + readerOutput.entities, + location.id, + ); + + const entities = await this.entitiesCatalog.entities({ + 'metadata.uid': writtenEntities.map(e => e.entityId), + }); + + return { location, entities }; } /** @@ -166,7 +170,7 @@ export class HigherOrderOperations implements HigherOrderOperation { try { await this.entitiesCatalog.batchAddOrUpdateEntities( - readerOutput.entities.map(e => e.entity), + readerOutput.entities, location.id, ); } catch (e) { diff --git a/plugins/catalog-backend/src/ingestion/LocationReaders.ts b/plugins/catalog-backend/src/ingestion/LocationReaders.ts index e23a4524b6..215bfc9dc1 100644 --- a/plugins/catalog-backend/src/ingestion/LocationReaders.ts +++ b/plugins/catalog-backend/src/ingestion/LocationReaders.ts @@ -18,6 +18,7 @@ import { UrlReader } from '@backstage/backend-common'; import { Entity, EntityPolicy, + EntityRelationSpec, ENTITY_DEFAULT_NAMESPACE, LocationSpec, } from '@backstage/catalog-model'; @@ -61,7 +62,10 @@ export class LocationReaders implements LocationReader { async read(location: LocationSpec): Promise { const { rulesEnforcer, logger } = this.options; - const output: ReadLocationResult = { entities: [], errors: [] }; + const output: ReadLocationResult = { + entities: [], + errors: [], + }; let items: CatalogProcessorResult[] = [result.location(location, false)]; for (let depth = 0; depth < MAX_DEPTH; ++depth) { @@ -75,11 +79,21 @@ export class LocationReaders implements LocationReader { await this.handleData(item, emit); } else if (item.type === 'entity') { if (rulesEnforcer.isAllowed(item.entity, item.location)) { - const entity = await this.handleEntity(item, emit); + const relations = Array(); + + const entity = await this.handleEntity(item, emitResult => { + if (emitResult.type === 'relation') { + relations.push(emitResult.relation); + return; + } + emit(emitResult); + }); + if (entity) { output.entities.push({ entity, location: item.location, + relations, }); } } else { @@ -118,11 +132,23 @@ export class LocationReaders implements LocationReader { ) { const { processors, logger } = this.options; + const validatedEmit: CatalogProcessorEmit = emitResult => { + if (emitResult.type === 'relation') { + throw new Error('readLocation may not emit entity relations'); + } + + emit(emitResult); + }; + for (const processor of processors) { if (processor.readLocation) { try { if ( - await processor.readLocation(item.location, item.optional, emit) + await processor.readLocation( + item.location, + item.optional, + validatedEmit, + ) ) { return; } @@ -145,10 +171,20 @@ export class LocationReaders implements LocationReader { ) { const { processors, logger } = this.options; + const validatedEmit: CatalogProcessorEmit = emitResult => { + if (emitResult.type === 'relation') { + throw new Error('parseData may not emit entity relations'); + } + + emit(emitResult); + }; + for (const processor of processors) { if (processor.parseData) { try { - if (await processor.parseData(item.data, item.location, emit)) { + if ( + await processor.parseData(item.data, item.location, validatedEmit) + ) { return; } } catch (e) { @@ -242,10 +278,18 @@ export class LocationReaders implements LocationReader { `Encountered error at location ${item.location.type} ${item.location.target}, ${item.error}`, ); + const validatedEmit: CatalogProcessorEmit = emitResult => { + if (emitResult.type === 'relation') { + throw new Error('handleError may not emit entity relations'); + } + + emit(emitResult); + }; + for (const processor of processors) { if (processor.handleError) { try { - await processor.handleError(item.error, item.location, emit); + await processor.handleError(item.error, item.location, validatedEmit); } catch (e) { const message = `Processor ${processor.constructor.name} threw an error while handling another error at ${item.location.type} ${item.location.target}, ${e}`; emit(result.generalError(item.location, message)); diff --git a/plugins/catalog-backend/src/ingestion/processors/OwnerRelationProcessor.ts b/plugins/catalog-backend/src/ingestion/processors/OwnerRelationProcessor.ts new file mode 100644 index 0000000000..4c67589789 --- /dev/null +++ b/plugins/catalog-backend/src/ingestion/processors/OwnerRelationProcessor.ts @@ -0,0 +1,74 @@ +/* + * Copyright 2020 Spotify AB + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +import { + Entity, + ENTITY_DEFAULT_NAMESPACE, + LocationSpec, + parseEntityRef, + ApiEntityV1alpha1, + ComponentEntityV1alpha1, + RELATION_OWNED_BY, + RELATION_OWNER_OF, + getEntityName, +} from '@backstage/catalog-model'; +import { CatalogProcessor, CatalogProcessorEmit } from './types'; +import * as result from './results'; + +const includedKinds = new Set(['api', 'component']); + +export class OwnerRelationProcessor implements CatalogProcessor { + async postProcessEntity( + entity: Entity, + _location: LocationSpec, + emit: CatalogProcessorEmit, + ): Promise { + if (!includedKinds.has(entity.kind.toLowerCase())) { + return entity; + } + const apiOrComponentEntity = entity as + | ApiEntityV1alpha1 + | ComponentEntityV1alpha1; + + const owner = apiOrComponentEntity.spec?.owner; + if (owner) { + const namespace = entity.metadata.namespace ?? ENTITY_DEFAULT_NAMESPACE; + + const selfRef = getEntityName(entity); + const ownerRef = parseEntityRef(owner, { + defaultKind: 'group', + defaultNamespace: namespace, + }); + + emit( + result.relation({ + source: selfRef, + type: RELATION_OWNED_BY, + target: ownerRef, + }), + ); + emit( + result.relation({ + source: ownerRef, + type: RELATION_OWNER_OF, + target: selfRef, + }), + ); + } + + return entity; + } +} diff --git a/plugins/catalog-backend/src/ingestion/processors/index.ts b/plugins/catalog-backend/src/ingestion/processors/index.ts index 59a87af555..92303f04a3 100644 --- a/plugins/catalog-backend/src/ingestion/processors/index.ts +++ b/plugins/catalog-backend/src/ingestion/processors/index.ts @@ -28,6 +28,7 @@ export { GithubOrgReaderProcessor } from './GithubOrgReaderProcessor'; export { GithubReaderProcessor } from './GithubReaderProcessor'; export { GitlabApiReaderProcessor } from './GitlabApiReaderProcessor'; export { GitlabReaderProcessor } from './GitlabReaderProcessor'; +export { OwnerRelationProcessor } from './OwnerRelationProcessor'; export { LocationRefProcessor } from './LocationEntityProcessor'; export { PlaceholderProcessor } from './PlaceholderProcessor'; export type { PlaceholderResolver } from './PlaceholderProcessor'; diff --git a/plugins/catalog-backend/src/ingestion/processors/results.ts b/plugins/catalog-backend/src/ingestion/processors/results.ts index 2718d871dd..158bc3d950 100644 --- a/plugins/catalog-backend/src/ingestion/processors/results.ts +++ b/plugins/catalog-backend/src/ingestion/processors/results.ts @@ -15,7 +15,11 @@ */ import { InputError, NotFoundError } from '@backstage/backend-common'; -import { Entity, LocationSpec } from '@backstage/catalog-model'; +import { + Entity, + EntityRelationSpec, + LocationSpec, +} from '@backstage/catalog-model'; import { CatalogProcessorResult } from './types'; export function notFoundError( @@ -67,3 +71,7 @@ export function entity( ): CatalogProcessorResult { return { type: 'entity', location: atLocation, entity: newEntity }; } + +export function relation(spec: EntityRelationSpec): CatalogProcessorResult { + return { type: 'relation', relation: spec }; +} diff --git a/plugins/catalog-backend/src/ingestion/processors/types.ts b/plugins/catalog-backend/src/ingestion/processors/types.ts index c9adf07eff..7821a61437 100644 --- a/plugins/catalog-backend/src/ingestion/processors/types.ts +++ b/plugins/catalog-backend/src/ingestion/processors/types.ts @@ -14,7 +14,11 @@ * limitations under the License. */ -import { Entity, LocationSpec } from '@backstage/catalog-model'; +import { + Entity, + EntityRelationSpec, + LocationSpec, +} from '@backstage/catalog-model'; export type CatalogProcessor = { /** @@ -113,6 +117,12 @@ export type CatalogProcessorEntityResult = { location: LocationSpec; }; +export type CatalogProcessorRelationResult = { + type: 'relation'; + relation: EntityRelationSpec; + entityRef?: string; +}; + export type CatalogProcessorErrorResult = { type: 'error'; error: Error; @@ -123,4 +133,5 @@ export type CatalogProcessorResult = | CatalogProcessorLocationResult | CatalogProcessorDataResult | CatalogProcessorEntityResult + | CatalogProcessorRelationResult | CatalogProcessorErrorResult; diff --git a/plugins/catalog-backend/src/ingestion/types.ts b/plugins/catalog-backend/src/ingestion/types.ts index 5bbec01e7e..c586e2d438 100644 --- a/plugins/catalog-backend/src/ingestion/types.ts +++ b/plugins/catalog-backend/src/ingestion/types.ts @@ -14,7 +14,12 @@ * limitations under the License. */ -import type { Entity, Location, LocationSpec } from '@backstage/catalog-model'; +import type { + Entity, + EntityRelationSpec, + Location, + LocationSpec, +} from '@backstage/catalog-model'; // // HigherOrderOperation @@ -53,6 +58,7 @@ export type ReadLocationResult = { export type ReadLocationEntity = { location: LocationSpec; entity: Entity; + relations: EntityRelationSpec[]; }; export type ReadLocationError = { diff --git a/plugins/catalog-backend/src/service/CatalogBuilder.test.ts b/plugins/catalog-backend/src/service/CatalogBuilder.test.ts index 3f7c613c00..bc71934488 100644 --- a/plugins/catalog-backend/src/service/CatalogBuilder.test.ts +++ b/plugins/catalog-backend/src/service/CatalogBuilder.test.ts @@ -128,12 +128,16 @@ describe('CatalogBuilder', () => { { apiVersion: 'av', kind: 'Component', - metadata: expect.objectContaining({ + metadata: { name: 'n', namespace: 'ns', post: 'p', replaced: 'tt2', - }), + uid: expect.any(String), + etag: expect.any(String), + generation: expect.any(Number), + }, + relations: [], }, ]); }); diff --git a/plugins/catalog-backend/src/service/CatalogBuilder.ts b/plugins/catalog-backend/src/service/CatalogBuilder.ts index e6fb540f16..4ec2a0834d 100644 --- a/plugins/catalog-backend/src/service/CatalogBuilder.ts +++ b/plugins/catalog-backend/src/service/CatalogBuilder.ts @@ -52,6 +52,7 @@ import { GithubReaderProcessor, GitlabApiReaderProcessor, GitlabReaderProcessor, + OwnerRelationProcessor, HigherOrderOperation, HigherOrderOperations, LocationReaders, @@ -334,6 +335,7 @@ export class CatalogBuilder { new YamlProcessor(), new CodeOwnersProcessor({ reader }), new LocationRefProcessor(), + new OwnerRelationProcessor(), new AnnotateLocationEntityProcessor(), ]; diff --git a/plugins/catalog-backend/src/service/router.test.ts b/plugins/catalog-backend/src/service/router.test.ts index a2b17ba360..cf843acf3c 100644 --- a/plugins/catalog-backend/src/service/router.test.ts +++ b/plugins/catalog-backend/src/service/router.test.ts @@ -32,10 +32,7 @@ describe('createRouter', () => { beforeAll(async () => { entitiesCatalog = { entities: jest.fn(), - addOrUpdateEntity: jest.fn(), - addEntities: jest.fn(), removeEntityByUid: jest.fn(), - setRelations: jest.fn(), batchAddOrUpdateEntities: jest.fn(), }; locationsCatalog = { @@ -173,7 +170,7 @@ describe('createRouter', () => { .set('Content-Type', 'application/json') .send(); - expect(entitiesCatalog.addOrUpdateEntity).not.toHaveBeenCalled(); + expect(entitiesCatalog.batchAddOrUpdateEntities).not.toHaveBeenCalled(); expect(response.status).toEqual(400); expect(response.text).toMatch(/body/); }); @@ -188,18 +185,24 @@ describe('createRouter', () => { }, }; - entitiesCatalog.addOrUpdateEntity.mockResolvedValue(entity); + entitiesCatalog.batchAddOrUpdateEntities.mockResolvedValue([ + { entityId: 'u' }, + ]); + entitiesCatalog.entities.mockResolvedValue([entity]); const response = await request(app) .post('/entities') .send(entity) .set('Content-Type', 'application/json'); - expect(entitiesCatalog.addOrUpdateEntity).toHaveBeenCalledTimes(1); - expect(entitiesCatalog.addOrUpdateEntity).toHaveBeenNthCalledWith( - 1, - entity, - ); + expect(entitiesCatalog.batchAddOrUpdateEntities).toHaveBeenCalledTimes(1); + expect(entitiesCatalog.batchAddOrUpdateEntities).toHaveBeenCalledWith([ + { entity, relations: [] }, + ]); + expect(entitiesCatalog.entities).toHaveBeenCalledTimes(1); + expect(entitiesCatalog.entities).toHaveBeenCalledWith({ + 'metadata.uid': 'u', + }); expect(response.status).toEqual(200); expect(response.body).toEqual(entity); }); diff --git a/plugins/catalog-backend/src/service/router.ts b/plugins/catalog-backend/src/service/router.ts index a262e9b5f4..c7d05afb2d 100644 --- a/plugins/catalog-backend/src/service/router.ts +++ b/plugins/catalog-backend/src/service/router.ts @@ -49,8 +49,13 @@ export async function createRouter( }) .post('/entities', async (req, res) => { const body = await requireRequestBody(req); - const result = await entitiesCatalog.addOrUpdateEntity(body as Entity); - res.status(200).send(result); + const [result] = await entitiesCatalog.batchAddOrUpdateEntities([ + { entity: body as Entity, relations: [] }, + ]); + const [entity] = await entitiesCatalog.entities({ + 'metadata.uid': result.entityId, + }); + res.status(200).send(entity); }) .get('/entities/by-uid/:uid', async (req, res) => { const { uid } = req.params; diff --git a/plugins/catalog/src/components/AboutCard/AboutCard.tsx b/plugins/catalog/src/components/AboutCard/AboutCard.tsx index ef04091c9f..fbbf5fa11c 100644 --- a/plugins/catalog/src/components/AboutCard/AboutCard.tsx +++ b/plugins/catalog/src/components/AboutCard/AboutCard.tsx @@ -14,7 +14,12 @@ * limitations under the License. */ -import { Entity, ENTITY_DEFAULT_NAMESPACE } from '@backstage/catalog-model'; +import { + Entity, + ENTITY_DEFAULT_NAMESPACE, + RELATION_OWNED_BY, + serializeEntityRef, +} from '@backstage/catalog-model'; import { Card, CardContent, @@ -146,7 +151,20 @@ export function AboutCard({ entity, variant }: AboutCardProps) { r.type === RELATION_OWNED_BY) + .map(({ target: { kind, name, namespace } }) => + // TODO(Rugvip): we want to provide some utils for this + serializeEntityRef({ + kind, + name, + namespace: + namespace === ENTITY_DEFAULT_NAMESPACE + ? undefined + : namespace, + }), + ) + .join(', ')} gridSizes={{ xs: 12, sm: 6, lg: 4 }} />