diff --git a/plugins/catalog-backend/migrations/20220616202842_refresh_keys.js b/plugins/catalog-backend/migrations/20220616202842_refresh_keys.js index ded53e4f67..b1b67f68f2 100644 --- a/plugins/catalog-backend/migrations/20220616202842_refresh_keys.js +++ b/plugins/catalog-backend/migrations/20220616202842_refresh_keys.js @@ -23,9 +23,9 @@ exports.up = async function up(knex) { 'This table contains relations between entities and keys to trigger refreshes with', ); table - .text('entity_ref') + .text('entity_id') .notNullable() - .references('entity_ref') + .references('entity_id') .inTable('refresh_state') .onDelete('CASCADE') .comment('A reference to the entity that the refresh key is tied to'); @@ -35,8 +35,7 @@ exports.up = async function up(knex) { .comment( 'A reference to a key which should be used to trigger a refresh on this entity', ); - table.unique(['entity_ref', 'key']); - table.index('entity_ref', 'refresh_keys_entity_ref_idx'); + table.index('entity_id', 'refresh_keys_entity_id_idx'); table.index('key', 'refresh_keys_key_idx'); }); }; @@ -46,7 +45,7 @@ exports.up = async function up(knex) { */ exports.down = async function down(knex) { await knex.schema.alterTable('refresh_keys', table => { - table.dropIndex([], 'refresh_keys_entity_ref_idx'); + table.dropIndex([], 'refresh_keys_entity_id_idx'); table.dropIndex([], 'refresh_keys_key_idx'); }); diff --git a/plugins/catalog-backend/src/database/DefaultProcessingDatabase.test.ts b/plugins/catalog-backend/src/database/DefaultProcessingDatabase.test.ts index f9f51b7a0d..d280f3a493 100644 --- a/plugins/catalog-backend/src/database/DefaultProcessingDatabase.test.ts +++ b/plugins/catalog-backend/src/database/DefaultProcessingDatabase.test.ts @@ -522,7 +522,7 @@ describe('Default Processing Database', () => { ); const refreshKeys = await knex('refresh_keys') - .where({ entity_ref: stringifyEntityRef(processedEntity) }) + .where({ entity_id: id }) .select(); expect(refreshKeys[0]).toEqual({ diff --git a/plugins/catalog-backend/src/database/DefaultProcessingDatabase.ts b/plugins/catalog-backend/src/database/DefaultProcessingDatabase.ts index 5b191dce32..f0779212be 100644 --- a/plugins/catalog-backend/src/database/DefaultProcessingDatabase.ts +++ b/plugins/catalog-backend/src/database/DefaultProcessingDatabase.ts @@ -79,6 +79,7 @@ export class DefaultProcessingDatabase implements ProcessingDatabase { errors, relations, deferredEntities, + refreshKeys, locationKey, } = options; const refreshResult = await tx('refresh_state') @@ -141,17 +142,19 @@ export class DefaultProcessingDatabase implements ProcessingDatabase { BATCH_SIZE, ); + // Delete old refresh keys + await tx('refresh_keys') + .where({ entity_id: id }) + .delete(); + // Insert the refresh keys for the processed entity - await Promise.all( - options.refreshKeys.map(k => { - return tx('refresh_keys') - .insert({ - entity_ref: sourceEntityRef, - key: k.key, - }) - .onConflict(['entity_ref', 'key']) - .ignore(); - }), + await tx.batchInsert( + 'refresh_keys', + refreshKeys.map(k => ({ + entity_id: id, + key: k.key, + })), + BATCH_SIZE, ); return { @@ -540,18 +543,18 @@ export class DefaultProcessingDatabase implements ProcessingDatabase { const { keys } = options; const updateResult = await tx('refresh_state') - .whereIn('entity_ref', function selectEntityRefs(tx2) { + .whereIn('entity_id', function selectEntityRefs(tx2) { tx2 .whereIn('key', keys) .select({ - entity_ref: 'refresh_keys.entity_ref', + entity_id: 'refresh_keys.entity_id', }) .from('refresh_keys'); }) .update({ next_update_at: tx.fn.now() }); if (updateResult === 0) { - throw new NotFoundError( + this.options.logger.info( `Failed to schedule ${JSON.stringify(keys)} for keys`, ); } diff --git a/plugins/catalog-backend/src/database/tables.ts b/plugins/catalog-backend/src/database/tables.ts index 33c3ab0ef3..b9fe12be11 100644 --- a/plugins/catalog-backend/src/database/tables.ts +++ b/plugins/catalog-backend/src/database/tables.ts @@ -44,7 +44,7 @@ export type DbRefreshStateRow = { }; export type DbRefreshKeysRow = { - entity_ref: string; + entity_id: string; key: string; }; diff --git a/plugins/catalog-backend/src/database/types.ts b/plugins/catalog-backend/src/database/types.ts index da7aa65714..8a839c0650 100644 --- a/plugins/catalog-backend/src/database/types.ts +++ b/plugins/catalog-backend/src/database/types.ts @@ -82,10 +82,6 @@ export type ReplaceUnprocessedEntitiesOptions = type: 'delta'; }; -export type RefreshKeyOptions = { - refreshKeys: RefreshKeyData[]; -}; - export type RefreshByKeyOptions = { keys: string[]; };