From f506a945753e5b2a8045bf41a7c5ad602f4e00cb Mon Sep 17 00:00:00 2001 From: Damon Kaswell Date: Mon, 28 Nov 2022 12:25:37 -0800 Subject: [PATCH] Use last ingestion cycle to generate total Signed-off-by: Damon Kaswell --- .../IncrementalIngestionDatabaseManager.ts | 45 +++++++++++++++---- 1 file changed, 36 insertions(+), 9 deletions(-) diff --git a/plugins/catalog-backend-module-incremental-ingestion/src/database/IncrementalIngestionDatabaseManager.ts b/plugins/catalog-backend-module-incremental-ingestion/src/database/IncrementalIngestionDatabaseManager.ts index 03a749d70c..fd99f1207d 100644 --- a/plugins/catalog-backend-module-incremental-ingestion/src/database/IncrementalIngestionDatabaseManager.ts +++ b/plugins/catalog-backend-module-incremental-ingestion/src/database/IncrementalIngestionDatabaseManager.ts @@ -113,6 +113,25 @@ export class IncrementalIngestionDatabaseManager { }); } + /** + * Finds the last ingestion record for the named provider. + * @param provider - string + * @returns IngestionRecord | undefined + */ + async getPreviousIngestionRecord(provider: string) { + return await this.client.transaction(async tx => { + const record = await tx('ingestions') + .where('provider_name', provider) + .andWhereNot('completion_ticket', 'open') + .first(); + if (!record) { + // This is the first time this entity provider has run. Return the current record. + return await this.getCurrentIngestionRecord(provider); + } + return record; + }); + } + /** * Removes all entries from `ingestion_marks_entities`, `ingestion_marks`, and `ingestions` * for prior ingestions that completed (i.e., have a `completion_ticket` value other than 'open'). @@ -286,16 +305,24 @@ export class IncrementalIngestionDatabaseManager { * @returns All entities to remove for this burst. */ async computeRemoved(provider: string, ingestionId: string) { + const previousIngestion = (await this.getPreviousIngestionRecord( + provider, + )) as IngestionRecord; return await this.client.transaction(async tx => { - const rows = await tx('final_entities') - .count({ total: '*' }) - .join('search', 'search.entity_id', 'final_entities.entity_id') - .where( - 'search.key', - `metadata.annotations.${INCREMENTAL_ENTITY_PROVIDER_ANNOTATION}`, - ) - .andWhere('search.value', provider); - const total = rows.reduce((acc, cur) => acc + (cur.total as number), 0); + let total = 0; + if (previousIngestion.id !== ingestionId) { + const rows = await tx('ingestion_mark_entities') + .count({ total: '*' }) + .join( + 'ingestion_marks', + 'ingestion_marks.id', + 'ingestion_mark_entities.ingestion_mark_id', + ) + .join('ingestions', 'ingestions.id', 'ingestion_marks.ingestion_id') + .where('ingestions.id', previousIngestion.id); + + total = rows.reduce((acc, cur) => acc + (cur.total as number), 0); + } const removed: { entity: string; ref: string }[] = await tx( 'final_entities', )