diff --git a/plugins/catalog-backend/src/database/types.ts b/plugins/catalog-backend/src/database/types.ts index 07241f0682..2d30f8bd86 100644 --- a/plugins/catalog-backend/src/database/types.ts +++ b/plugins/catalog-backend/src/database/types.ts @@ -157,6 +157,14 @@ export interface ProcessingDatabase { */ refresh(txOpaque: Transaction, options: RefreshOptions): Promise; + /** + * Schedules a refresh for every entity that has a matching set of refresh key stored for it. + */ + refreshByRefreshKeys( + txOpaque: Transaction, + options: RefreshByKeyOptions, + ): Promise; + /** * Lists all ancestors of a given entityRef. * diff --git a/plugins/catalog-backend/src/processing/connectEntityProviders.ts b/plugins/catalog-backend/src/processing/connectEntityProviders.ts index 70a43de589..8325e2bcac 100644 --- a/plugins/catalog-backend/src/processing/connectEntityProviders.ts +++ b/plugins/catalog-backend/src/processing/connectEntityProviders.ts @@ -61,6 +61,16 @@ class Connection implements EntityProviderConnection { } } + async refresh(keys: string[]): Promise { + const db = this.config.processingDatabase; + + await db.transaction(async (tx: any) => { + return db.refreshByRefreshKeys(tx, { + keys, + }); + }); + } + private check(entities: Entity[]) { for (const entity of entities) { try {