diff --git a/plugins/catalog-backend/src/next/NextCatalogBuilder.ts b/plugins/catalog-backend/src/next/NextCatalogBuilder.ts index 6685296312..944036e48c 100644 --- a/plugins/catalog-backend/src/next/NextCatalogBuilder.ts +++ b/plugins/catalog-backend/src/next/NextCatalogBuilder.ts @@ -263,11 +263,11 @@ export class NextCatalogBuilder { const db = new CommonDatabase(dbClient, logger); - const processingDatabase = new DefaultProcessingDatabase( - dbClient, + const processingDatabase = new DefaultProcessingDatabase({ + database: dbClient, logger, - this.refreshIntervalSeconds, - ); + refreshIntervalSeconds: this.refreshIntervalSeconds, + }); const integrations = ScmIntegrations.fromConfig(config); const orchestrator = new DefaultCatalogProcessingOrchestrator({ processors, diff --git a/plugins/catalog-backend/src/next/database/DefaultProcessingDatabase.test.ts b/plugins/catalog-backend/src/next/database/DefaultProcessingDatabase.test.ts index 5843ef6c34..0f83311944 100644 --- a/plugins/catalog-backend/src/next/database/DefaultProcessingDatabase.test.ts +++ b/plugins/catalog-backend/src/next/database/DefaultProcessingDatabase.test.ts @@ -39,7 +39,11 @@ describe('Default Processing Database', () => { await DatabaseManager.createDatabase(knex); return { knex, - db: new DefaultProcessingDatabase(knex, logger, 100), + db: new DefaultProcessingDatabase({ + database: knex, + logger, + refreshIntervalSeconds: 100, + }), }; } diff --git a/plugins/catalog-backend/src/next/database/DefaultProcessingDatabase.ts b/plugins/catalog-backend/src/next/database/DefaultProcessingDatabase.ts index 10c213d218..f4804f9118 100644 --- a/plugins/catalog-backend/src/next/database/DefaultProcessingDatabase.ts +++ b/plugins/catalog-backend/src/next/database/DefaultProcessingDatabase.ts @@ -44,9 +44,11 @@ const BATCH_SIZE = 50; export class DefaultProcessingDatabase implements ProcessingDatabase { constructor( - private readonly database: Knex, - private readonly logger: Logger, - private readonly refreshIntervalSeconds: number, + private readonly options: { + database: Knex; + logger: Logger; + refreshIntervalSeconds: number; + }, ) {} async updateProcessedEntity( @@ -273,7 +275,7 @@ export class DefaultProcessingDatabase implements ProcessingDatabase { .whereIn('target_entity_ref', toRemove) .delete(); - this.logger.debug( + this.options.logger.debug( `removed, ${removedCount} entities: ${JSON.stringify(toRemove)}`, ); } @@ -379,11 +381,11 @@ export class DefaultProcessingDatabase implements ProcessingDatabase { next_update_at: tx.client.config.client === 'sqlite3' ? tx.raw(`datetime('now', ?)`, [ - `${this.refreshIntervalSeconds} seconds`, + `${this.options.refreshIntervalSeconds} seconds`, ]) : tx.raw( `now() + interval '${Number( - this.refreshIntervalSeconds, + this.options.refreshIntervalSeconds, )} seconds'`, ), }); @@ -413,7 +415,7 @@ export class DefaultProcessingDatabase implements ProcessingDatabase { try { let result: T | undefined = undefined; - await this.database.transaction( + 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: // https://github.com/knex/knex/blob/e37aeaa31c8ef9c1b07d2e4d3ec6607e557d800d/lib/transaction.js#L136 @@ -427,7 +429,7 @@ export class DefaultProcessingDatabase implements ProcessingDatabase { return result!; } catch (e) { - this.logger.debug(`Error during transaction, ${e}`); + this.options.logger.debug(`Error during transaction, ${e}`); if ( /SQLITE_CONSTRAINT: UNIQUE/.test(e.message) ||