diff --git a/packages/backend/src/plugins/catalog.ts b/packages/backend/src/plugins/catalog.ts index 7651836e2a..63a3e53c81 100644 --- a/packages/backend/src/plugins/catalog.ts +++ b/packages/backend/src/plugins/catalog.ts @@ -20,6 +20,7 @@ import { createRouter, NextCatalogBuilder, runPeriodically, + createNextRouter, } from '@backstage/plugin-catalog-backend'; import { Router } from 'express'; import { PluginEnvironment } from '../types'; @@ -36,18 +37,18 @@ export default async function createPlugin( const builder = new NextCatalogBuilder(env); const { entitiesCatalog, - locationsCatalog, locationAnalyzer, processingEngine, + locationService, } = await builder.build(); // TODO(jhaals): run and manage in background. await processingEngine.start(); - return await createRouter({ + return await createNextRouter({ entitiesCatalog, - locationsCatalog, locationAnalyzer, + locationService, logger: env.logger, config: env.config, }); diff --git a/plugins/catalog-backend/src/next/DefaultCatalogProcessingOrchestrator.ts b/plugins/catalog-backend/src/next/DefaultCatalogProcessingOrchestrator.ts index 8fbb19e536..9b2d6550d1 100644 --- a/plugins/catalog-backend/src/next/DefaultCatalogProcessingOrchestrator.ts +++ b/plugins/catalog-backend/src/next/DefaultCatalogProcessingOrchestrator.ts @@ -97,11 +97,8 @@ export class DefaultCatalogProcessingOrchestrator async process( request: EntityProcessingRequest, ): Promise { - const { entity } = request; - - const result = await this.processSingleEntity(entity); - - return result; + // TODO: implement dryRun/eager + return this.processSingleEntity(request.entity); } private async processSingleEntity( diff --git a/plugins/catalog-backend/src/next/NextCatalogBuilder.ts b/plugins/catalog-backend/src/next/NextCatalogBuilder.ts index 92f3822055..43f321248e 100644 --- a/plugins/catalog-backend/src/next/NextCatalogBuilder.ts +++ b/plugins/catalog-backend/src/next/NextCatalogBuilder.ts @@ -63,11 +63,12 @@ import { } from '../ingestion/processors/PlaceholderProcessor'; import { defaultEntityDataParser } from '../ingestion/processors/util/parse'; import { LocationAnalyzer } from '../ingestion/types'; -import { CatalogProcessingEngine } from '../next/types'; +import { CatalogProcessingEngine, LocationService } from '../next/types'; import { ConfigLocationEntityProvider } from './ConfigLocationEntityProvider'; import { DefaultProcessingDatabase } from './database/DefaultProcessingDatabase'; import { DefaultCatalogProcessingEngine } from './DefaultCatalogProcessingEngine'; import { DefaultCatalogProcessingOrchestrator } from './DefaultCatalogProcessingOrchestrator'; +import { DefaultLocationService } from './DefaultLocationService'; import { DefaultLocationStore } from './DefaultLocationStore'; import { NextEntitiesCatalog } from './NextEntitiesCatalog'; import { Stitcher } from './Stitcher'; @@ -233,6 +234,7 @@ export class NextCatalogBuilder { locationsCatalog: LocationsCatalog; locationAnalyzer: LocationAnalyzer; processingEngine: CatalogProcessingEngine; + locationService: LocationService; }> { const { config, database, logger } = this.env; @@ -274,12 +276,16 @@ export class NextCatalogBuilder { const locationsCatalog = new DatabaseLocationsCatalog(db); const locationAnalyzer = new RepoLocationAnalyzer(logger); - + const locationService = new DefaultLocationService( + locationStore, + orchestrator, + ); return { entitiesCatalog, locationsCatalog, locationAnalyzer, processingEngine, + locationService, }; } diff --git a/plugins/catalog-backend/src/next/NextEntitiesCatalog.ts b/plugins/catalog-backend/src/next/NextEntitiesCatalog.ts index e87bb01e1f..13cbf39c7b 100644 --- a/plugins/catalog-backend/src/next/NextEntitiesCatalog.ts +++ b/plugins/catalog-backend/src/next/NextEntitiesCatalog.ts @@ -18,37 +18,128 @@ import { Knex } from 'knex'; import { DbFinalEntitiesRow } from './Stitcher'; import { EntitiesCatalog } from '../catalog'; import { EntitiesRequest, EntitiesResponse } from '../catalog/types'; +import { DbRefreshStateRow } from './database/DefaultProcessingDatabase'; +import { + DbEntitiesSearchRow, + DbPageInfo, + EntityPagination, +} from '../database/types'; +import { InputError } from '@backstage/errors'; + +function parsePagination( + input?: EntityPagination, +): { limit?: number; offset?: number } { + if (!input) { + return {}; + } + + let { limit, offset } = input; + + if (input.after !== undefined) { + let cursor; + try { + const json = Buffer.from(input.after, 'base64').toString('utf8'); + cursor = JSON.parse(json); + } catch { + throw new InputError('Malformed after cursor, could not be parsed'); + } + if (cursor.limit !== undefined) { + if (!Number.isInteger(cursor.limit)) { + throw new InputError('Malformed after cursor, limit was not an number'); + } + limit = cursor.limit; + } + if (cursor.offset !== undefined) { + if (!Number.isInteger(cursor.offset)) { + throw new InputError('Malformed after cursor, offset was not a number'); + } + offset = cursor.offset; + } + } + + return { limit, offset }; +} + +function stringifyPagination(input: { limit: number; offset: number }) { + const json = JSON.stringify({ limit: input.limit, offset: input.offset }); + const base64 = Buffer.from(json, 'utf8').toString('base64'); + return base64; +} export class NextEntitiesCatalog implements EntitiesCatalog { constructor(private readonly database: Knex) {} async entities(request?: EntitiesRequest): Promise { - if (request?.fields) { - throw new Error('Fields extraction is not implemented'); - } - if (request?.pagination) { - throw new Error('Pagination is not implemented'); - } - if (request?.filter) { - throw new Error('Filters are not implemented'); + const db = this.database; + + let entitiesQuery = db('final_entities'); + + for (const singleFilter of request?.filter?.anyOf ?? []) { + entitiesQuery = entitiesQuery.orWhere(function singleFilterFn() { + for (const { key, matchValueIn } of singleFilter.allOf) { + // NOTE(freben): This used to be a set of OUTER JOIN, which may seem to + // make a lot of sense. However, it had abysmal performance on sqlite + // when datasets grew large, so we're using IN instead. + const matchQuery = db('search') + .select('entity_id') + .where(function keyFilter() { + this.andWhere({ key: key.toLowerCase() }); + if (matchValueIn) { + if (matchValueIn.length === 1) { + this.andWhere({ value: matchValueIn[0].toLowerCase() }); + } else if (matchValueIn.length > 1) { + this.andWhere( + 'value', + 'in', + matchValueIn.map(v => v.toLowerCase()), + ); + } + } + }); + this.andWhere('entity_id', 'in', matchQuery); + } + }); } - const dbResponse = await this.database( - 'final_entities', - ).select(); + // TODO: move final_entities to use entity_Ref + entitiesQuery = entitiesQuery + .select('final_entities.*') + .orderBy('entity_id', 'asc'); - const entities = dbResponse.map(e => JSON.parse(e.final_entity)); + const { limit, offset } = parsePagination(request?.pagination); + if (limit !== undefined) { + entitiesQuery = entitiesQuery.limit(limit + 1); + } + if (offset !== undefined) { + entitiesQuery = entitiesQuery.offset(offset); + } + + let rows = await entitiesQuery; + + let pageInfo: DbPageInfo; + if (limit === undefined || rows.length <= limit) { + pageInfo = { hasNextPage: false }; + } else { + rows = rows.slice(0, -1); + pageInfo = { + hasNextPage: true, + endCursor: stringifyPagination({ + limit, + offset: (offset ?? 0) + limit, + }), + }; + } return { - entities, - pageInfo: { - hasNextPage: false, - }, + entities: rows.map(e => JSON.parse(e.final_entity)), + pageInfo, }; } - async removeEntityByUid(_uid: string): Promise { - throw new Error('Not implemented'); + async removeEntityByUid(uid: string): Promise { + await this.database('refresh_state') + .where('entity_id', uid) + .delete(); } async batchAddOrUpdateEntities(): Promise { diff --git a/plugins/catalog-backend/src/next/NextRouter.ts b/plugins/catalog-backend/src/next/NextRouter.ts new file mode 100644 index 0000000000..433abd97e9 --- /dev/null +++ b/plugins/catalog-backend/src/next/NextRouter.ts @@ -0,0 +1,164 @@ +/* + * 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 { errorHandler } from '@backstage/backend-common'; +import { + analyzeLocationSchema, + locationSpecSchema, +} from '@backstage/catalog-model'; +import { Config } from '@backstage/config'; +import { NotFoundError } from '@backstage/errors'; +import express from 'express'; +import Router from 'express-promise-router'; +import { Logger } from 'winston'; +import yn from 'yn'; +import { EntitiesCatalog } from '../catalog'; +import { LocationAnalyzer } from '../ingestion/types'; +import { + basicEntityFilter, + parseEntityFilterParams, + parseEntityPaginationParams, + parseEntityTransformParams, +} from '../service/request'; +import { LocationService } from './types'; +import { disallowReadonlyMode, validateRequestBody } from '../service/util'; + +export interface RouterOptions { + entitiesCatalog?: EntitiesCatalog; + locationAnalyzer?: LocationAnalyzer; + locationService: LocationService; + logger: Logger; + config: Config; +} + +export async function createNextRouter( + options: RouterOptions, +): Promise { + const { + entitiesCatalog, + locationAnalyzer, + locationService, + config, + logger, + } = options; + + const router = Router(); + router.use(express.json()); + + const readonlyEnabled = + config.getOptionalBoolean('catalog.readonly') || false; + if (readonlyEnabled) { + logger.info('Catalog is running in readonly mode'); + } + + if (entitiesCatalog) { + router + .get('/entities', async (req, res) => { + const { entities, pageInfo } = await entitiesCatalog.entities({ + filter: parseEntityFilterParams(req.query), + fields: parseEntityTransformParams(req.query), + pagination: parseEntityPaginationParams(req.query), + }); + + // Add a Link header to the next page + if (pageInfo.hasNextPage) { + const url = new URL(`http://ignored${req.url}`); + url.searchParams.delete('offset'); + url.searchParams.set('after', pageInfo.endCursor); + res.setHeader('link', `<${url.pathname}${url.search}>; rel="next"`); + } + + // TODO(freben): encode the pageInfo in the response + res.json(entities); + }) + .get('/entities/by-uid/:uid', async (req, res) => { + const { uid } = req.params; + const { entities } = await entitiesCatalog.entities({ + filter: basicEntityFilter({ 'metadata.uid': uid }), + }); + if (!entities.length) { + throw new NotFoundError(`No entity with uid ${uid}`); + } + res.status(200).json(entities[0]); + }) + .delete('/entities/by-uid/:uid', async (req, res) => { + const { uid } = req.params; + await entitiesCatalog.removeEntityByUid(uid); + res.status(204).end(); + }) + .get('/entities/by-name/:kind/:namespace/:name', async (req, res) => { + const { kind, namespace, name } = req.params; + const { entities } = await entitiesCatalog.entities({ + filter: basicEntityFilter({ + kind: kind, + 'metadata.namespace': namespace, + 'metadata.name': name, + }), + }); + if (!entities.length) { + throw new NotFoundError( + `No entity named '${name}' found, with kind '${kind}' in namespace '${namespace}'`, + ); + } + res.status(200).json(entities[0]); + }); + } + + if (locationService) { + router + .post('/locations', async (req, res) => { + const input = await validateRequestBody(req, locationSpecSchema); + const dryRun = yn(req.query.dryRun, { default: false }); + + // when in dryRun addLocation is effectively a read operation so we don't + // need to disallow readonly + if (!dryRun) { + disallowReadonlyMode(readonlyEnabled); + } + + const output = await locationService.createLocation(input, dryRun); + res.status(201).json(output); + }) + .get('/locations', async (_req, res) => { + const output = await locationService.listLocations(); + res.status(200).json(output); + }) + + .get('/locations/:id', async (req, res) => { + const { id } = req.params; + const output = await locationService.getLocation(id); + res.status(200).json(output); + }) + .delete('/locations/:id', async (req, res) => { + disallowReadonlyMode(readonlyEnabled); + + const { id } = req.params; + await locationService.deleteLocation(id); + res.status(204).end(); + }); + } + + if (locationAnalyzer) { + router.post('/analyze-location', async (req, res) => { + const input = await validateRequestBody(req, analyzeLocationSchema); + const output = await locationAnalyzer.analyzeLocation(input); + res.status(200).json(output); + }); + } + + router.use(errorHandler()); + return router; +} diff --git a/plugins/catalog-backend/src/next/Stitcher.ts b/plugins/catalog-backend/src/next/Stitcher.ts index bdc028d923..2bddf41ac9 100644 --- a/plugins/catalog-backend/src/next/Stitcher.ts +++ b/plugins/catalog-backend/src/next/Stitcher.ts @@ -84,14 +84,14 @@ export class Stitcher { relationTarget: 'relations.target_entity_ref', }) .from('refresh_state') - .leftJoin('incoming_references', {}) + .where({ 'refresh_state.entity_ref': entityRef }) + .crossJoin(tx.raw('incoming_references')) .leftOuterJoin('final_entities', { 'final_entities.entity_id': 'refresh_state.entity_id', }) .leftOuterJoin('relations', { 'relations.source_entity_ref': 'refresh_state.entity_ref', }) - .where({ 'refresh_state.entity_ref': entityRef }) .orderBy('relationType', 'asc') .orderBy('relationTarget', 'asc'); diff --git a/plugins/catalog-backend/src/next/index.ts b/plugins/catalog-backend/src/next/index.ts index d4eeca4c81..76867e05c8 100644 --- a/plugins/catalog-backend/src/next/index.ts +++ b/plugins/catalog-backend/src/next/index.ts @@ -15,3 +15,4 @@ */ export { NextCatalogBuilder } from './NextCatalogBuilder'; +export { createNextRouter } from './NextRouter';