From f9f4f87a78f0b4ab6d6d0a4080f9cc00dda00ef8 Mon Sep 17 00:00:00 2001 From: Patrik Oldsberg Date: Fri, 9 Dec 2022 16:07:26 +0100 Subject: [PATCH] catalog-backend: split out pruned delete and refresh by key operations MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: Fredrik Adelöw Signed-off-by: Patrik Oldsberg --- .../src/database/DefaultProcessingDatabase.ts | 14 +- .../src/database/DefaultProviderDatabase.ts | 161 +++--------------- .../deleteWithEagerPruningOfChildren.ts | 155 +++++++++++++++++ .../provider/refreshByRefreshKeys.ts | 43 +++++ .../refreshState/checkLocationKeyConflict.ts | 12 +- .../refreshState/insertUnprocessedEntity.ts | 16 +- .../refreshState/updateUnprocessedEntity.ts | 14 +- 7 files changed, 250 insertions(+), 165 deletions(-) create mode 100644 plugins/catalog-backend/src/database/operations/provider/deleteWithEagerPruningOfChildren.ts create mode 100644 plugins/catalog-backend/src/database/operations/provider/refreshByRefreshKeys.ts diff --git a/plugins/catalog-backend/src/database/DefaultProcessingDatabase.ts b/plugins/catalog-backend/src/database/DefaultProcessingDatabase.ts index 6e77a6a0d8..1127deed7e 100644 --- a/plugins/catalog-backend/src/database/DefaultProcessingDatabase.ts +++ b/plugins/catalog-backend/src/database/DefaultProcessingDatabase.ts @@ -322,24 +322,24 @@ export class DefaultProcessingDatabase implements ProcessingDatabase { const entityRef = stringifyEntityRef(entity); const hash = generateStableHash(entity); - const updated = await updateUnprocessedEntity( + const updated = await updateUnprocessedEntity({ tx, entity, hash, locationKey, - ); + }); if (updated) { stateReferences.push(entityRef); continue; } - const inserted = await insertUnprocessedEntity( + const inserted = await insertUnprocessedEntity({ tx, entity, hash, - this.options.logger, locationKey, - ); + logger: this.options.logger, + }); if (inserted) { stateReferences.push(entityRef); continue; @@ -348,11 +348,11 @@ export class DefaultProcessingDatabase implements ProcessingDatabase { // If the row can't be inserted, we have a conflict, but it could be either // because of a conflicting locationKey or a race with another instance, so check // whether the conflicting entity has the same entityRef but a different locationKey - const conflictingKey = await checkLocationKeyConflict( + const conflictingKey = await checkLocationKeyConflict({ tx, entityRef, locationKey, - ); + }); if (conflictingKey) { this.options.logger.warn( `Detected conflicting entityRef ${entityRef} already referenced by ${conflictingKey} and now also ${locationKey}`, diff --git a/plugins/catalog-backend/src/database/DefaultProviderDatabase.ts b/plugins/catalog-backend/src/database/DefaultProviderDatabase.ts index 9dbdc5f516..f1073fd984 100644 --- a/plugins/catalog-backend/src/database/DefaultProviderDatabase.ts +++ b/plugins/catalog-backend/src/database/DefaultProviderDatabase.ts @@ -22,6 +22,8 @@ import lodash from 'lodash'; import { v4 as uuid } from 'uuid'; import type { Logger } from 'winston'; import { rethrowError } from './conversion'; +import { deleteWithEagerPruningOfChildren } from './operations/provider/deleteWithEagerPruningOfChildren'; +import { refreshByRefreshKeys } from './operations/provider/refreshByRefreshKeys'; import { checkLocationKeyConflict } from './operations/refreshState/checkLocationKeyConflict'; import { insertUnprocessedEntity } from './operations/refreshState/insertUnprocessedEntity'; import { updateUnprocessedEntity } from './operations/refreshState/updateUnprocessedEntity'; @@ -51,10 +53,10 @@ export class DefaultProviderDatabase implements ProviderDatabase { async transaction(fn: (tx: Transaction) => Promise): Promise { try { let result: T | undefined = undefined; - await this.options.database.transaction( async tx => { - // We can't return here, as knex swallows the return type in case the transaction is rolled back: + // We can't return here, as knex swallows the return type in case the + // transaction is rolled back: // https://github.com/knex/knex/blob/e37aeaa31c8ef9c1b07d2e4d3ec6607e557d800d/lib/transaction.js#L136 result = await fn(tx); }, @@ -63,7 +65,6 @@ export class DefaultProviderDatabase implements ProviderDatabase { doNotRejectOnRollback: true, }, ); - return result!; } catch (e) { this.options.logger.debug(`Error during transaction, ${e}`); @@ -76,128 +77,14 @@ export class DefaultProviderDatabase implements ProviderDatabase { options: ReplaceUnprocessedEntitiesOptions, ): Promise { const tx = txOpaque as Knex.Transaction; - const { toAdd, toUpsert, toRemove } = await this.createDelta(tx, options); if (toRemove.length) { - let removedCount = 0; - const rootId = () => { - if (tx.client.config.client.includes('mysql')) { - return tx.raw('CAST(NULL as UNSIGNED INT)', []); - } - - return tx.raw('CAST(NULL as INT)', []); - }; - for (const refs of lodash.chunk(toRemove, 1000)) { - /* - WITH RECURSIVE - -- All the nodes that can be reached downwards from our root - descendants(root_id, entity_ref) AS ( - SELECT id, target_entity_ref - FROM refresh_state_references - WHERE source_key = "R1" AND target_entity_ref = "A" - UNION - SELECT descendants.root_id, target_entity_ref - FROM descendants - JOIN refresh_state_references ON source_entity_ref = descendants.entity_ref - ), - -- All the nodes that can be reached upwards from the descendants - ancestors(root_id, via_entity_ref, to_entity_ref) AS ( - SELECT CAST(NULL as INT), entity_ref, entity_ref - FROM descendants - UNION - SELECT - CASE WHEN source_key IS NOT NULL THEN id ELSE NULL END, - source_entity_ref, - ancestors.to_entity_ref - FROM ancestors - JOIN refresh_state_references ON target_entity_ref = ancestors.via_entity_ref - ) - -- Start out with all of the descendants - SELECT descendants.entity_ref - FROM descendants - -- Expand with all ancestors that point to those, but aren't the current root - LEFT OUTER JOIN ancestors - ON ancestors.to_entity_ref = descendants.entity_ref - AND ancestors.root_id IS NOT NULL - AND ancestors.root_id != descendants.root_id - -- Exclude all lines that had such a foreign ancestor - WHERE ancestors.root_id IS NULL; - */ - removedCount += await tx('refresh_state') - .whereIn('entity_ref', function orphanedEntityRefs(orphans) { - return ( - orphans - // All the nodes that can be reached downwards from our root - .withRecursive('descendants', function descendants(outer) { - return outer - .select({ root_id: 'id', entity_ref: 'target_entity_ref' }) - .from('refresh_state_references') - .where('source_key', options.sourceKey) - .whereIn('target_entity_ref', refs) - .union(function recursive(inner) { - return inner - .select({ - root_id: 'descendants.root_id', - entity_ref: - 'refresh_state_references.target_entity_ref', - }) - .from('descendants') - .join('refresh_state_references', { - 'descendants.entity_ref': - 'refresh_state_references.source_entity_ref', - }); - }); - }) - // All the nodes that can be reached upwards from the descendants - .withRecursive('ancestors', function ancestors(outer) { - return outer - .select({ - root_id: rootId(), - via_entity_ref: 'entity_ref', - to_entity_ref: 'entity_ref', - }) - .from('descendants') - .union(function recursive(inner) { - return inner - .select({ - root_id: tx.raw( - 'CASE WHEN source_key IS NOT NULL THEN id ELSE NULL END', - [], - ), - via_entity_ref: 'source_entity_ref', - to_entity_ref: 'ancestors.to_entity_ref', - }) - .from('ancestors') - .join('refresh_state_references', { - target_entity_ref: 'ancestors.via_entity_ref', - }); - }); - }) - // Start out with all of the descendants - .select('descendants.entity_ref') - .from('descendants') - // Expand with all ancestors that point to those, but aren't the current root - .leftOuterJoin('ancestors', function keepaliveRoots() { - this.on( - 'ancestors.to_entity_ref', - '=', - 'descendants.entity_ref', - ); - this.andOnNotNull('ancestors.root_id'); - this.andOn('ancestors.root_id', '!=', 'descendants.root_id'); - }) - .whereNull('ancestors.root_id') - ); - }) - .delete(); - - await tx('refresh_state_references') - .where('source_key', '=', options.sourceKey) - .whereIn('target_entity_ref', refs) - .delete(); - } - + const removedCount = await deleteWithEagerPruningOfChildren({ + tx, + entityRefs: toRemove, + sourceKey: options.sourceKey, + }); this.options.logger.debug( `removed, ${removedCount} entities: ${JSON.stringify(toRemove)}`, ); @@ -258,15 +145,20 @@ export class DefaultProviderDatabase implements ProviderDatabase { const entityRef = stringifyEntityRef(entity); try { - let ok = await updateUnprocessedEntity(tx, entity, hash, locationKey); + let ok = await updateUnprocessedEntity({ + tx, + entity, + hash, + locationKey, + }); if (!ok) { - ok = await insertUnprocessedEntity( + ok = await insertUnprocessedEntity({ tx, entity, hash, - this.options.logger, locationKey, - ); + logger: this.options.logger, + }); } if (ok) { @@ -277,11 +169,11 @@ export class DefaultProviderDatabase implements ProviderDatabase { target_entity_ref: entityRef, }); } else { - const conflictingKey = await checkLocationKeyConflict( + const conflictingKey = await checkLocationKeyConflict({ tx, entityRef, locationKey, - ); + }); if (conflictingKey) { this.options.logger.warn( `Source ${options.sourceKey} detected conflicting entityRef ${entityRef} already referenced by ${conflictingKey} and now also ${locationKey}`, @@ -302,18 +194,7 @@ export class DefaultProviderDatabase implements ProviderDatabase { options: RefreshByKeyOptions, ) { const tx = txOpaque as Knex.Transaction; - const { keys } = options; - - await tx('refresh_state') - .whereIn('entity_id', function selectEntityRefs(tx2) { - tx2 - .whereIn('key', keys) - .select({ - entity_id: 'refresh_keys.entity_id', - }) - .from('refresh_keys'); - }) - .update({ next_update_at: tx.fn.now() }); + await refreshByRefreshKeys({ tx, keys: options.keys }); } private async createDelta( diff --git a/plugins/catalog-backend/src/database/operations/provider/deleteWithEagerPruningOfChildren.ts b/plugins/catalog-backend/src/database/operations/provider/deleteWithEagerPruningOfChildren.ts new file mode 100644 index 0000000000..5e2aa482a3 --- /dev/null +++ b/plugins/catalog-backend/src/database/operations/provider/deleteWithEagerPruningOfChildren.ts @@ -0,0 +1,155 @@ +/* + * Copyright 2022 The Backstage Authors + * + * 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 { Knex } from 'knex'; +import lodash from 'lodash'; +import { DbRefreshStateReferencesRow, DbRefreshStateRow } from '../../tables'; + +/** + * Given a number of entity refs originally created by a given entity provider + * (source key), remove those entities from the refresh state, and at the same + * time recursively remove every child that is a direct or indirect result of + * processing those entities, if they would have otherwise become orphaned by + * the removal of their parents. + */ +export async function deleteWithEagerPruningOfChildren(options: { + tx: Knex.Transaction; + entityRefs: string[]; + sourceKey: string; +}): Promise { + const { tx, entityRefs, sourceKey } = options; + let removedCount = 0; + + const rootId = () => + tx.raw( + tx.client.config.client.includes('mysql') + ? 'CAST(NULL as UNSIGNED INT)' + : 'CAST(NULL as UNSIGNED INT)', + [], + ); + + // Split up the operation by (large) chunks, so that we do not hit database + // limits for the number of permitted bindings on a precompiled statement + for (const refs of lodash.chunk(entityRefs, 1000)) { + /* + WITH RECURSIVE + -- All the nodes that can be reached downwards from our root + descendants(root_id, entity_ref) AS ( + SELECT id, target_entity_ref + FROM refresh_state_references + WHERE source_key = "R1" AND target_entity_ref = "A" + UNION + SELECT descendants.root_id, target_entity_ref + FROM descendants + JOIN refresh_state_references ON source_entity_ref = descendants.entity_ref + ), + -- All the nodes that can be reached upwards from the descendants + ancestors(root_id, via_entity_ref, to_entity_ref) AS ( + SELECT CAST(NULL as INT), entity_ref, entity_ref + FROM descendants + UNION + SELECT + CASE WHEN source_key IS NOT NULL THEN id ELSE NULL END, + source_entity_ref, + ancestors.to_entity_ref + FROM ancestors + JOIN refresh_state_references ON target_entity_ref = ancestors.via_entity_ref + ) + -- Start out with all of the descendants + SELECT descendants.entity_ref + FROM descendants + -- Expand with all ancestors that point to those, but aren't the current root + LEFT OUTER JOIN ancestors + ON ancestors.to_entity_ref = descendants.entity_ref + AND ancestors.root_id IS NOT NULL + AND ancestors.root_id != descendants.root_id + -- Exclude all lines that had such a foreign ancestor + WHERE ancestors.root_id IS NULL; + */ + removedCount += await tx('refresh_state') + .whereIn('entity_ref', function orphanedEntityRefs(orphans) { + return ( + orphans + // All the nodes that can be reached downwards from our root + .withRecursive('descendants', function descendants(outer) { + return outer + .select({ root_id: 'id', entity_ref: 'target_entity_ref' }) + .from('refresh_state_references') + .where('source_key', sourceKey) + .whereIn('target_entity_ref', refs) + .union(function recursive(inner) { + return inner + .select({ + root_id: 'descendants.root_id', + entity_ref: 'refresh_state_references.target_entity_ref', + }) + .from('descendants') + .join('refresh_state_references', { + 'descendants.entity_ref': + 'refresh_state_references.source_entity_ref', + }); + }); + }) + // All the nodes that can be reached upwards from the descendants + .withRecursive('ancestors', function ancestors(outer) { + return outer + .select({ + root_id: rootId(), + via_entity_ref: 'entity_ref', + to_entity_ref: 'entity_ref', + }) + .from('descendants') + .union(function recursive(inner) { + return inner + .select({ + root_id: tx.raw( + 'CASE WHEN source_key IS NOT NULL THEN id ELSE NULL END', + [], + ), + via_entity_ref: 'source_entity_ref', + to_entity_ref: 'ancestors.to_entity_ref', + }) + .from('ancestors') + .join('refresh_state_references', { + target_entity_ref: 'ancestors.via_entity_ref', + }); + }); + }) + // Start out with all of the descendants + .select('descendants.entity_ref') + .from('descendants') + // Expand with all ancestors that point to those, but aren't the current root + .leftOuterJoin('ancestors', function keepaliveRoots() { + this.on('ancestors.to_entity_ref', '=', 'descendants.entity_ref'); + this.andOnNotNull('ancestors.root_id'); + this.andOn('ancestors.root_id', '!=', 'descendants.root_id'); + }) + .whereNull('ancestors.root_id') + ); + }) + .delete(); + + // Delete the references that originate only from this entity provider. Note + // that there may be more than one entity provider making a "claim" for a + // given root entity, if they emit with the same location key. + await tx('refresh_state_references') + .where('source_key', '=', sourceKey) + .whereIn('target_entity_ref', refs) + .delete(); + } + + return removedCount; +} diff --git a/plugins/catalog-backend/src/database/operations/provider/refreshByRefreshKeys.ts b/plugins/catalog-backend/src/database/operations/provider/refreshByRefreshKeys.ts new file mode 100644 index 0000000000..dc34c9cdf0 --- /dev/null +++ b/plugins/catalog-backend/src/database/operations/provider/refreshByRefreshKeys.ts @@ -0,0 +1,43 @@ +/* + * Copyright 2022 The Backstage Authors + * + * 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 { Knex } from 'knex'; +import { DbRefreshStateRow } from '../../tables'; + +/** + * Schedules a future refresh of entities, by so called "refresh keys" that may + * be associated with one or more entities. Note that this does not mean that + * the refresh happens immediately, but rather that their scheduling time gets + * moved up the queue and will get picked up eventually by the regular + * processing loop. + */ +export async function refreshByRefreshKeys(options: { + tx: Knex.Transaction; + keys: string[]; +}): Promise { + const { tx, keys } = options; + + await tx('refresh_state') + .whereIn('entity_id', function selectEntityRefs(inner) { + inner + .whereIn('key', keys) + .select({ + entity_id: 'refresh_keys.entity_id', + }) + .from('refresh_keys'); + }) + .update({ next_update_at: tx.fn.now() }); +} diff --git a/plugins/catalog-backend/src/database/operations/refreshState/checkLocationKeyConflict.ts b/plugins/catalog-backend/src/database/operations/refreshState/checkLocationKeyConflict.ts index a256b36729..97d376f9a4 100644 --- a/plugins/catalog-backend/src/database/operations/refreshState/checkLocationKeyConflict.ts +++ b/plugins/catalog-backend/src/database/operations/refreshState/checkLocationKeyConflict.ts @@ -23,11 +23,13 @@ import { DbRefreshStateRow } from '../../tables'; * * @returns The conflicting key if there is one. */ -export async function checkLocationKeyConflict( - tx: Knex.Transaction, - entityRef: string, - locationKey?: string, -): Promise { +export async function checkLocationKeyConflict(options: { + tx: Knex.Transaction; + entityRef: string; + locationKey?: string; +}): Promise { + const { tx, entityRef, locationKey } = options; + const row = await tx('refresh_state') .select('location_key') .where('entity_ref', entityRef) diff --git a/plugins/catalog-backend/src/database/operations/refreshState/insertUnprocessedEntity.ts b/plugins/catalog-backend/src/database/operations/refreshState/insertUnprocessedEntity.ts index bf34fd2ce8..8e1c6edc87 100644 --- a/plugins/catalog-backend/src/database/operations/refreshState/insertUnprocessedEntity.ts +++ b/plugins/catalog-backend/src/database/operations/refreshState/insertUnprocessedEntity.ts @@ -25,13 +25,15 @@ import { isDatabaseConflictError } from '@backstage/backend-common'; * Attempts to insert a new refresh state row for the given entity, returning * true if successful and false if there was a conflict. */ -export async function insertUnprocessedEntity( - tx: Knex.Transaction, - entity: Entity, - hash: string, - logger: Logger, - locationKey?: string, -): Promise { +export async function insertUnprocessedEntity(options: { + tx: Knex.Transaction; + entity: Entity; + hash: string; + locationKey?: string; + logger: Logger; +}): Promise { + const { tx, entity, hash, logger, locationKey } = options; + const entityRef = stringifyEntityRef(entity); const serializedEntity = JSON.stringify(entity); diff --git a/plugins/catalog-backend/src/database/operations/refreshState/updateUnprocessedEntity.ts b/plugins/catalog-backend/src/database/operations/refreshState/updateUnprocessedEntity.ts index 084454de66..e26db4429b 100644 --- a/plugins/catalog-backend/src/database/operations/refreshState/updateUnprocessedEntity.ts +++ b/plugins/catalog-backend/src/database/operations/refreshState/updateUnprocessedEntity.ts @@ -24,12 +24,14 @@ import { DbRefreshStateRow } from '../../tables'; * * Updating the entity will also cause it to be scheduled for immediate processing. */ -export async function updateUnprocessedEntity( - tx: Knex.Transaction, - entity: Entity, - hash: string, - locationKey?: string, -): Promise { +export async function updateUnprocessedEntity(options: { + tx: Knex.Transaction; + entity: Entity; + hash: string; + locationKey?: string; +}): Promise { + const { tx, entity, hash, locationKey } = options; + const entityRef = stringifyEntityRef(entity); const serializedEntity = JSON.stringify(entity);